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
相关产品推荐
相关产品推荐

