You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.23 15:36:06