多线程与多进程结合的任务调度优化方案咨询
多线程与多进程结合优化任务调度方案
问题背景
刚接触多线程和多进程,已理解基本概念,但不清楚如何合理结合二者(比如线程内用进程、进程内用线程)。尝试混用后反而导致效率下降,现在需要优化任务框架,确保所有任务在指定时间(比如每小时00:15)前完成,同时处理好任务的并发/并行和依赖关系。
任务框架代码
class JobManager(): def __init__(self): # 初始化属性:线程、成员变量、套接字等 def job1(self): # 快速计算:无IO、无循环 # 最优耗时1分钟 def job2(self): # 重型计算:涉及大量内存数据,适合多进程并行(循环密集) # 最优耗时2分钟 def job3(self): # 通过REST API获取数据并更新内存数据 # 耗时3分钟(受API速率限制无法优化) def job4(self): # 开启WebSocket监听消息并计算,需后台持续运行,更新内存数据 def job5(self): # 开启另一个WebSocket监听不同来源消息并计算,需后台持续运行,更新内存数据 def job6(self): # 快速计算,依赖job1、job2完成 def run_threaded(job_func): job_thread = threading.Thread(target=job_func) job_thread.start() def main(): job_manager = JobManager() # 启动后台线程(比如job4、job5) job_manager.some_thread1.start() job_manager.some_thread2.start() scheduler = schedule.Scheduler() scheduler.every().hour.at("00:15").do(run_threaded, job_manager.job6) scheduler.every().hour.at("00:14").do(run_threaded, job_manager.job1) scheduler.every().hour.at("00:13").do(run_threaded, job_manager.job2) scheduler.every().hour.at("00:12").do(run_threaded, job_manager.job3) while True: scheduler.run_pending() time.sleep(1) if __name__ == "__main__": main()
最优结合方案
1. 按任务特性匹配并发模型
- IO密集/后台长任务(job3、job4、job5):用多线程
- job4、job5是WebSocket长连接,属于IO阻塞型任务,放在独立线程后台运行,既不占用主线程资源,又能方便共享内存更新数据。
- job3是API请求,受速率限制,大部分时间在等待响应,用线程执行可避免阻塞其他任务。
- CPU密集型任务(job2):用多进程
- job2是循环密集的重型计算,Python的GIL会限制多线程并行效率,改用多进程可利用多核CPU拆分任务并行执行,压缩耗时。
- 轻量计算任务(job1、job6):直接在主线程或单线程执行,避免频繁创建销毁线程/进程的上下文切换开销。
2. 处理任务依赖关系
job6依赖job1、job2的结果,不能靠定时触发,需用同步机制确保前置任务完成后再执行:
- 用
concurrent.futures的wait方法监听job1、job2的完成状态,待两者都结束后再启动job6。 - 示例:用线程池执行job1,进程池执行job2,通过
wait阻塞等待结果,再触发job6。
3. 优化调度逻辑
原定时设置不合理(比如job3需3分钟,00:12触发刚好00:15完成,但易与其他任务冲突超时),调整为:
- 提前触发长耗时任务:job3在00:09触发(3分钟后00:12完成),job2在00:11触发(2分钟后00:13完成),job1在00:13触发(1分钟后00:14完成),00:14执行job6,确保00:15前全部完成。
- 用线程/进程池管理任务,复用资源,减少创建销毁的开销。
4. 内存数据同步注意事项
- 多线程共享
JobManager内存数据时,对共享数据的读写加threading.Lock,避免竞态条件。 - 多进程无法直接共享内存,job2需读取主进程数据时,用
multiprocessing.Queue或共享内存对象传递;任务结果要更新主进程内存时,同样通过IPC机制传递。
调整后的核心代码示例
import threading import multiprocessing from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor, wait import schedule import time class JobManager(): def __init__(self): self.data_lock = threading.Lock() self.shared_data = {} # 初始化WebSocket守护线程(随主进程退出) self.ws_thread4 = threading.Thread(target=self.job4, daemon=True) self.ws_thread5 = threading.Thread(target=self.job5, daemon=True) def job1(self): # 快速计算逻辑 time.sleep(60) with self.data_lock: self.shared_data["job1_result"] = "result1" def job2(self, input_data): # 重型计算:拆分子任务并行执行 def sub_task(task_data): time.sleep(30) return task_data * 2 with ProcessPoolExecutor(max_workers=multiprocessing.cpu_count()) as executor: futures = [executor.submit(sub_task, i) for i in input_data] results = [f.result() for f in futures] return results def job3(self): # API请求逻辑 time.sleep(180) with self.data_lock: self.shared_data["job3_data"] = "api_data" def job4(self): # WebSocket监听逻辑(持续运行) while True: time.sleep(5) with self.data_lock: self.shared_data["ws4_update"] = time.time() def job5(self): # WebSocket监听逻辑(持续运行) while True: time.sleep(5) with self.data_lock: self.shared_data["ws5_update"] = time.time() def job6(self): # 依赖job1、job2结果的计算 with self.data_lock: job1_res = self.shared_data.get("job1_result") job2_res = self.shared_data.get("job2_result") if job1_res and job2_res: time.sleep(10) print("job6 completed") def run_hourly_jobs(job_manager): # 先执行job3(IO密集,线程) with ThreadPoolExecutor(max_workers=1) as executor: job3_future = executor.submit(job_manager.job3) job3_future.result() # 执行job2(CPU密集,进程) with ProcessPoolExecutor(max_workers=multiprocessing.cpu_count()) as executor: input_data = job_manager.shared_data.get("job3_data", [1,2,3,4]) job2_future = executor.submit(job_manager.job2, input_data) # 执行job1(轻量计算,线程) with ThreadPoolExecutor(max_workers=1) as executor: job1_future = executor.submit(job_manager.job1) # 等待job1、job2完成,更新结果到共享内存 wait([job1_future, job2_future]) with job_manager.data_lock: job_manager.shared_data["job2_result"] = job2_future.result() # 执行job6 job_manager.job6() def main(): job_manager = JobManager() # 启动后台WebSocket线程 job_manager.ws_thread4.start() job_manager.ws_thread5.start() # 每小时00:09触发批量任务,确保00:15前完成 scheduler = schedule.Scheduler() scheduler.every().hour.at("00:09").do(run_hourly_jobs, job_manager) while True: scheduler.run_pending() time.sleep(1) if __name__ == "__main__": main()
关键总结
- 按需选择并发模型:CPU密集用多进程,IO密集/后台长任务用多线程,轻量计算直接执行。
- 明确依赖关系:用同步机制而非定时来确保前置任务完成。
- 保障数据安全:多线程共享数据加锁,多进程用IPC传递数据。
- 复用资源:用线程/进程池减少创建销毁的开销。
内容的提问来源于stack exchange,提问作者NicFit_88
相关产品推荐
相关产品推荐

