You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何实现每次使用生成的币种符号运行币安直播流函数?

币安实时涨幅币种WebSocket监控实现

核心需求拆解

  • 涨幅榜扫描模块:定期遍历币安USDT现货交易对,计算5分钟周期的价格涨幅,筛选出涨幅最高的币种
  • 实时K线拉取模块:针对涨幅榜输出的新币种,通过币安WebSocket接口拉取1分钟K线数据,完成K线后记录数据到DataFrame

原代码问题与优化方案

原代码实现了基础逻辑,但存在几个关键问题:

  1. 异步函数中使用同步time.sleep阻塞事件循环
  2. 注释标注取Top3涨幅币种,但实际代码只取了Top1
  3. 全局DataFrame在多任务场景下存在线程安全风险
  4. 异步消息处理函数调用未加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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.08 07:10:08