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

如何用ThreadPoolExecutor在生成器执行期间穿插普通方法调用?

解决Flask-Executor中生成器与普通函数结果穿插输出的问题

问题原因

你当前的代码逻辑存在关键问题:当generator()被提交到线程池后,线程仅返回生成器对象本身,并未开始迭代。调用future.result()拿到生成器后,会在主线程中同步遍历所有yield值,此时as_completed循环被完全阻塞,无法处理已完成的func1/func2任务,导致普通函数结果必须等生成器全部输出后才能返回。

解决方案

要实现生成器输出与普通函数结果的穿插,需将生成器的迭代过程放到异步线程中,通过线程安全队列传递每个生成值,同时主线程同时监听普通任务的完成状态和队列的新值。

修改后的代码示例

import concurrent.futures
import queue
from app import executor

# 线程安全队列,用于传递生成器输出
gen_queue = queue.Queue()
# 标记生成器执行完成的特殊值
GEN_FINISHED = object()

def generator_worker():
    """在异步线程中迭代生成器,并将结果放入队列"""
    for val in external_api():
        gen_queue.put(val)
    # 生成器执行完毕后放入结束标记
    gen_queue.put(GEN_FINISHED)

def func1():
    return 1

def func2():
    return 2

def stream():
    # 提交生成器的工作线程
    gen_task = executor.submit(generator_worker)
    # 提交普通函数任务
    task1 = executor.submit(func1)
    task2 = executor.submit(func2)

    # 跟踪待处理的普通任务
    pending_tasks = {task1, task2}

    while pending_tasks or not gen_queue.empty() or gen_task.running():
        # 短超时轮询普通任务的完成状态
        done, pending = concurrent.futures.wait(
            pending_tasks,
            timeout=0.1,
            return_when=concurrent.futures.FIRST_COMPLETED
        )
        # 输出已完成的普通任务结果
        for task in done:
            yield task.result()
        pending_tasks = pending

        # 输出队列中生成器的新值
        while not gen_queue.empty():
            value = gen_queue.get()
            if value is GEN_FINISHED:
                break
            yield value

    # 确保生成器任务完全结束
    gen_task.result()

代码说明

  1. 生成器异步迭代:用generator_worker代替直接提交生成器,它会在Flask-Executor的线程中完整迭代生成器,把每个结果放入线程安全队列,完成后放入结束标记。
  2. 双路监听:主线程通过concurrent.futures.wait设置短超时,同时检查普通任务的完成状态和队列中的生成值,两边的结果都会被及时输出,不会互相阻塞。
  3. 循环终止条件:确保所有普通任务完成、队列清空、生成器线程结束后才退出循环,避免遗漏结果。

注意事项

  • 队列采用线程安全的queue.Queue,无需额外处理锁的问题。
  • 超时时间(示例中为0.1秒)可根据业务场景调整,平衡响应速度和CPU占用。
  • 如果external_api()是IO密集型操作,该方案能充分发挥异步线程的优势,不会阻塞主线程。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 21:17:52