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

FastAPI自定义异步日志器报错:RuntimeError: 事件循环已运行

FastAPI自定义异步日志Handler问题解决

我有一个FastAPI应用,想要实现自定义日志器,通过异步请求将日志发送到指定主机,但遇到了两个问题:

第一个错误

RuntimeError: This event loop is already running

对应的初始代码:

import asyncio
import functools
import logging
import requests

class HttpHandler(logging.Handler):
    """
    Logging Handler for fluent.
    Sends HTTP POST requests 
    """
    def __init__(self, tag: str, host: str, port: int):
        self.tag = tag
        self.host = host
        self.port = port
        self.loop = asyncio.get_event_loop()
        self.loop.run_forever()
    
        logging.Handler.__init__(self)
    
    async def emit(self, record: logging.LogRecord):
        payload = {} # some data
        future = self.loop.run_in_executor(None, functools.partial(requests.post, data={
            "url": self.host,
            "json": payload
        }))
        response = await future

FastAPI中使用方式:

@app.on_event("startup")
async def startup():
    # replace fastapi handlers
    logger.handlers = ServiceLogger.get_logger_instance().handlers # special service that gives me custom logger

@app.post("/some_url")
async def foo(body=Depends(get_body)):
    logger.info("TEXT")

更新后的警告

尝试改用aiohttp修改代码后:

async def emit(self, record: logging.LogRecord):
        payload = {
            ...
            "message": record.message # log message
        }
        await aiohttp.post(self.host, json=payload)

出现如下警告:

RuntimeWarning: coroutine 'HttpHandler.emit' was never awaited
self.emit(record)
RuntimeWarning: Enable tracemalloc to get the object allocation traceback


问题根源分析

  1. 初始代码中self.loop.run_forever()会直接启动事件循环,但FastAPI已经在运行自己的事件循环,因此触发RuntimeError。
  2. logging.Handler的emit方法默认是同步设计,若定义为async def emit,日志系统调用时不会自动await,导致协程未执行的警告。

正确实现方式

方案1:同步emit + 异步任务提交(推荐)

利用FastAPI已有的事件循环,在同步emit方法中提交异步任务,无需自行管理循环:

import asyncio
import logging
import aiohttp

class AsyncHttpHandler(logging.Handler):
    def __init__(self, tag: str, host: str, port: int):
        super().__init__()
        self.tag = tag
        self.host = f"{host}:{port}"  # 拼接完整日志接收地址
        self.session = None

    # 在应用启动时初始化aiohttp会话,复用连接提升性能
    async def setup(self):
        self.session = aiohttp.ClientSession()

    # 在应用关闭时关闭会话,释放资源
    async def cleanup(self):
        if self.session:
            await self.session.close()

    def emit(self, record: logging.LogRecord):
        # 将异步日志发送任务提交到当前事件循环
        asyncio.create_task(self._async_emit(record))

    async def _async_emit(self, record: logging.LogRecord):
        if not self.session:
            return  # 会话未初始化则跳过发送
        payload = {
            "tag": self.tag,
            "message": record.getMessage(),
            "level": record.levelname,
            "timestamp": record.created
        }
        try:
            async with self.session.post(self.host, json=payload):
                pass  # 根据需求添加响应处理逻辑
        except Exception:
            # 捕获发送异常,避免日志系统崩溃影响主应用
            self.handleError(record)

在FastAPI中正确初始化

@app.on_event("startup")
async def startup():
    custom_handler = AsyncHttpHandler(tag="my_fastapi_app", host="http://your-log-server", port=8080)
    await custom_handler.setup()
    # 替换默认日志处理器
    logger = logging.getLogger("uvicorn")
    logger.handlers = [custom_handler]

@app.on_event("shutdown")
async def shutdown():
    for handler in logger.handlers:
        if isinstance(handler, AsyncHttpHandler):
            await handler.cleanup()

@app.post("/some_url")
async def foo(body=Depends(get_body)):
    logger.info("处理请求:/some_url")

方案2:线程池执行同步请求(简单但效率较低)

若不想使用异步HTTP客户端,可通过线程池包装同步requests请求,避免阻塞事件循环:

import logging
import requests
from concurrent.futures import ThreadPoolExecutor

class SyncHttpHandler(logging.Handler):
    def __init__(self, tag: str, host: str, port: int):
        super().__init__()
        self.tag = tag
        self.host = f"{host}:{port}"
        self.executor = ThreadPoolExecutor(max_workers=5)

    def emit(self, record: logging.LogRecord):
        payload = {
            "tag": self.tag,
            "message": record.getMessage()
        }
        # 提交到线程池执行,不阻塞主线程
        self.executor.submit(self._send_log, payload, record)

    def _send_log(self, payload, record):
        try:
            requests.post(self.host, json=payload)
        except Exception:
            self.handleError(record)

关键注意事项

  • 不要在Handler的__init__中调用loop.run_forever(),FastAPI已维护自身事件循环。
  • logging模块为同步设计,自定义emit必须是同步方法;若需异步逻辑,用asyncio.create_task提交到事件循环,或用线程池包装。
  • 使用aiohttp时务必复用ClientSession,避免每次请求创建新连接。
  • 必须捕获日志发送异常,防止日志系统错误影响主应用运行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 08:52:44