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

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("所有任务执行完成")

关键细节说明

  1. 队列设计:
    • task_queue:预处理进程向其中放入带文件名标识的待API任务
    • result_queue:API进程向其中放入带文件名标识的响应结果
  2. 进程退出逻辑:预处理进程全部完成后,向task_queue放入None信号,API进程收到后主动退出,避免进程挂起
  3. 响应匹配:预处理进程通过文件名精准匹配自己的API响应,确保每个进程只处理自身任务的后续工作
  4. 进程安全:multiprocessing.Queue自带锁机制,无需额外处理多进程竞争问题

之前尝试未生效的可能原因

大概率是没处理好任务标识匹配和进程退出信号:

  • 未给任务添加唯一标识,导致预处理进程无法准确获取对应自己的API响应
  • 未正确发送进程结束信号,导致进程一直阻塞在队列读取操作

内容的提问来源于stack exchange,提问作者Mr Anonymous

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 23:12:23