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

多线程与多进程结合的任务调度优化方案咨询

多线程与多进程结合优化任务调度方案

问题背景

刚接触多线程和多进程,已理解基本概念,但不清楚如何合理结合二者(比如线程内用进程、进程内用线程)。尝试混用后反而导致效率下降,现在需要优化任务框架,确保所有任务在指定时间(比如每小时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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 10:01:21