如何实现每次使用生成的币种符号运行币安直播流函数?
币安实时涨幅币种WebSocket监控实现
核心需求拆解
- 涨幅榜扫描模块:定期遍历币安USDT现货交易对,计算5分钟周期的价格涨幅,筛选出涨幅最高的币种
- 实时K线拉取模块:针对涨幅榜输出的新币种,通过币安WebSocket接口拉取1分钟K线数据,完成K线后记录数据到DataFrame
原代码问题与优化方案
原代码实现了基础逻辑,但存在几个关键问题:
- 异步函数中使用同步
time.sleep阻塞事件循环 - 注释标注取Top3涨幅币种,但实际代码只取了Top1
- 全局DataFrame在多任务场景下存在线程安全风险
- 异步消息处理函数调用未加
await
优化后完整代码
import json import asyncio import websockets import pandas as pd from binance.client import Client import numpy as np # 替换为你的币安API密钥 API_KEY = "your_api_key_here" API_SECRET = "your_api_secret_here" client = Client(api_key=API_KEY, api_secret=API_SECRET, tld='com', testnet=True) interval = '1m' # 存储K线数据的DataFrame(生产环境建议改用数据库) kline_data = pd.DataFrame(columns=['Date','Coin','Open', 'High', 'Low', 'Close', 'Volume', 'no_trades','Quote_Volume','Complete']) async def top_gainer_coins(top_n=3): try: print("更新涨幅榜数据...") # 异步延迟,避免阻塞事件循环 await asyncio.sleep(90) exchange_info = client.get_exchange_info() spot_symbols = [ symbol['symbol'] for symbol in exchange_info['symbols'] if symbol['isSpotTradingAllowed'] and symbol['quoteAsset'] == 'USDT' and symbol['status'] == 'TRADING' ] all_data = [] for symbol in spot_symbols: try: # 获取最近12根5分钟K线 klines = client.get_klines(symbol=symbol, interval='5m', limit=12) data = pd.DataFrame(klines).copy() data = data.iloc[:, [0, 4, 5, 8]] data.columns = ['Date', 'Close', 'Volume', 'NoTrades'] data['Date'] = pd.to_datetime(data['Date'], unit='ms') for col in data.columns[1:]: data[col] = pd.to_numeric(data[col]) # 计算最新一根K线的涨幅 data['pricegain'] = (data.Close.pct_change() * 100).fillna(0) data['VGain'] = (data.Volume.pct_change() * 100).fillna(0) data['TPS'] = (data.NoTrades / 60).fillna(0) data.insert(1, 'Symbol', symbol) # 只保留最后一根完整K线的涨幅数据 data = data.iloc[-1:] data = data.replace([np.inf, -np.inf], np.nan) all_data.append(data) except Exception as e: print(f"获取{symbol}数据失败: {e}") continue if all_data: all_data = pd.concat(all_data).dropna() all_data = all_data.sort_values(by='pricegain', ascending=False) all_data.reset_index(drop=True, inplace=True) all_data.index += 1 # 获取Top N涨幅币种 top_gainers = all_data.head(top_n) top_gainer_symbols = top_gainers['Symbol'].tolist() print(f"当前Top {top_n}涨幅币种: {top_gainer_symbols}") print(top_gainers.to_string()) return top_gainer_symbols except Exception as e: print(f"涨幅榜更新失败: {e}") return [] async def process_kline_message(msg): try: msg_data = json.loads(msg) kline = msg_data['k'] start_time = pd.to_datetime(kline['t'], unit='ms') symbol = kline['s'] open_price = float(kline['o']) high = float(kline['h']) low = float(kline['l']) close = float(kline['c']) volume = float(kline['v']) no_trades = float(kline['n']) quote_volume = float(kline['q']) complete = kline['x'] if complete: global kline_data # 添加完整K线数据 kline_data.loc[len(kline_data)] = [ start_time, symbol, open_price, high, low, close, volume, no_trades, quote_volume, complete ] print(f"已记录{symbol}完整1分钟K线:\n{kline_data.tail(1)}") except Exception as e: print(f"处理K线消息失败: {e}") async def run_ws(symbol): print(f"启动{symbol}实时K线监控...") try: ws_symbol = symbol.lower() socket_url = f"wss://stream.binance.com:9443/ws/{ws_symbol}@kline_{interval}" async with websockets.connect(socket_url) as ws: while True: response = await ws.recv() await process_kline_message(response) # 避免消息处理过快,可根据需求调整 await asyncio.sleep(0.1) except Exception as e: print(f"{symbol}连接异常: {e}") async def main(): # 存储当前运行的WebSocket任务,避免重复连接 active_tasks = {} while True: try: new_symbols = await top_gainer_coins(top_n=3) if new_symbols: # 关闭不在新涨幅榜中的币种连接 for symbol in list(active_tasks.keys()): if symbol not in new_symbols: active_tasks[symbol].cancel() del active_tasks[symbol] print(f"停止{symbol}的K线监控") # 启动新币种的WebSocket连接 for symbol in new_symbols: if symbol not in active_tasks: task = asyncio.create_task(run_ws(symbol)) active_tasks[symbol] = task # 后台运行所有任务,不阻塞主循环 await asyncio.gather(*active_tasks.values(), return_exceptions=True) except Exception as e: print(f"主循环异常: {e}") await asyncio.sleep(30) if __name__ == "__main__": try: asyncio.run(main()) except KeyboardInterrupt: print("程序已终止") except Exception as e: print(f"程序启动失败: {e}")
关键优化说明
- 替换同步
time.sleep为异步asyncio.sleep,避免阻塞整个事件循环 - 修正Top3涨幅币种的筛选逻辑,匹配需求描述
- 改进消息处理函数的异步调用方式,确保逻辑正确
- 添加活跃任务管理,自动关闭不在涨幅榜中的币种连接
- 优化K线数据记录逻辑,只保存完整的1分钟K线
内容的提问来源于stack exchange,提问作者Suliman Zunni
相关产品推荐
相关产品推荐

