Python异步生成器已实现__anext__却报错未实现的问题排查
问题原因
你错误地将异步迭代器的__aiter__方法定义成了异步函数(添加了async关键字)。根据Python异步迭代器的规范:
__aiter__必须是同步方法,直接返回实现了__anext__的迭代器对象(这里就是你的SubEventStream实例)- 只有
__anext__需要定义为异步方法,用于返回下一个异步迭代值
因为你给__aiter__加了async,调用它时会返回一个协程对象,而不是你的SubEventStream实例。async for尝试从这个协程对象获取__anext__方法,自然找不到,就抛出了TypeError。
修改后的代码
import asyncio class SubEventStream(): def __init__(self) -> None: self.queue = asyncio.Queue() def __aiter__(self): # 去掉async关键字 return self async def __anext__(self): return await self.pop() async def append(self, request): return await self.queue.put(request) async def pop(self): r = await self.queue.get() self.queue.task_done() return r def create_append_tasks(ls, q): return [ asyncio.create_task(q.append(i)) for i in ls ] async def append_tasks(q): tasks = create_append_tasks(('a', 'b', 'c', 'd', 'e'), q) return await asyncio.gather(*tasks) async def run(): q = SubEventStream() await append_tasks(q) async for v in q: print(v) asyncio.run(run())
补充说明
修改后执行代码,会依次输出a、b、c、d、e,但当前代码会在输出完所有元素后一直阻塞——因为队列空了之后__anext__会持续等待新元素。如果需要迭代完现有元素后自动停止,可以在__anext__中判断队列状态并抛出StopAsyncIteration异常:
async def __anext__(self): if self.queue.empty(): raise StopAsyncIteration return await self.pop()
这个调整会让迭代仅处理队列中已有的元素,不会等待后续添加的内容,可根据你的实际需求选择是否添加。
内容的提问来源于stack exchange,提问作者Jamie Marshall
相关产品推荐
相关产品推荐

