【发布时间】:2022-09-23 02:51:47
【问题描述】:
我正在尝试同时从 websocket 服务器的两个不同通道中提取数据。特别是,我正在从 FTX 交易所提取两对加密货币 BTC/USD 和 ZRX/USD 的订单簿。 Websocket 发送订单簿的初始快照,然后发送更新。在每条消息中,我都会应用更新。
我想将订单簿存储在字典中(也许有更好的数据结构?)并且我正在使用multprocessing 在两个不同的进程中运行两个 websocket 连接。我正在定义两个共享内存字典,asks 和bids,它将把询价和出价存储为字典的字典。
例如,asks = {\'ZRX/USD\':{best_askZRX: quantity_bestZRX, second_best_askZRX: quantity_second_bestZRX,..},\'BTC/USD\':{best_askBTC: quantity_bestBTC, second_best_askBTC: quantity_second_bestBTC,..} } 和类似的出价。
在 main 函数中,我定义了字典和两个进程
manager = Manager()
asks = manager.dict()
bids = manager.dict()
asks[\'ZRX/USD\'] = manager.dict()
bids[\'ZRX/USD\'] = manager.dict()
asks[\'BTC/USD\'] = manager.dict()
bids[\'BTC/USD\'] = manager.dict()
extract_data1 = Process(target=get_orderbook, args=(\'ZRX/USD\', asks, bids))
extract_data2 = Process(target=get_orderbook, args=(\'BTC/USD\', asks, bids))
每个进程调用代码中定义的函数get_orderbook。然后我运行这些进程。但是,这些进程不会并行运行,而是 extract_data1 先启动,extract_data2 仅在我手动终止 extract_data1 时启动。请参阅下面的代码。这里发生了什么?我究竟做错了什么?
from multiprocessing import Process, Manager
import websocket, json
import time
def on_open(ws, market):
print(\'opened connection\')
subscribe_message = {\'op\': \'subscribe\', \'channel\': \'orderbook\', \'market\': market}
print(subscribe_message)
ws.send(json.dumps(subscribe_message))
def init_orderbook(ws, asks, bids, market, asks_snapshot, bids_snapshot):
proxy_asks, proxy_bids = asks, bids
for level in asks_snapshot:
proxy_asks[level[0]] = level[1]
for level in bids_snapshot:
proxy_bids[level[0]] = level[1]
asks, bids = proxy_asks, proxy_bids
def apply_changes(ws, asks, bids, market, asks_snapshot, bids_snapshot):
proxy_asks, proxy_bids = asks, bids
for level in asks_snapshot:
if level[1] == 0:
del proxy_asks[level[0]]
else:
proxy_asks[level[0]] = level[1]
for level in bids_snapshot:
if level[1] == 0:
del proxy_bids[level[0]]
else:
proxy_bids[level[0]] = level[1]
asks, bids = proxy_asks, proxy_bids
def on_message(ws, message, asks, bids):
js = json.loads(message)
if js[\'type\'] == \'subscribed\':
print(\'Subscribed, \', js)
elif js[\'type\'] == \'update\':
market = js[\'market\']
update_asks = js[\'data\'][\'asks\']
update_bids = js[\'data\'][\'asks\']
apply_changes(ws, asks[market], bids[market], market, update_asks, update_bids)
elif js[\'type\'] == \'partial\':
print(\'Storing snapshot, \', js)
market = js[\'market\']
update_asks = js[\'data\'][\'asks\']
update_bids = js[\'data\'][\'asks\']
init_orderbook(ws, asks[market], bids[market], market, update_asks, update_bids)
def get_orderbook(market, asks, bids):
socket = \"wss://ftx.com/ws/\"
ws = websocket.WebSocketApp(socket)
ws.on_open = lambda *x: on_open(*x, market)
ws.on_message = lambda ws, msg: on_message(ws, msg, asks, bids)
ws.run_forever()
if __name__ == \'__main__\':
manager = Manager()
asks = manager.dict()
bids = manager.dict()
asks[\'ZRX/USD\'] = manager.dict()
bids[\'ZRX/USD\'] = manager.dict()
asks[\'BTC/USD\'] = manager.dict()
bids[\'BTC/USD\'] = manager.dict()
extract_data1 = Process(target=get_orderbook, args=(\'ZRX/USD\', asks, bids))
extract_data2 = Process(target=get_orderbook, args=(\'BTC/USD\', asks, bids))
extract_data1.run()
extract_data2.run()
extract_data1.join()
extract_data2.join()
-
难道
wss://ftx.com/ws/一次只接受一个连接?
标签: python python-3.x multiprocessing