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

Python日志阻塞问题:如何实现Datadog非阻塞日志传输?

解决Python异步环境下Datadog日志阻塞问题

问题原因

Python标准logging库的默认处理器是同步的,你的DatadogCustomLogHandler发送日志到Datadog时执行的是同步IO操作。在异步任务(比如http_get)中调用logger.info时,会阻塞整个事件循环——所有并发任务都要等待同步IO完成才能继续,这就是耗时从0.97秒暴涨到8秒的核心原因。

解决方案

方案1:用QueueHandler+QueueListener异步化日志处理(标准库方案,改动最小)

利用Python标准库的QueueHandler和QueueListener,将日志发送操作转移到后台线程执行,完全不阻塞异步事件循环。

修改日志初始化函数:

import logging
from logging.handlers import QueueHandler, QueueListener
import queue
import os
import socket

os.environ["DD_API_KEY"] = 'XXXXX'
os.environ["DD_SITE"] = 'XXXXX'
host_name = socket.gethostname()

def init_datadog_logging(service_name: str = None, env_name: str = None, min_log_level: int = logging.INFO):
    tags = f'service: {service_name}, host: {host_name}, environment: {env_name}'
    root_logger = logging.getLogger()
    root_logger.setLevel(min_log_level)

    # 移除默认的同步处理器,避免重复输出
    for handler in root_logger.handlers[:]:
        root_logger.removeHandler(handler)

    if env_name in ('UAT', 'Production'):
        # 创建无界队列存放日志
        log_queue = queue.Queue(-1)
        # 初始化原有的Datadog自定义处理器
        datadog_handler = DatadogCustomLogHandler(tags=tags, service=service_name, level=min_log_level)
        # 启动QueueListener,在后台线程处理日志发送
        queue_listener = QueueListener(log_queue, datadog_handler)
        queue_listener.start()
        # 给根日志器添加QueueHandler,将日志转发到队列
        root_logger.addHandler(QueueHandler(log_queue))

原理:所有日志输出会先被QueueHandler存入队列,QueueListener在独立线程中从队列取出日志,调用DatadogCustomLogHandler完成发送,异步任务的事件循环不会被阻塞。

方案2:自定义异步Datadog日志处理器(更贴合异步环境)

自己实现基于aiohttp的异步日志处理器,直接在事件循环中异步发送日志,避免线程切换开销。

import asyncio
import aiohttp
import logging
from logging import Handler, LogRecord
import os
import socket

os.environ["DD_API_KEY"] = 'XXXXX'
os.environ["DD_SITE"] = 'XXXXX'
host_name = socket.gethostname()

class AsyncDatadogLogHandler(Handler):
    def __init__(self, tags: str, service: str, level=logging.NOTSET):
        super().__init__(level)
        self.api_key = os.environ["DD_API_KEY"]
        self.site = os.environ["DD_SITE"]
        self.tags = [tag.strip() for tag in tags.split(",")]
        self.service = service
        self.session = aiohttp.ClientSession()
        self.log_url = f"https://http-intake.logs.{self.site}/v1/input"
        # 启动批量发送任务(减少HTTP请求次数)
        asyncio.create_task(self._batch_send())
        self.batch_queue = asyncio.Queue(maxsize=50)
        self.flush_interval = 3  # 每3秒批量发送一次

    def emit(self, record: LogRecord):
        # 将日志任务提交到异步队列,不阻塞当前线程
        log_entry = {
            "message": self.format(record),
            "service": self.service,
            "tags": self.tags,
            "host": host_name,
            "level": record.levelname.lower()
        }
        asyncio.create_task(self.batch_queue.put(log_entry))

    async def _batch_send(self):
        while True:
            batch = []
            # 攒够批量大小或到刷新时间就发送
            try:
                for _ in range(self.batch_queue.maxsize):
                    entry = await asyncio.wait_for(self.batch_queue.get(), timeout=self.flush_interval)
                    batch.append(entry)
            except asyncio.TimeoutError:
                pass
            
            if batch:
                headers = {
                    "DD-API-KEY": self.api_key,
                    "Content-Type": "application/json"
                }
                try:
                    async with self.session.post(self.log_url, json=batch, headers=headers):
                        pass
                except Exception as e:
                    # 发送失败则将日志放回队列
                    for entry in batch:
                        await self.batch_queue.put(entry)
                    self.handleError(None)

    async def close(self):
        # 程序退出前发送剩余日志
        remaining = []
        while not self.batch_queue.empty():
            remaining.append(await self.batch_queue.get())
        if remaining:
            headers = {
                "DD-API-KEY": self.api_key,
                "Content-Type": "application/json"
            }
            await self.session.post(self.log_url, json=remaining, headers=headers)
        await self.session.close()
        await super().close()

# 修改日志初始化函数
def init_datadog_logging(service_name: str = None, env_name: str = None, min_log_level: int = logging.INFO):
    tags = f'service: {service_name}, host: {host_name}, environment: {env_name}'
    root_logger = logging.getLogger()
    root_logger.setLevel(min_log_level)

    for handler in root_logger.handlers[:]:
        root_logger.removeHandler(handler)

    if env_name in ('UAT', 'Production'):
        async_handler = AsyncDatadogLogHandler(tags=tags, service=service_name, level=min_log_level)
        root_logger.addHandler(async_handler)

额外优化:复用aiohttp ClientSession

你的http_get函数每次都创建新的ClientSession,这会带来不必要的开销。改成复用main函数中的session:

logger = logging.getLogger(__name__)

async def http_get(session: aiohttp.ClientSession, url: str) -> JSONObject:
    logger.info(f'Getting data from {url}')
    async with session.get(url) as response:
        return await response.json()

# main_async中修改任务创建逻辑
async def main_async() -> None:
    time_before = time.perf_counter()
    async with aiohttp.ClientSession() as session:
        tasks = [asyncio.create_task(http_get(session, BASE_URL + date)) for date in dates]
        responses = await asyncio.gather(*tasks)

    logger.info(f'Total time taken: {time.perf_counter() - time_before}')

这样能减少TCP连接的创建和销毁开销,进一步提升异步任务的执行效率。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 05:07:02