如何使用Python的concurrent.futures处理流式任务?
流式任务场景下的concurrent.futures实现方案
常规用法回顾
concurrent.futures是Python中用于异步执行任务的标准库,常规批量任务处理示例如下:
from time import sleep from random import random from concurrent.futures import ThreadPoolExecutor from concurrent.futures import wait # 自定义任务:随机休眠后输出名称 def task(name): sleep(random()) print(name) # 启动线程池并批量提交任务 with ThreadPoolExecutor(10) as executor: futures = [executor.submit(task, i) for i in range(10)] wait(futures) print('All tasks are done!')
常规模式的局限是:任务需一次性批量提交,之后只能通过wait()等待全部完成,或as_completed()遍历已完成任务,无法应对任务逐个/批量流式到达、需持续提交并按完成顺序处理的场景。
基于concurrent.futures的流式任务实现
通过维护线程安全的任务集合,配合独立线程监听任务完成状态,可实现流式提交+按完成顺序处理的需求。核心思路是:
- 主线程负责接收流式任务并提交至线程池,将返回的
Future对象加入线程安全集合 - 单独启动一个线程,通过
as_completed()持续监听集合中的任务,一旦有任务完成立即处理
实现代码
from time import sleep from random import random from concurrent.futures import ThreadPoolExecutor, as_completed import threading # 自定义任务:随机休眠后返回结果 def task(name): sleep(random()) return f"Task {name} completed" def process_completed(futures_set, lock): """持续监听并处理已完成的任务""" while True: # 线程安全地获取当前待监听的任务列表 with lock: current_futures = list(futures_set) if not current_futures: sleep(0.1) continue # 遍历已完成的任务 for future in as_completed(current_futures): try: result = future.result() print(result) # 线程安全地移除已处理的任务 with lock: futures_set.remove(future) except Exception as e: print(f"Task failed: {str(e)}") # 每次处理一个任务后重新获取任务列表,确保能监听新提交的任务 break def main(): futures = set() futures_lock = threading.Lock() executor = ThreadPoolExecutor(max_workers=10) # 启动任务处理线程(守护线程,随主线程退出) processing_thread = threading.Thread( target=process_completed, args=(futures, futures_lock), daemon=True ) processing_thread.start() # 模拟流式任务到达:分批、间隔提交任务 for i in range(15): sleep(random() * 0.5) # 模拟任务间隔到达 future = executor.submit(task, i) with futures_lock: futures.add(future) print(f"Submitted task {i}") # 等待所有任务处理完成 while True: with futures_lock: if not futures: break sleep(0.1) executor.shutdown() if __name__ == "__main__": main()
关键细节
- 使用线程锁
threading.Lock()保护futures集合的读写操作,避免多线程冲突 - 处理线程每次完成一个任务后重新获取任务列表,确保能监听到新提交的任务
- 采用守护线程处理任务结果,避免主线程退出时残留线程
替代解决方案
如果concurrent.futures的实现方式过于繁琐,可考虑以下方案:
- asyncio 结合 ThreadPoolExecutor:
利用asyncio的异步特性处理任务提交和结果监听,更适合IO密集型流式任务,示例如下:import asyncio from time import sleep from random import random from concurrent.futures import ThreadPoolExecutor def task(name): sleep(random()) return f"Task {name} completed" async def main(): executor = ThreadPoolExecutor(max_workers=10) loop = asyncio.get_event_loop() tasks = [] # 模拟流式提交任务 for i in range(15): await asyncio.sleep(random() * 0.5) # 将同步任务包装为异步任务 async_task = loop.run_in_executor(executor, task, i) tasks.append(async_task) print(f"Submitted task {i}") # 按完成顺序处理结果 for completed in asyncio.as_completed(tasks): result = await completed print(result) executor.shutdown() asyncio.run(main()) - 第三方任务队列库:
如Celery(分布式任务)、RQ(简单本地任务队列),适合需要持久化任务、分布式处理的场景,可天然支持流式任务的提交和结果回调。
内容的提问来源于stack exchange,提问作者SupernoobBran
相关产品推荐
相关产品推荐

