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"]