asyncio已运行事件循环添加协程失败原因及实现方案咨询
问题解答
1. 前两种方法失效的核心原因
run_until_complete仅会等待传入的指定协程执行完成就终止事件循环,不会主动等待协程内部动态提交的其他任务。你在get_items协程末尾调用asyncio.ensure_future提交重试协程时,初始的get_items协程已经执行完所有逻辑,run_until_complete检测到目标协程返回后直接关闭了事件循环,刚提交的重试协程还没被调度就被销毁,因此无报错也无执行效果。asyncio.run_coroutine_threadsafe是跨线程提交任务的专用接口,同线程内调用没有意义,且只要事件循环已经停止,无论哪种提交方式的任务都不会执行。- 附加代码错误提示:你当前代码中
asyncio.gather的参数return_exception应为return_exceptions(缺少后缀s),入口判断的__names__应为__name__,如果gather返回了异常对象,直接取r[1]会抛出索引错误。
2. 无需重启事件循环的实现方案
方案1:内置重试循环(最适合当前场景)
直接将重试逻辑封装在同一个协程内部,不需要动态提交任务,初始协程会自动处理所有重试直到无失败条目:
async def get_items(self, item_type, item_names): step_size = 500 pending_names = item_names.copy() # 只要有未处理的条目就一直循环 while pending_names: current_batch = pending_names[:step_size] pending_names = pending_names[step_size:] # 生成批次任务 tasks = [ self.get_item_fetching_function(item_type)(name) for name in current_batch ] batch_results = await asyncio.gather(*tasks, return_exceptions=True) success_data = [] failed_names = [] # 分离成功结果和失败条目 for name, res in zip(current_batch, batch_results): # 可根据实际需求调整失败判断条件 if isinstance(res, Exception) or res.empty: failed_names.append(name) continue success_data.append(res) # 合并成功结果 if success_data: self.result = pd.concat([self.result] + success_data) # 失败条目塞回待处理队列 pending_names.extend(failed_names)
调用时无需改动原有wrapper_function逻辑,run_until_complete会等待所有重试完成后再终止循环。
方案2:通用动态任务调度方案
如果确实需要动态提交任务的通用能力,可以通过任务计数+事件通知的方式控制循环生命周期:
class DataSource: def __init__(self): self.loop = asyncio.get_event_loop() self._pending_task_count = 0 self._all_done_event = asyncio.Event() # 其他原有初始化逻辑不变 async def _wrap_task(self, coro): """任务包装器,自动计数""" self._pending_task_count += 1 try: await coro finally: self._pending_task_count -= 1 if self._pending_task_count == 0: self._all_done_event.set() def wrapper_function(self, item_type): # 提交初始任务 self.loop.create_task( self._wrap_task(self.get_items(item_type, self.item_names_to_process)) ) # 等待所有任务完成再停止循环 self.loop.run_until_complete(self._all_done_event.wait()) # 在get_items末尾提交重试任务时,用_wrap_task包裹即可 async def get_items(self, item_type, item_names): # 原有处理逻辑不变 # ... # 提交重试任务 self.loop.create_task( self._wrap_task(self.get_items(item_type, self.unprocessed)) )
内容的提问来源于stack exchange,提问作者stucash
相关产品推荐
相关产品推荐

