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

如何使用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的实现方式过于繁琐,可考虑以下方案:

  1. 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())
    
  2. 第三方任务队列库:
    如Celery(分布式任务)、RQ(简单本地任务队列),适合需要持久化任务、分布式处理的场景,可天然支持流式任务的提交和结果回调。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 03:37:24