使用multiprocessing运行3个仅tokens不同的tick任务,是否建议切换为multithreading?
结论
非常建议你切换为multithreading实现,切换后内存占用会显著降低,运行速度不会下降,甚至会比多进程版本更稳定高效。
原因说明
- 你的任务属于典型的IO密集型场景:核心操作分别是WebSocket接收行情tick、执行数据库UPDATE语句,全程大部分时间是CPU等待IO响应,Python的GIL全局解释器锁在IO阻塞阶段会自动释放,多线程完全可以实现并行处理,不会出现CPU密集型场景下的GIL争抢问题。
- 多进程会为每个子进程生成独立的Python解释器副本、独立内存空间,你当前3个进程相当于重复存储了3份导入的依赖库、公共变量、函数逻辑,内存开销极大。多线程共享同一进程的内存空间,仅保留一份公共资源,3个线程的内存开销仅为原多进程方案的1/3甚至更低。
- 多进程的启动、销毁、上下文切换开销远高于线程,切换为多线程后,你完全不需要担心性能下降,反而会因为减少了进程间调度开销,tick接收和入库的延迟会更低。
代码优化建议
你不需要写3个重复的tick_A/tick_B/tick_C函数,直接参数化即可,改写后的参考代码如下:
from database_function import * from kiteconnect import KiteTicker import pandas as pd from datetime import datetime, timedelta import schedule import time from threading import Thread def tick_task(offset): # credentials代码保持不变 # 按参数动态拉取对应范围的tokens if offset == 0: tokens = [x[0] for x in db_fetchquery("SELECT zerodha FROM script ORDER BY id ASC LIMIT 50")] else: tokens = [x[0] for x in db_fetchquery(f"SELECT zerodha FROM script ORDER BY id ASC OFFSET ({offset}) ROWS FETCH NEXT (50) ROWS ONLY")] ##### 等待到8:59启动的逻辑保持不变 ########### t = datetime.today() future = datetime(t.year,t.month,t.day,8,59) if ((future-t).total_seconds()) < 0: future = datetime(t.year,t.month,t.day,t.hour,t.minute,(t.second+2)) time.sleep((future-t).total_seconds()) ##### 等待逻辑结束 ########### def on_ticks(ws, ticks): global ltp ltp = ticks[0]["last_price"] for tick in ticks: print(f"{tick['instrument_token']}") db_runquery(f'UPDATE SCRIPT SET ltp = {tick["last_price"]} WHERE zerodha = {tick["instrument_token"]}') def on_connect(ws, response): ws.subscribe(tokens) ws.set_mode(ws.MODE_LTP,tokens) kws.on_ticks = on_ticks kws.on_connect = on_connect kws.connect(threaded=True) ##### 15:32停止的逻辑保持不变 ####### end_time = datetime(t.year,t.month,t.day,15,32) while True: schedule.run_pending() if datetime.now() > end_time: break ##### 停止逻辑结束 ####### if __name__ == '__main__': def runInParallel(*args): threads = [] for offset in args: t = Thread(target=tick_task, args=(offset,)) t.start() threads.append(t) for t in threads: t.join() runInParallel(0, 50, 100)
注意事项
需要确认你使用的db_runquery对应的数据库连接是否为线程安全:如果每个线程会独立创建数据库连接则无需额外处理,如果是全局共享的单数据库连接,需要给更新操作加线程锁,避免并发写入出现异常。
内容的提问来源于stack exchange,提问作者teneji
相关产品推荐
相关产品推荐

