Python多进程实现两组并行任务:10预处理+4受限API调用协同
多进程优化方案:API请求限制下的任务并行实现
需求明确
- 10个进程负责预处理任务:读取S3对象、预处理,生成待API请求数据
- 4个进程负责API调用任务:严格控制同时4个请求,接收预处理数据调用API,返回响应结果
- 预处理进程拿到API响应后,完成后续S3存储等工作
- 核心要求:两组进程并行执行,避免API请求超限导致无响应
可行实现方案
用multiprocessing.Queue实现进程间的任务传递和结果回传,这是进程安全的通信方式,完全适配你的需求。关键是给每个任务加唯一标识(比如文件名),让预处理进程能准确匹配到自己的API响应。
完整代码实现
import multiprocessing import numpy as np # 替换为实际API调用逻辑 def call_external_api(data): # 此处编写真实的API请求代码 return {"filename": data["filename"], "api_response": "API处理结果"} # API进程逻辑:监听任务队列,处理后写入结果队列 def api_worker(task_queue, result_queue): while True: task = task_queue.get() # 收到结束信号时退出进程 if task is None: break # 调用外部API response = call_external_api(task) # 将响应写入结果队列 result_queue.put(response) # 预处理+后续处理进程逻辑 def preprocess_worker(file_list, task_queue, result_queue): for filename in file_list: # 1. 从S3读取对象(替换为实际逻辑) s3_raw_data = f"读取S3中的{filename}数据" # 2. 数据预处理(替换为实际逻辑) preprocessed_data = {"filename": filename, "processed_data": s3_raw_data} # 3. 将预处理结果送入API任务队列 task_queue.put(preprocessed_data) # 4. 等待匹配自己任务的API响应 while True: result = result_queue.get() if result["filename"] == filename: # 5. 后续工作:将结果存储至S3(替换为实际逻辑) print(f"{filename} 完成后续存储,API响应:{result['api_response']}") break if __name__ == "__main__": myList = ["abc.xml", "def.xml", "ghi.xml", "jkl.xml"] # 实际为上千个对象 num_preprocess_workers = 10 num_api_workers = 4 # 将文件列表拆分给10个预处理进程 split_file_lists = np.array_split(myList, num_preprocess_workers) # 创建进程安全的通信队列 task_queue = multiprocessing.Queue() result_queue = multiprocessing.Queue() # 启动4个API处理进程 api_processes = [] for _ in range(num_api_workers): p = multiprocessing.Process(target=api_worker, args=(task_queue, result_queue)) p.start() api_processes.append(p) # 启动10个预处理进程 preprocess_processes = [] for sub_list in split_file_lists: # 将numpy数组转换为普通列表 sub_list = sub_list.tolist() p = multiprocessing.Process(target=preprocess_worker, args=(sub_list, task_queue, result_queue)) p.start() preprocess_processes.append(p) # 等待所有预处理进程完成任务提交 for p in preprocess_processes: p.join() # 向API进程发送结束信号 for _ in range(num_api_workers): task_queue.put(None) # 等待所有API进程处理完毕并退出 for p in api_processes: p.join() print("所有任务执行完成")
关键细节说明
- 队列设计:
task_queue:预处理进程向其中放入带文件名标识的待API任务result_queue:API进程向其中放入带文件名标识的响应结果
- 进程退出逻辑:预处理进程全部完成后,向
task_queue放入None信号,API进程收到后主动退出,避免进程挂起 - 响应匹配:预处理进程通过文件名精准匹配自己的API响应,确保每个进程只处理自身任务的后续工作
- 进程安全:
multiprocessing.Queue自带锁机制,无需额外处理多进程竞争问题
之前尝试未生效的可能原因
大概率是没处理好任务标识匹配和进程退出信号:
- 未给任务添加唯一标识,导致预处理进程无法准确获取对应自己的API响应
- 未正确发送进程结束信号,导致进程一直阻塞在队列读取操作
内容的提问来源于stack exchange,提问作者Mr Anonymous
相关产品推荐
相关产品推荐

