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

如何并行运行TradingView数据采集与K线形态检测Python脚本

解决Python多进程重复启动module1及并行任务执行问题

问题根源

用multiprocessing.Pool时重复启动module1,核心原因是Pool创建子进程会重新导入主模块,如果主模块未做入口判断,子进程会重复执行module1的启动逻辑。此外需要确保两个长期运行的任务独立并行,且形态触发时能执行后续操作。

分步解决方案

1. 修正主进程入口,阻止子进程重复执行

将main.py的启动逻辑放在if __name__ == '__main__':判断内,这是Python多进程的标准写法,避免子进程重复初始化:

import multiprocessing
from module1 import run_tick_collector
from module2 import run_pinbar_detector

def main():
    # 用Process创建独立守护进程,而非Pool(Pool更适合批量短任务)
    tick_process = multiprocessing.Process(target=run_tick_collector, daemon=True)
    detector_process = multiprocessing.Process(target=run_pinbar_detector, daemon=True)

    tick_process.start()
    detector_process.start()

    # 阻塞等待进程运行(按Ctrl+C可终止)
    tick_process.join()
    detector_process.join()

if __name__ == '__main__':
    main()

2. 实现终端清晰输出

给两个模块的输出加前缀标识,避免终端输出混乱:

  • module1输出示例:
def run_tick_collector():
    # WebSocket连接及数据存储逻辑...
    while True:
        tick_data = get_websocket_tick()
        save_to_sql(tick_data)
        print(f"[Tick采集] 新增数据: {tick_data}")
  • module2输出示例:
import time

def run_pinbar_detector():
    # 循环读取SQL数据、聚合周期、检测形态...
    while True:
        raw_data = read_from_sql()
        aggregated_data = aggregate_to_timeframe(raw_data, timeframe='1m')
        signal = detect_pinbar(aggregated_data)
        
        if signal:
            print(f"[Pinbar检测] 触发{signal}形态!")
            execute_post_task(signal)  # 形态满足时执行后续任务
        else:
            print(f"[Pinbar检测] 当前无信号")
        
        time.sleep(10)  # 按需调整轮询间隔

3. 形态触发后的后续任务处理

如果后续任务需要和主进程/其他模块交互,用multiprocessing.Queue传递信号:

  • main.py扩展:
import multiprocessing
from module1 import run_tick_collector
from module2 import run_pinbar_detector
from module3 import run_post_task

def main():
    signal_queue = multiprocessing.Queue()

    tick_process = multiprocessing.Process(target=run_tick_collector, daemon=True)
    detector_process = multiprocessing.Process(target=run_pinbar_detector, args=(signal_queue,), daemon=True)
    post_process = multiprocessing.Process(target=run_post_task, args=(signal_queue,), daemon=True)

    tick_process.start()
    detector_process.start()
    post_process.start()

    tick_process.join()
    detector_process.join()
    post_process.join()

if __name__ == '__main__':
    main()
  • module2传递信号:
def run_pinbar_detector(signal_queue):
    while True:
        # ...数据聚合和形态检测逻辑
        if signal:
            print(f"[Pinbar检测] 触发{signal}形态!")
            signal_queue.put(signal)
  • module3执行后续任务:
def run_post_task(signal_queue):
    while True:
        signal = signal_queue.get()
        print(f"[后续任务] 收到{signal}信号,执行操作...")
        # 此处写入具体逻辑,如下单、通知等

为什么不用multiprocessing.Pool?

Pool是为批量短任务设计的进程池管理工具,而你的场景是两个长期运行的守护进程,用multiprocessing.Process单独创建进程更灵活,能避免Pool的进程复用逻辑导致的意外重复执行问题。

内容的提问来源于stack exchange,提问作者lazybot

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 18:45:16