如何向asyncio.as_completed的任务列表追加任务?求解
问题分析与解决方案
问题原因
asyncio.as_completed只会在初始化时遍历传入的可迭代对象(比如你的tasks列表),生成对应任务的完成迭代器。后续你往tasks中追加的新任务,不会被这个已生成的迭代器纳入监测范围——所以循环处理完第一个任务就直接结束,而新追加的任务仍处于pending状态,当事件循环终止时就会抛出Task was destroyed but it is pending!错误。
另外你的create_task_form_url函数返回值类型标注有误,应该改为asyncio.Future而非[Future]。
正确实现方式
方案一:用集合+asyncio.wait动态管理任务
通过asyncio.wait配合FIRST_COMPLETED参数,每次只处理最先完成的任务,同时动态将新任务加入待处理集合:
import asyncio from asyncio import get_running_loop async def async_get_results(url, session): # 模拟URL请求逻辑 await asyncio.sleep(1) return f"Result from {url}" async def create_task_form_url(aio_session, form_url: str) -> asyncio.Future: loop = get_running_loop() task = loop.create_task(async_get_results(form_url, aio_session)) return task async def dynamic_task_handler(aio_session, initial_url): pending = set() # 添加初始任务 initial_task = await create_task_form_url(aio_session, initial_url) pending.add(initial_task) # 这里可以根据需求设置终止条件,比如处理N个任务后停止 task_limit = 5 processed_count = 0 while pending and processed_count < task_limit: # 等待第一个完成的任务 done, pending = await asyncio.wait(pending, return_when=asyncio.FIRST_COMPLETED) for task in done: processed_count +=1 # 获取任务结果 result = await task print(result) # 根据业务逻辑添加新任务 if processed_count < task_limit: new_task = await create_task_form_url(aio_session, initial_url) pending.add(new_task) print(f"当前待处理任务数: {len(pending)}") # 运行示例 async def main(): async with asyncio.ClientSession() as session: await dynamic_task_handler(session, "https://example.com") if __name__ == "__main__": asyncio.run(main())
方案二:用队列实现任务动态调度
如果需要更灵活的任务管理(比如多消费者),可以用asyncio.Queue来承载任务,通过消费者协程动态处理并添加新任务:
import asyncio from asyncio import get_running_loop async def async_get_results(url, session): await asyncio.sleep(1) return f"Result from {url}" async def create_task_form_url(aio_session, form_url: str) -> asyncio.Future: loop = get_running_loop() task = loop.create_task(async_get_results(form_url, aio_session)) return task async def queue_based_handler(aio_session, initial_url): task_queue = asyncio.Queue() # 放入初始任务 initial_task = await create_task_form_url(aio_session, initial_url) await task_queue.put(initial_task) processed_count = 0 task_limit = 5 async def consumer(): nonlocal processed_count while processed_count < task_limit: task = await task_queue.get() try: processed_count +=1 result = await task print(result) # 添加新任务到队列 if processed_count < task_limit: new_task = await create_task_form_url(aio_session, initial_url) await task_queue.put(new_task) print(f"当前队列任务数: {task_queue.qsize()}") finally: task_queue.task_done() # 启动消费者协程 consumer_task = asyncio.create_task(consumer()) # 等待队列所有任务处理完成 await task_queue.join() # 取消消费者协程 consumer_task.cancel() try: await consumer_task except asyncio.CancelledError: pass async def main(): async with asyncio.ClientSession() as session: await queue_based_handler(session, "https://example.com") if __name__ == "__main__": asyncio.run(main())
关键注意事项
- 必须设置明确的任务终止条件(比如任务数量限制、无新任务可添加等),否则会无限循环创建任务
- 避免直接修改
as_completed初始化时传入的可迭代对象,因为它不会感知到后续的修改
内容的提问来源于stack exchange,提问作者jake wong
相关产品推荐
相关产品推荐

