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就退出。
关键注意事项
- 分块下载是核心:不管用哪种方案,下载逻辑一定要拆成分块处理,这样暂停时不会丢失太多进度,也能更快响应控制信号。
- 避免长时间阻塞:子进程的循环里不要有长时间的阻塞操作(比如不带stream的requests.get()),否则会导致控制信号响应延迟。
- 断点续传的补充:如果需要断点续传,要把每个文件的下载进度存在共享存储里(比如Manager字典或本地文件),恢复时通过HTTP的Range请求续传。
内容的提问来源于stack exchange,提问作者K. Macieja

