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

Python中如何实现3个顺序函数的独立流水线式并发执行?

Python 实现流水线式并发处理

你要的是典型的流水线并发,核心是让每个处理阶段(A、B、C)独立运行,前一个阶段处理完数据就立刻传给下一个,不用等整个流程走完再处理下一个数据。用Python的queue模块配合threading就能轻松实现,下面给你一步步讲:

核心思路

  • 用三个线程分别负责函数A、B、C的执行
  • 用队列(queue.Queue)作为阶段之间的缓冲区:
    • 输入队列:存放待处理的X数据
    • A→B队列:A处理完的数据传给B
    • B→C队列:B处理完的数据传给C
  • 每个线程从自己的输入队列取数据,处理后放到输出队列,直到收到结束信号

完整代码示例

假设你的A、B、C是处理单个数据的函数,我们先模拟耗时操作,方便直观看到提速效果:

import threading
import queue
import time

# 模拟你的实际处理函数
def func_a(x):
    time.sleep(0.5)  # 模拟IO/计算耗时
    return f"A处理完: {x}"

def func_b(a_result):
    time.sleep(0.5)
    return f"B处理完: {a_result}"

def func_c(b_result):
    time.sleep(0.5)
    print(f"最终输出Y: {b_result}")
    return f"C处理完: {b_result}"

# 通用的阶段线程逻辑
def stage_worker(input_q, output_q, func):
    while True:
        item = input_q.get()
        # 收到None就终止线程
        if item is None:
            input_q.task_done()
            break
        # 处理当前数据
        result = func(item)
        # 传给下一个阶段(如果有输出队列)
        if output_q is not None:
            output_q.put(result)
        input_q.task_done()

def main():
    # 创建各阶段的传递队列
    input_queue = queue.Queue()
    a_to_b_queue = queue.Queue()
    b_to_c_queue = queue.Queue()

    # 启动三个阶段的线程(daemon=True让线程随主进程退出)
    threading.Thread(target=stage_worker, args=(input_queue, a_to_b_queue, func_a), daemon=True).start()
    threading.Thread(target=stage_worker, args=(a_to_b_queue, b_to_c_queue, func_b), daemon=True).start()
    threading.Thread(target=stage_worker, args=(b_to_c_queue, None, func_c), daemon=True).start()

    # 放入待处理的数据
    data_list = [1, 2, 3, 4, 5]
    for x in data_list:
        input_queue.put(x)

    # 等待所有数据处理完成
    input_queue.join()
    a_to_b_queue.join()
    b_to_c_queue.join()

    # 给每个线程发送终止信号
    input_queue.put(None)
    a_to_b_queue.put(None)
    b_to_c_queue.put(None)

if __name__ == "__main__":
    start_time = time.time()
    main()
    print(f"总耗时: {time.time() - start_time:.2f}秒")

代码解释

  • stage_worker是通用线程函数,负责从输入队列取数、调用处理函数、传递结果,收到None就停止运行
  • join()方法用来等待队列中所有任务都处理完毕
  • 串行处理5个数据的总耗时是5*(0.5+0.5+0.5)=7.5秒,用流水线后总耗时约0.5*3 + 0.5*(5-1)=3.5秒,提速效果明显

注意事项

  1. 场景适配:

    • 如果你的A/B/C是IO密集型(比如读写文件、网络请求),上面的多线程方案完全够用
    • 如果是CPU密集型(比如大量计算),Python的GIL会限制多线程效率,建议换成multiprocessing模块:把threading.Thread改成multiprocessing.Process,queue.Queue改成multiprocessing.Queue,逻辑基本一致
  2. 异常处理:实际使用时要给处理函数加异常捕获,避免单个数据处理失败导致整个线程挂掉

  3. 队列限流:如果某个阶段处理速度慢,队列可能堆积大量数据,可以给Queue设置maxsize参数,让生产者线程等待,避免内存占用过高

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 22:30:52