答案是生成器惰性求值与异步任务调度冲突导致状态共享问题,需通过转列表、islice分片或异步生成器解耦。核心在于每个任务须有独立数据源,避免对同一生成器多次调用next()引发状态错乱或提前耗尽。

这个问题本质是混淆了生成器的惰性求值特性与异步任务的调度时机,导致所有协程共享同一个生成器状态或误用 next() 触发了意外的全量消费。
生成器在循环中被重复调用 next 的陷阱
生成器对象是一次性可迭代对象,每次调用 next() 都推进其内部状态。若在 for 循环或并发任务创建过程中反复对同一生成器实例调用 next()(尤其未做隔离),就会提前耗尽它,后续任务拿到的可能是 StopIteration 或默认值,甚至因状态错乱而全部读到末尾值。
- 错误写法:把一个生成器对象传给多个
asyncio.create_task,每个 task 内部都调用next(gen) - 后果:第一个 task 拿到第 1 个值,第二个 task 拿到第 2 个值……最后所有 task 实际处理的是不同元素;但如果逻辑误写成“所有 task 都等循环结束再统一取值”,就可能全取到最后一个 yield 出来的值
- 关键点:生成器不是数据容器,而是执行状态机;
next()是有副作用的操作
正确解耦生成器与异步任务的三种方式
核心原则:每个异步任务应拥有独立的数据源或明确的索引边界,避免共享生成器实例。
-
方式一:转为列表或元组 —— 适合数据量小、可全量加载的场景
用list(gen)或[*gen]提前展开,再按索引分发给各 task:data = [*my_generator()],然后asyncio.create_task(process(data[i])) -
方式二:用 itertools.islice 分片 —— 适合大数据流、需分批处理
from itertools import islice,每个 task 处理一段:chunk = list(islice(gen, start, end)),确保不重叠、不遗漏 -
方式三:改用异步生成器 + async for —— 真正适配 asyncio 的流式处理
定义async def async_data_stream():,用yield返回 awaitable 对象,task 内用async for item in async_data_stream():,天然支持暂停/恢复和事件循环调度
检查是否已触发全量消费的快速验证法
运行时加一行诊断代码,就能暴露问题:
- 在生成器函数里每 yield 前打印序号:
print(f"[gen] yielding #{i}") - 在 task 启动时打印获取的值:
val = next(gen); print(f"[task] got {val}") - 如果发现所有 task 打印的序号连续递增(如 1→2→3→4…),说明是共享生成器;如果全打印同一个序号(如全是 10),说明生成器已被提前耗尽,只剩最后一个缓存值
典型修复示例:从错误到安全
原错误代码(所有 task 共享一个 gen):
gen = (x for x in range(5)) tasks = [asyncio.create_task(worker(next(gen))) for _ in range(3)] # 错!next(gen) 被调用3次,gen 已走到第3个值
修复后(每个 task 独立 slice):
gen = (x for x in range(5)) data = list(gen) # 展开为列表 tasks = [asyncio.create_task(worker(data[i])) for i in range(min(3, len(data)))]
或更省内存的切片写法:
from itertools import islice gen = (x for x in range(5)) chunks = [list(islice(gen, i, i+1)) for i in range(3)] # 每个 chunk 是长度为1的列表 tasks = [asyncio.create_task(worker(chunk[0])) for chunk in chunks if chunk]
不复杂但容易忽略











