本文目录导读:

日志聚合、时间同步、日志级别管理和集中存储,下面提供几种常用的解决方案:
基础方案:结构化日志 + 文件存储
import logging
import json
from datetime import datetime
import socket
class DistributedLogger:
def __init__(self, node_id, log_dir="./logs"):
self.node_id = node_id
self.log_dir = log_dir
# 创建logger
self.logger = logging.getLogger(f"distributed-{node_id}")
self.logger.setLevel(logging.DEBUG)
# 文件处理器
fh = logging.FileHandler(f"{log_dir}/sync_{node_id}_{datetime.now().strftime('%Y%m%d')}.log")
fh.setLevel(logging.DEBUG)
# JSON格式化
formatter = logging.Formatter(
'{"timestamp": "%(asctime)s", "node": "%(node)s", '
'"level": "%(levelname)s", "message": "%(message)s", '
'"extra": %(extra)s}'
)
fh.setFormatter(formatter)
self.logger.addHandler(fh)
def log_sync(self, level, message, **extra):
"""记录同步日志"""
extra_dict = {
"sync_id": extra.get("sync_id", ""),
"operation": extra.get("operation", ""),
"status": extra.get("status", ""),
"duration": extra.get("duration", 0)
}
self.logger.info(
message,
extra={
"node": self.node_id,
"extra": json.dumps(extra_dict)
}
)
# 使用示例
node_logger = DistributedLogger("node-001")
node_logger.log_sync("INFO", "同步开始",
sync_id="sync-001",
operation="data_replication",
status="started")
ELK Stack 方案(生产环境推荐)
import logging
import logstash
from elasticsearch import Elasticsearch
from datetime import datetime
class ELKDistributedLogger:
def __init__(self, node_id, logstash_host='localhost', logstash_port=5000):
self.node_id = node_id
self.logger = logging.getLogger(f'distributed-{node_id}')
self.logger.setLevel(logging.INFO)
# Logstash handler
logstash_handler = logstash.LogstashHandler(
logstash_host,
logstash_port,
version=1
)
self.logger.addHandler(logstash_handler)
def log_sync_event(self, event_type, data):
"""记录同步事件到ELK"""
log_data = {
'node_id': self.node_id,
'event_type': event_type,
'timestamp': datetime.utcnow().isoformat(),
'data': data,
'service': 'distributed-sync'
}
if event_type == 'sync_start':
self.logger.info('同步开始', extra=log_data)
elif event_type == 'sync_complete':
self.logger.info('同步完成', extra=log_data)
elif event_type == 'sync_error':
self.logger.error('同步错误', extra=log_data)
# 使用示例
elk_logger = ELKDistributedLogger('node-001')
elk_logger.log_sync_event('sync_start', {
'sync_id': 'sync-001',
'source': 'db-primary',
'target': 'db-replica'
})
使用 Redis 做日志聚合
import redis
import json
import threading
from datetime import datetime
class RedisLogAggregator:
def __init__(self, node_id, redis_host='localhost', redis_port=6379):
self.node_id = node_id
self.redis_client = redis.Redis(
host=redis_host,
port=redis_port,
decode_responses=True
)
self.log_queue = f"sync_logs:{node_id}"
# 启动后台发送线程
self.buffer = []
self.buffer_lock = threading.Lock()
self.send_thread = threading.Thread(target=self._batch_send)
self.send_thread.daemon = True
self.send_thread.start()
def log_sync(self, level, message, **kwargs):
"""记录同步日志到Redis"""
log_entry = {
'timestamp': datetime.utcnow().isoformat(),
'node_id': self.node_id,
'level': level,
'message': message,
'extra': kwargs,
'host': socket.gethostname()
}
with self.buffer_lock:
self.buffer.append(log_entry)
def _batch_send(self):
"""批量发送日志到Redis"""
while True:
time.sleep(2) # 每2秒批量发送一次
with self.buffer_lock:
if self.buffer:
# 使用pipeline批量插入
pipe = self.redis_client.pipeline()
for entry in self.buffer:
pipe.lpush(self.log_queue, json.dumps(entry))
pipe.execute()
self.buffer.clear()
def get_recent_logs(self, count=100):
"""获取最近的日志"""
logs = self.redis_client.lrange(self.log_queue, 0, count-1)
return [json.loads(log) for log in logs]
# 使用示例
redis_logger = RedisLogAggregator('node-001')
redis_logger.log_sync('INFO', '数据同步开始',
sync_id='sync-001',
source_table='orders',
target_table='orders_backup')
分布式追踪方案(使用OpenTelemetry)
from opentelemetry import trace
from opentelemetry.exporter.otlp.proto.grpc.exporter import OTLPSpanExporter
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
import time
class DistributedTracingLogger:
def __init__(self, service_name, otlp_endpoint="localhost:4317"):
# 配置OpenTelemetry
trace.set_tracer_provider(TracerProvider())
tracer_provider = trace.get_tracer_provider()
otlp_exporter = OTLPSpanExporter(endpoint=otlp_endpoint)
span_processor = BatchSpanProcessor(otlp_exporter)
tracer_provider.add_span_processor(span_processor)
self.tracer = trace.get_tracer(__name__)
self.service_name = service_name
def trace_sync_operation(self, operation_name, func):
"""使用追踪记录同步操作"""
with self.tracer.start_as_current_span(
operation_name,
attributes={
"service.name": self.service_name,
"operation.type": "sync",
"sync.id": str(uuid.uuid4())
}
) as span:
try:
start_time = time.time()
result = func()
duration = time.time() - start_time
span.set_attribute("sync.duration", duration)
span.set_attribute("sync.status", "success")
return result
except Exception as e:
span.set_attribute("sync.status", "error")
span.set_attribute("sync.error", str(e))
span.set_status(trace.Status(trace.StatusCode.ERROR))
raise
# 使用示例
tracing_logger = DistributedTracingLogger("data-sync-service")
def sync_data():
# 实际的同步逻辑
time.sleep(1)
return "同步成功"
result = tracing_logger.trace_sync_operation("sync_orders_data", sync_data)
完整的生产级方案
import logging
import json
import time
import uuid
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime
import requests
class ProductionDistributedLogger:
def __init__(self, config):
self.config = config
self.node_id = f"{config['hostname']}-{config['process_id']}"
# 初始化多个日志后端
self.loggers = []
self._init_local_logger()
self._init_remote_logger()
# 日志缓冲区
self.buffer = []
self.buffer_size = config.get('buffer_size', 100)
self.flush_interval = config.get('flush_interval', 5)
# 启动自动刷新
self._start_auto_flush()
def _init_local_logger(self):
"""初始化本地文件日志"""
local_logger = logging.getLogger(f"local-{self.node_id}")
handler = logging.FileHandler(
f"{self.config['log_dir']}/sync_{datetime.now().strftime('%Y%m%d')}.log"
)
handler.setFormatter(logging.Formatter(
'%(asctime)s - %(levelname)s - %(message)s'
))
local_logger.addHandler(handler)
self.loggers.append(local_logger)
def _init_remote_logger(self):
"""初始化远程日志收集器"""
self.remote_endpoint = self.config.get('remote_log_endpoint')
def _start_auto_flush(self):
"""启动自动刷新线程"""
def flush_worker():
while True:
time.sleep(self.flush_interval)
self.flush()
import threading
thread = threading.Thread(target=flush_worker, daemon=True)
thread.start()
def log_sync_event(self, event_type, metadata=None):
"""记录同步事件"""
log_entry = {
'id': str(uuid.uuid4()),
'timestamp': datetime.utcnow().isoformat(),
'node_id': self.node_id,
'event_type': event_type,
'metadata': metadata or {},
'sync_id': metadata.get('sync_id') if metadata else None,
'duration_ms': metadata.get('duration_ms') if metadata else None
}
# 添加到缓冲区
self.buffer.append(log_entry)
# 如果缓冲区满,立即刷新
if len(self.buffer) >= self.buffer_size:
self.flush()
# 本地记录
for logger in self.loggers:
if isinstance(logger, logging.Logger):
logger.info(json.dumps(log_entry))
def flush(self):
"""刷新缓冲区到远程"""
if not self.buffer:
return
with ThreadPoolExecutor(max_workers=1) as executor:
future = executor.submit(self._send_to_remote, self.buffer.copy())
self.buffer.clear()
def _send_to_remote(self, logs):
"""发送日志到远程服务"""
if not self.remote_endpoint:
return
try:
response = requests.post(
self.remote_endpoint,
json={'logs': logs},
timeout=5
)
response.raise_for_status()
except Exception as e:
logging.error(f"发送日志到远程失败: {e}")
# 重新添加失败的日志到缓冲区
with threading.Lock():
self.buffer.extend(logs)
# 生产环境使用示例
config = {
'hostname': 'server-01',
'process_id': 12345,
'log_dir': '/var/log/app/sync',
'remote_log_endpoint': 'http://log-aggregator:8080/api/logs',
'buffer_size': 200,
'flush_interval': 10
}
logger = ProductionDistributedLogger(config)
# 记录同步操作
sync_id = str(uuid.uuid4())
start_time = time.time()
logger.log_sync_event('sync_start', {
'sync_id': sync_id,
'source': 'postgresql://primary',
'target': 'postgresql://replica',
'tables': ['orders', 'users']
})
# 执行同步...
time.sleep(2)
duration_ms = (time.time() - start_time) * 1000
logger.log_sync_event('sync_complete', {
'sync_id': sync_id,
'duration_ms': duration_ms,
'records_synced': 1500,
'status': 'success'
})
最佳实践建议
- 使用结构化日志格式:JSON格式便于解析和分析
- 包含关键信息:时间戳、节点ID、操作ID、状态等
- 实现日志聚合:使用ELK、Splunk等工具集中管理
- 设置合理的日志级别:避免日志量过大
- 实现缓冲机制:减少I/O操作
- 考虑日志安全:加密敏感信息
- 监控日志系统:确保日志收集不成为性能瓶颈
选择哪种方案取决于你的具体需求、基础设施和预算,对于小型项目,文件日志+ELK就足够了;对于大型分布式系统,建议使用OpenTelemetry或专门的日志平台。