Python脚本如何记录分布式同步执行日志

wen python案例 33

本文目录导读:

Python脚本如何记录分布式同步执行日志

  1. 基础方案:结构化日志 + 文件存储
  2. ELK Stack 方案(生产环境推荐)
  3. 使用 Redis 做日志聚合
  4. 分布式追踪方案(使用OpenTelemetry)
  5. 完整的生产级方案
  6. 最佳实践建议

日志聚合、时间同步、日志级别管理和集中存储,下面提供几种常用的解决方案:

基础方案:结构化日志 + 文件存储

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'
})

最佳实践建议

  1. 使用结构化日志格式:JSON格式便于解析和分析
  2. 包含关键信息:时间戳、节点ID、操作ID、状态等
  3. 实现日志聚合:使用ELK、Splunk等工具集中管理
  4. 设置合理的日志级别:避免日志量过大
  5. 实现缓冲机制:减少I/O操作
  6. 考虑日志安全:加密敏感信息
  7. 监控日志系统:确保日志收集不成为性能瓶颈

选择哪种方案取决于你的具体需求、基础设施和预算,对于小型项目,文件日志+ELK就足够了;对于大型分布式系统,建议使用OpenTelemetry或专门的日志平台。

抱歉,评论功能暂时关闭!