Gestire più connessioni websocket

Sep 15 2020

Ho il seguente codice di base, che si connette a un server websocket e riceve alcuni dati:

import websocket, json, time

def process_message(ws, msg):
    message = json.loads(msg)
    print(message)

def on_error(ws, error):
    print('Error', e)

def on_close(ws):
    print('Closing')

def on_open(ws):
    def run(*args):
        Subs = []
       
        tradeStr=  """{"method": "SUBSCRIBE", "params":%s, "id": 1}"""%(json.dumps(Subs))
        ws.send(tradeStr)

    thread.start_new_thread(run, ())

def Connect():
    websocket.enableTrace(False)
    ws = websocket.WebSocketApp("wss://myurl", on_message = process_message, on_error = on_error, on_close = on_close)
    ws.on_open = on_open
    ws.run_forever()

Connect()

Ora, vorrei creare più connessioni a server diversi e ricevere dati contemporaneamente nello stesso script. Ho provato quanto segue:

def run(url):

    def process_message(ws, msg):
        message = json.loads(msg)
        print(message)

    def on_error(ws, error):
        print('Error', e)

    def on_close(ws):
        print('Closing')

    def on_open(ws):
        def run(*args):
            Subs = []
           
            tradeStr=  """{"method": "SUBSCRIBE", "params":%s, "id": 1}"""%(json.dumps(Subs))
            ws.send(tradeStr)

        thread.start_new_thread(run, ())

    def Connect():
        websocket.enableTrace(False)
        ws = websocket.WebSocketApp(url, on_message = process_message, on_error = on_error, on_close = on_close)
        ws.on_open = on_open
        ws.run_forever()

    Connect()

threading.Thread(target=run, kwargs={'url': 'url1'}).start()
threading.Thread(target=run, kwargs={'url': 'url2'}).start()
threading.Thread(target=run, kwargs={'url': 'url3'}).start()

Ora, questo codice funziona, ma mi sto collegando a URL diversi e sto trasmettendo dati da tutti loro, ma mi è sembrata una soluzione "hacky". Inoltre non so se quello che sto facendo potrebbe essere una cattiva pratica oppure no. Ogni connessione invierà circa 600/700 piccoli dizionari JSON e devo aggiornare ogni record nel db.

Quindi la mia domanda è: questa implementazione va bene? Dato che funziona con i thread, può creare problemi a lungo termine? Dovrei fare un'altra libreria come Tornado?

Risposte

3 hjpotter92 Sep 19 2020 at 02:44

Benvenuto nella community di revisione del codice. Pochi pensieri iniziali, quando si scrive codice Python; seguendo la guida allo stile PEP-8 rende il codice più manutenibile. Che equivarrebbe a (ma non limitato a):

  1. Funzioni e variabili denominate in lower_snake_case.
  2. Evitare spazi bianchi nei parametri / argomenti delle funzioni.

Andando avanti, puoi inserire il codice di gestione del tuo websocket nel suo sottomodulo. Per es.

class WebSocket:
    def __init__(self, url):
        self._url = url

    def close(self, ws):
        # handle closed connection callback
        pass

    def connect(self):
        self.ws = websocket.WebSocketApp(self._url, on_close=self.close)
        self.ws.on_open = self.open  # for eg.
        self.ws.run_forever()
    .
    .
    .

Se vuoi che supporti il ​​threading, puoi estendere questa classe con la threading.Threadsottoclasse (puoi anche scegliere di sottoclassare in multiprocessing.Processseguito):

class ThreadedWebSocket(WebSocket, threading.Thread):
    def run(self):
        # This simply calls `self.connect()`

in questo modo, se vuoi inizializzare i thread:

url1_thread = ThreadedWebSocket(url=url1)
url1_thread.start()
url1_thread.join()

La sottoclasse è principalmente una questione di preferenza.


Nella on_openrichiamata, stai scaricando un dict su json e quindi inserendolo in un modello di stringa. Quanto segue ottiene lo stesso risultato ed è più intuitivo (imo):

def on_open(ws):
    def run(*args):
        subs = []
        trade =  {
            "method": "SUBSCRIBE",
            "params": subs,
            "id": 1
        }
        ws.send(json.dumps(trade))

Venendo alla tua domanda sulle alternative di libreria e sulle prestazioni, l'I / O di rete è generalmente una chiamata di blocco. Questo è stato uno dei motivi principali per cui ho suggerito di utilizzare il websocketspacchetto. L'altro essere; utilizziamo websocket nella nostra organizzazione da oltre un anno in produzione. E ha funzionato in modo davvero impressionante per noi, gestendo un throughput di pochi GB di dati ogni giorno con una singola macchina di basso livello.

3 Coupcoup Sep 19 2020 at 03:26

Questo è lo stesso approccio che hai usato, appena riscritto come approccio oop per pacmaninbw con una classe che socketConneredita dawebsocket.WebSocketApp

Poiché il on_message, on_errore on_closefunzioni definite sono così brevi / semplice è possibile impostare direttamente nella init sovraccarico con funzioni lambda invece di scrivere su ciascuno separatamente e passandoli come parametri

Si formano nuove connessioni con ws = SocketConn('url'). Dovresti anche eseguire il loop su un elenco o un insieme di URL se ti connetti a più (cosa che probabilmente stai facendo nel tuo codice effettivo ma che vale la pena notare per ogni evenienza)

import websocket, json, threading, _thread

class SocketConn(websocket.WebSocketApp):
    def __init__(self, url):
        #WebSocketApp sets on_open=None, pass socketConn.on_open to init
        super().__init__(url, on_open=self.on_open)

        self.on_message = lambda ws, msg: print(json.loads(msg))
        self.on_error = lambda ws, e: print('Error', e)
        self.on_close = lambda ws: print('Closing')

        self.run_forever()

    def on_open(self, ws):
        def run(*args):
            Subs = []
            
            tradeStr={"method":"SUBSCRIBE", "params":Subs, "id":1}
            ws.send(json.dumps(tradeStr))

        _thread.start_new_thread(run, ())



if __name__ == '__main__':
    websocket.enableTrace(False)

    urls = {'url1', 'url2', 'url3'}


    for url in urls:
        # args=(url,) stops the url from being unpacked as chars
        threading.Thread(target=SocketConn, args=(url,)).start()