如何用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()
代码说明
- 生成器异步迭代:用
generator_worker代替直接提交生成器,它会在Flask-Executor的线程中完整迭代生成器,把每个结果放入线程安全队列,完成后放入结束标记。 - 双路监听:主线程通过
concurrent.futures.wait设置短超时,同时检查普通任务的完成状态和队列中的生成值,两边的结果都会被及时输出,不会互相阻塞。 - 循环终止条件:确保所有普通任务完成、队列清空、生成器线程结束后才退出循环,避免遗漏结果。
注意事项
- 队列采用线程安全的
queue.Queue,无需额外处理锁的问题。 - 超时时间(示例中为0.1秒)可根据业务场景调整,平衡响应速度和CPU占用。
- 如果
external_api()是IO密集型操作,该方案能充分发挥异步线程的优势,不会阻塞主线程。
内容的提问来源于stack exchange,提问作者jefe23984
相关产品推荐
相关产品推荐

