如何并行运行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
相关产品推荐
相关产品推荐

