import atexit import os import sys from dotenv import load_dotenv from fluent import sender from loguru import logger load_dotenv() # 日志格式 STDOUT_FORMAT = ( "{time:YYYY-MM-DD HH:mm:ss.SSS} | " "{level: <8} | " "{name}:{function}:{line} - " "{message}" ) FILE_FORMAT = ( "{time:YYYY-MM-DD HH:mm:ss.SSS} | " "{level: <8} | " "{name}:{function}:{line} - {message}" ) fluent_sender = sender.FluentSender( tag=os.getenv("FLUENTD_TOPIC"), host=os.getenv("FLUENTD_HOST"), port=int(os.getenv("FLUENTD_PORT", 24224)), buffer_max_size=8 * 1024 * 1024, timeout=3.0, retry_timeout=60 ) def log_to_fluent(message): try: record = message.record # 构建结构化日志数据 log_data = { 'topic': os.getenv("FLUENTD_TOPIC"), 'timestamp': record['time'].timestamp(), 'level': record['level'].name.lower(), 'message': record['message'], 'source': f"{record['file'].path}:{record['line']}", 'module': record['module'], 'function': record['function'], 'process_id': record['process'].id, 'thread_id': record['thread'].id, **record['extra'] } if not fluent_sender.emit(os.getenv("FLUENTD_TOPIC"), log_data): print(f"Fluentd 发送失败: {fluent_sender.last_error}") except Exception as e: print(f"日志处理异常: {str(e)}") # 移除默认处理器 logger.remove() # 添加控制台处理器 logger.add( sink=sys.stdout, level="INFO", format=STDOUT_FORMAT, colorize=True, backtrace=True, # 显示完整异常堆栈 diagnose=True, # 显示详细异常信息 ) logger.add( log_to_fluent, level="INFO", # 处理 INFO 及以上级别 format="{message}", # 原始消息(实际使用结构化数据) backtrace=True, # 启用堆栈回溯 diagnose=True # 显示诊断信息 ) atexit.register(fluent_sender.close) # 导出配置好的logger __all__ = ["logger"]