python-binance库user_timeout参数不生效,如何实现每秒存1条数据?
问题原因
你误解了user_timeout参数的作用:该参数是WebSocket连接无消息超时断开阈值,单位为毫秒,设置为1000代表如果1秒内没有收到币安服务端推送的消息就主动断开连接,和消息推送频率、数据库写入频率完全无关,因此配置后不会生效。币安交易WebSocket会在每笔成交发生时推送消息,BTCUSDT交易活跃时每秒推送数十条属于正常情况。
实现每秒写入1条记录的方案
推荐使用定时采样方案:不阻塞WebSocket消息接收,仅将每秒内最新的成交数据写入数据库,既可以达到每秒写1条的要求,也不会丢失最新的成交价格信息。
修改后的完整代码
import pandas as pd import sqlalchemy import asyncio from binance import AsyncClient, BinanceSocketManager # 存储每秒内最新的成交数据 latest_frame = None def createframe(msg): df = pd.DataFrame([msg]) df = df.loc[:, ['s', 'E', 'p']] df.columns = ['symbol', 'Time', 'Price'] df.Price = df.Price.astype(float) df.Time = pd.to_datetime(df.Time, unit='ms') return df async def write_to_db(): """每秒定时写入最新数据到数据库的协程""" global latest_frame while True: if latest_frame is not None: print(latest_frame) latest_frame.to_sql('BTCUSDT', 'sqlite:///stream.db', if_exists='append', index=False) await asyncio.sleep(1) async def main(): client = await AsyncClient.create() bm = BinanceSocketManager(client) # 启动写入协程 asyncio.create_task(write_to_db()) ts = bm.trade_socket('BTCUSDT') async with ts as tscm: while True: res = await tscm.recv() # 只更新最新数据,不直接写入 global latest_frame latest_frame = createframe(res) await client.close_connection() if __name__ == "__main__": loop = asyncio.get_event_loop() loop.run_until_complete(main())
可选简化方案(不推荐)
如果你不需要保留每秒内的中间价格更新,也可以在每次写入后增加1秒等待实现,缺点是会导致WebSocket消息积压,长时间运行可能出现内存占用过高的问题:
# 只需修改main函数内的接收逻辑即可 async with ts as tscm: while True: res = await tscm.recv() frame = createframe(res) print(frame) frame.to_sql('BTCUSDT', 'sqlite:///stream.db', if_exists='append', index=False) # 等待1秒再处理下一条消息 await asyncio.sleep(1)
内容的提问来源于stack exchange,提问作者Jan Nęciński
相关产品推荐
相关产品推荐

