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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 14:48:01