websocket multiple stream python

Viewed 560

Here is my problem : I'm trying to get the price of three pairs of crypto on Binance using websocket

My problem is that I have difficulties in managing the fact that there are multiple streams (from the three pairs) whereas with just one it works perfectly well.

Here is my code:

import websocket, json, pprint, talib, numpy

cc = ["btcusdt","adausdt","solusdt"]
interval = "1m"


for i in cc:
    socket = f'wss://stream.binance.com:9443/ws/{i}@kline_{interval}'
    
    def on_message(ws, message):
        json_message = json.loads(message)
        candle = json_message['k']
        close = candle["c"]
        print(i,': ',close)

    def on_close(ws):
        print("Connection Closed")


    wsapp = websocket.WebSocketApp(socket, on_message=on_message, on_close=on_close)
    wsapp.run_forever()

This is the results that I'm getting each second :

btcusdt :  47009.01000000
btcusdt :  47009.02000000
btcusdt :  47004.00000000
...

However this is not what I want, it takes into account the first pair only. I am looking for this kind of result :

btcusdt : #price
adausdt : #price
solusdt : #price
btcusdt : #price
adausdt : #price
solusdt : #price
...

I know that my problem comes from the wsapp.run_forever() because by doing that it takes the first pair only without looping to the others. But I don't know how to manage it and I would be very glad if somebody has a solution on managing this issue.

Thank you very much

1 Answers

I couldn't get what you want through your code but I made one that does what you need, check it out if it's helpful. See you later.

import json
import websocket
import sqlalchemy
import pandas as pd

#websocket.enableTrace(True) #DEBUG
engine = sqlalchemy.create_engine('sqlite:///CESTA_MOEDAS_raw.db') #Chamada SQL
pair_coins = 'ethusdt@kline_1m/btcusdt@kline_1m/bnbusdt@kline_1m/ethbtc@kline_1m'

def on_open(ws):
    print("open")

def on_message(ws, message):
        json_message = json.loads(message)
        candle = json_message['data']['k']
        df = pd.DataFrame([candle])

        #print(df.keys()) # DEBUG: 
        #print(df.values()) # DEBUG:

        df = df.loc[:, ['t', 'T', 's', 'i', 'f', 'L', 'o', 'c', 'h', 'l', 'v', 'n', 'x', 'q', 'V', 'Q', 'B']]  
        df.columns = ['start_time', 'close_time', 'Symbol', 'Interval', ' First_trade_ID', 'Last_trade_ID', 'Open_price', 'Close_price', 'High_price', 'Low_price', 'Base_asset_volume', 'Number_of_trades', 'Candle_close_price', 'Quote_asset_volume', 'Taker_buy_base_asset_volume', 'Taker_buy_quote_asset_volume', 'Ignore']  
        df = df.loc[:, ['close_time', 'Symbol', 'Close_price']]
        df.Close_price = df.Close_price.astype(float)  
        df.close_time = pd.to_datetime(df.close_time, unit='ms')  
        print(df) 

        #Salvando os dados no DbSQL
        frame = df
        frame.to_sql('CESTA_MOEDAS', engine, if_exists='append', index=False)

def on_close(ws, close_status_code, close_msg):
    print("closed")

SOCK = f"wss://stream.binance.com:9443/stream?streams={pair_coins}"
ws = websocket.WebSocketApp(SOCK, on_open=on_open, on_close=on_close, on_message=on_message)

ws.run_forever()
Related