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

实现带异步emit的logging.Handler?解决进程退出警告问题

解决Asyncio日志处理器进程退出时的协程未等待警告问题

你遇到的这个警告本质是进程退出时,还有通过asyncio.run_coroutine_threadsafe提交的日志写入协程没来得及执行完成,Python解释器就终止了,所以抛出了这个提示。咱们一步步来优化你的实现:

1. 先修正原代码的小问题

原代码里的_async_emit方法漏掉了self参数,这会直接导致运行报错,先补上这个细节:

async def _async_emit(self, message):
    await my_async_write_function(message)

2. 优雅处理待完成的日志任务(推荐方案)

要彻底解决退出警告,核心是在进程关闭前,等待所有已提交的日志写入协程执行完毕。我们可以维护一个任务集合来跟踪所有提交的协程,然后在处理器关闭时统一等待:

import asyncio
import logging
from typing import Set

class AsyncEmitLogHandler(logging.Handler):
    def __init__(self):
        self.loop = asyncio.get_running_loop()
        self.pending_tasks: Set[asyncio.Future] = set()
        super().__init__()

    def emit(self, record):
        self.format(record)
        coro = self._async_emit(record.message)
        future = asyncio.run_coroutine_threadsafe(coro, loop=self.loop)
        # 将线程安全的Future转为asyncio原生Future,方便跟踪
        task = asyncio.wrap_future(future)
        self.pending_tasks.add(task)
        # 任务完成后自动从集合中移除,避免内存泄漏
        task.add_done_callback(self.pending_tasks.discard)

    async def _async_emit(self, message):
        await my_async_write_function(message)

    def close(self):
        # 关闭处理器时,等待所有待完成的日志任务
        if self.pending_tasks:
            # 因为close是同步方法,用run_until_complete等待所有任务完成
            self.loop.run_until_complete(asyncio.gather(*self.pending_tasks))
        super().close()

关键细节说明:

  • 用pending_tasks集合跟踪所有提交的日志任务,任务完成后通过add_done_callback自动移除,避免内存泄漏
  • 重写close方法,在处理器关闭阶段统一等待所有待执行的协程,从根源解决未等待的警告
  • asyncio.wrap_future将run_coroutine_threadsafe返回的concurrent.futures.Future转换为asyncio原生Future,方便统一管理

3. 临时方案:抑制特定警告

如果因为业务限制无法修改关闭逻辑,也可以通过Python的警告过滤器来抑制这个特定警告,不过这属于治标不治本的临时方案:

import warnings

# 精准抑制"coroutine was never awaited"警告
warnings.filterwarnings(
    "ignore",
    category=RuntimeWarning,
    message="coroutine '.*_async_emit' was never awaited"
)

4. 额外优化:处理日志写入异常

原代码没有处理emit过程中的异常,建议加上异常捕获,避免日志写入失败影响主进程:

def emit(self, record):
    try:
        self.format(record)
        coro = self._async_emit(record.message)
        future = asyncio.run_coroutine_threadsafe(coro, loop=self.loop)
        task = asyncio.wrap_future(future)
        self.pending_tasks.add(task)
        task.add_done_callback(self.pending_tasks.discard)
    except Exception:
        # 用logging内置的错误处理逻辑处理异常
        self.handleError(record)

内容的提问来源于stack exchange,提问作者Russell Owen

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 18:50:29