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

Python多进程:主进程与多子进程间的异步通信实现需求

主进程与多子进程异步通信及下载暂停实现方案

我之前做过类似的批量下载工具,用multiprocessing.JoinableQueue的时候确实踩过没法暂停的坑——子进程调用get()后会一直阻塞等待任务,根本没法响应主进程的暂停信号。要解决这个问题,核心是给子进程设计一个可中断的任务循环,让它们能同时监听任务和控制指令。下面几个方案都是我实际验证过的,分享给你:

方案1:双队列分离任务与控制指令

给每个子进程绑定两个队列:一个用来存下载链接的任务队列,另一个用来传递PAUSE/RESUME/EXIT这类控制指令的控制队列。子进程用multiprocessing.connection.wait()同时监听两个队列的消息,这样既不会一直卡在任务队列的get()上,又能及时响应主进程的控制信号。

示例代码如下:

import multiprocessing
import time
import requests

def worker(task_queue, control_queue):
    paused = False
    while True:
        # 同时监听两个队列,短超时保证及时响应信号
        ready_queues = multiprocessing.connection.wait([task_queue, control_queue], timeout=0.1)
        
        # 先处理控制指令
        if control_queue in ready_queues:
            cmd = control_queue.get()
            if cmd == 'PAUSE':
                paused = True
                print(f"Worker {multiprocessing.current_process().pid} 已暂停")
            elif cmd == 'RESUME':
                paused = False
                print(f"Worker {multiprocessing.current_process().pid} 已恢复")
            elif cmd == 'EXIT':
                print(f"Worker {multiprocessing.current_process().pid} 正在退出")
                break
        
        # 非暂停状态下处理下载任务
        if not paused and task_queue in ready_queues:
            url = task_queue.get()
            try:
                print(f"开始下载: {url}")
                # 分块下载,每块检查一次暂停状态
                response = requests.get(url, stream=True)
                filename = url.split('/')[-1]
                with open(filename, 'wb') as f:
                    for chunk in response.iter_content(chunk_size=1024):
                        if paused:
                            print(f"{url} 下载已暂停")
                            break  # 可以记录进度,后续实现断点续传
                        f.write(chunk)
                task_queue.task_done()
            except Exception as e:
                print(f"下载失败 {url}: {str(e)}")
                task_queue.task_done()

if __name__ == '__main__':
    task_queue = multiprocessing.JoinableQueue()
    control_queue = multiprocessing.Queue()
    
    # 启动3个工作进程
    worker_count = 3
    workers = [
        multiprocessing.Process(target=worker, args=(task_queue, control_queue))
        for _ in range(worker_count)
    ]
    for w in workers:
        w.start()
    
    # 添加测试下载任务
    test_urls = [
        "https://example.com/file1",
        "https://example.com/file2",
        "https://example.com/file3"
    ]
    for url in test_urls:
        task_queue.put(url)
    
    # 模拟主进程发送控制指令
    time.sleep(2)
    print("主进程发送暂停指令")
    control_queue.put('PAUSE')
    
    time.sleep(3)
    print("主进程发送恢复指令")
    control_queue.put('RESUME')
    
    # 等待所有任务完成,发送退出指令
    task_queue.join()
    for _ in range(worker_count):
        control_queue.put('EXIT')
    for w in workers:
        w.join()

这个方案的优势是逻辑清晰,任务和控制指令完全分离,子进程能精准响应每个信号。如果需要断点续传,还可以在暂停时把当前下载进度存在multiprocessing.Manager的共享字典里,恢复时从断点继续。

方案2:用Event实现全局暂停/恢复

如果不需要精细控制单个进程,只需要全局暂停所有下载,用multiprocessing.Event会更简洁。主进程通过设置/清除Event来控制所有子进程的状态,子进程在下载循环中定期检查Event的状态。

示例代码:

import multiprocessing
import time
import requests
import os

def worker(task_queue, pause_event):
    while True:
        # 非阻塞获取任务,避免卡在get()上无法响应Event
        try:
            url = task_queue.get(block=False)
            # 用None作为退出标记
            if url is None:
                break
        except multiprocessing.queues.Empty:
            time.sleep(0.1)
            continue
        
        print(f"开始下载: {url}")
        try:
            filename = url.split('/')[-1]
            downloaded_size = 0
            # 支持断点续传:读取已下载的文件大小
            try:
                downloaded_size = os.path.getsize(filename)
                headers = {'Range': f'bytes={downloaded_size}-'}
                response = requests.get(url, stream=True, headers=headers)
            except FileNotFoundError:
                response = requests.get(url, stream=True)
            
            with open(filename, 'ab') as f:
                for chunk in response.iter_content(chunk_size=1024):
                    # 暂停时进入循环等待,直到Event被清除
                    while pause_event.is_set():
                        print(f"{filename} 下载暂停中...")
                        time.sleep(1)
                    f.write(chunk)
            task_queue.task_done()
            print(f"下载完成: {url}")
        except Exception as e:
            print(f"下载失败 {url}: {str(e)}")
            task_queue.task_done()

if __name__ == '__main__':
    task_queue = multiprocessing.JoinableQueue()
    pause_event = multiprocessing.Event()
    
    worker_count = 3
    workers = [
        multiprocessing.Process(target=worker, args=(task_queue, pause_event))
        for _ in range(worker_count)
    ]
    for w in workers:
        w.start()
    
    # 添加测试任务
    test_urls = [
        "https://example.com/file1",
        "https://example.com/file2",
        "https://example.com/file3"
    ]
    for url in test_urls:
        task_queue.put(url)
    
    # 模拟全局暂停/恢复
    time.sleep(2)
    print("全局暂停所有下载")
    pause_event.set()
    
    time.sleep(3)
    print("全局恢复所有下载")
    pause_event.clear()
    
    # 等待任务完成,发送退出信号
    task_queue.join()
    for _ in range(worker_count):
        task_queue.put(None)
    for w in workers:
        w.join()

这个方案代码更简洁,适合全局控制的场景。需要注意的是,子进程获取任务时必须用非阻塞的get(block=False),否则队列空的时候子进程会一直阻塞,没法响应Event的变化。

方案3:共享字典实现精细化进程控制

如果需要单独控制某个进程,或者实时查看每个进程的下载进度,可以用multiprocessing.Manager创建一个共享字典,主进程通过修改字典里的进程状态来控制子进程,子进程定期检查自己的状态。

比如共享字典可以这样设计:

from multiprocessing import Manager
manager = Manager()
worker_status = manager.dict()
# 子进程启动后把自己的ID和初始状态存入字典
worker_status[multiprocessing.current_process().pid] = 'RUNNING'

子进程在循环里每次都检查worker_status[自己的PID],如果状态是PAUSED就暂停下载,RESUMED就继续,EXIT就退出。

关键注意事项

  1. 分块下载是核心:不管用哪种方案,下载逻辑一定要拆成分块处理,这样暂停时不会丢失太多进度,也能更快响应控制信号。
  2. 避免长时间阻塞:子进程的循环里不要有长时间的阻塞操作(比如不带stream的requests.get()),否则会导致控制信号响应延迟。
  3. 断点续传的补充:如果需要断点续传,要把每个文件的下载进度存在共享存储里(比如Manager字典或本地文件),恢复时通过HTTP的Range请求续传。

内容的提问来源于stack exchange,提问作者K. Macieja

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 12:27:55