如何将KuCoin WebSocket接收的K线数据保存至变量后存入数据库
实现方案
核心需求是跨脚本共享实时K线数据,推荐两种落地方式,优先选Redis方案,适配性最高。
方案1:Redis作为中间缓存(最推荐,跨语言/跨机器都支持)
实现逻辑
- Redis是内存数据库,读写速度快,支持多种数据结构,刚好适配K线数据的存储、去重、跨进程读取需求
- 收到WebSocket推送的K线后直接写入Redis,其他所有脚本直接连Redis读取即可,完全解耦
- 还可以设置数据过期时间,避免内存占用过高
前置准备
安装依赖:
pip install redis kucoin-python
本地/服务器安装Redis服务,启动后默认端口6379即可使用。
修改后的WebSocket接收代码
import asyncio import redis from kucoin.client import WsToken from kucoin.ws_client import KucoinWsClient # 初始化Redis连接,全局复用 r = redis.Redis(host='127.0.0.1', port=6379, db=0, decode_responses=True) # 存K线的Redis键名 KLINE_KEY = "kline:SLP-USDT:30min" async def kline_msg(msg): if msg["topic"] == "/market/candles:SLP-USDT_30min": kline_data = msg["data"] # K线第一个字段是时间戳,用有序集合存,score设为时间戳自动去重排序 timestamp = int(kline_data[0]) # 把K线转成字符串存储,也可以用JSON序列化更安全 r.zadd(KLINE_KEY, {str(kline_data): timestamp}) # 可选:只保留最近1000根K线,避免内存占用过高 r.zremrangebyrank(KLINE_KEY, 0, -1001) print(f"已保存K线,时间戳:{timestamp}") async def wsocket(): client = WsToken() ws_client = await KucoinWsClient.create(None, client, kline_msg, private=False) await ws_client.subscribe("/market/candles:SLP-USDT_30min") while True: await asyncio.sleep(60) if __name__ == "__main__": loop = asyncio.get_event_loop() loop.run_until_complete(wsocket())
其他脚本读取K线示例
不管是存数据库、算指标还是绘图的脚本,直接连Redis读即可:
import redis import json r = redis.Redis(host='127.0.0.1', port=6379, db=0, decode_responses=True) KLINE_KEY = "kline:SLP-USDT:30min" # 读取最近100根K线 recent_klines = r.zrange(KLINE_KEY, -100, -1) # 转成列表格式处理 for line in recent_klines: kline = eval(line) # 存的时候如果用JSON序列化这里替换为json.loads更安全 print(kline) # 后续直接做入库、指标计算即可
方案2:本地共享内存(无额外服务依赖,仅适配同机器Python进程)
如果不想安装Redis,可以用Python自带的multiprocessing.Manager创建共享队列,把WebSocket进程和其他业务进程放在同一个启动脚本里,共享队列里的数据:
示例代码
import asyncio from multiprocessing import Process, Manager from kucoin.client import WsToken from kucoin.ws_client import KucoinWsClient # WebSocket接收进程逻辑 async def ws_process(shared_queue): async def kline_msg(msg): if msg["topic"] == "/market/candles:SLP-USDT_30min": shared_queue.put(msg["data"]) client = WsToken() ws_client = await KucoinWsClient.create(None, client, kline_msg, private=False) await ws_client.subscribe("/market/candles:SLP-USDT_30min") while True: await asyncio.sleep(60) def ws_runner(shared_queue): asyncio.run(ws_process(shared_queue)) # 业务处理进程逻辑(比如入库、算指标) def business_process(shared_queue): while True: if not shared_queue.empty(): kline = shared_queue.get() print(f"拿到K线:{kline}") # 这里写你的入库、指标计算逻辑 if __name__ == "__main__": # 创建共享队列 manager = Manager() kline_queue = manager.Queue(maxsize=1000) # 启动两个进程 p1 = Process(target=ws_runner, args=(kline_queue,)) p2 = Process(target=business_process, args=(kline_queue,)) p1.start() p2.start() p1.join() p2.join()
后续操作建议
- 数据入库可以用SQLite(轻量)或者MySQL/PostgreSQL(存储量大),单独写一个消费脚本从Redis/共享队列读数据批量写入即可,不要在WebSocket回调里直接写库,避免阻塞消息接收
- 指标计算可以用TA-Lib库,直接读取K线数据计算即可
- 可视化可以用Matplotlib或者Echarts,从数据库/Redis读数据后渲染即可
内容的提问来源于stack exchange,提问作者Max Power
相关产品推荐
相关产品推荐

