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秒,提速效果明显
注意事项
场景适配:
- 如果你的A/B/C是IO密集型(比如读写文件、网络请求),上面的多线程方案完全够用
- 如果是CPU密集型(比如大量计算),Python的GIL会限制多线程效率,建议换成
multiprocessing模块:把threading.Thread改成multiprocessing.Process,queue.Queue改成multiprocessing.Queue,逻辑基本一致
异常处理:实际使用时要给处理函数加异常捕获,避免单个数据处理失败导致整个线程挂掉
队列限流:如果某个阶段处理速度慢,队列可能堆积大量数据,可以给
Queue设置maxsize参数,让生产者线程等待,避免内存占用过高
内容的提问来源于stack exchange,提问作者Adi
相关产品推荐
相关产品推荐

