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

wen python案例 25

本文目录导读:

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

  1. 基础日志记录方案
  2. 带进度和性能统计的高级方案
  3. 结构化JSON日志(适合日志分析)
  4. 快速启动模板
  5. 推荐实践

为Python分片同步脚本记录执行日志,推荐使用内置的logging模块,以下是几种实用的日志记录方案:

基础日志记录方案

import logging
import time
from datetime import datetime
# 配置日志
logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(levelname)s - %(message)s',
    handlers=[
        logging.FileHandler(f'sync_{datetime.now().strftime("%Y%m%d")}.log'),
        logging.StreamHandler()  # 同时输出到控制台
    ]
)
logger = logging.getLogger(__name__)
def sync_shard(shard_id, data):
    """模拟分片同步函数"""
    try:
        logger.info(f"开始同步分片 {shard_id}")
        logger.debug(f"分片 {shard_id} 数据量: {len(data)} 条")
        # 模拟同步过程
        time.sleep(2)
        # 模拟可能的错误
        if shard_id == 3:
            raise Exception("连接超时")
        logger.info(f"分片 {shard_id} 同步完成")
        return True
    except Exception as e:
        logger.error(f"分片 {shard_id} 同步失败: {str(e)}", exc_info=True)
        return False
# 使用示例
shards = [
    {"id": 1, "data": list(range(100))},
    {"id": 2, "data": list(range(200))},
    {"id": 3, "data": list(range(150))},
    {"id": 4, "data": list(range(300))}
]
for shard in shards:
    sync_shard(shard["id"], shard["data"])

带进度和性能统计的高级方案

import logging
import time
import json
from datetime import datetime
from logging.handlers import RotatingFileHandler
class SyncLogger:
    def __init__(self, log_dir="logs"):
        import os
        os.makedirs(log_dir, exist_ok=True)
        # 配置日志记录器
        self.logger = logging.getLogger('sync_logger')
        self.logger.setLevel(logging.DEBUG)
        # 文件日志 - 带轮转
        file_handler = RotatingFileHandler(
            f'{log_dir}/sync_{datetime.now().strftime("%Y%m%d")}.log',
            maxBytes=10*1024*1024,  # 10MB
            backupCount=5,
            encoding='utf-8'
        )
        file_handler.setLevel(logging.DEBUG)
        # 错误日志 - 单独文件
        error_handler = RotatingFileHandler(
            f'{log_dir}/sync_error_{datetime.now().strftime("%Y%m%d")}.log',
            maxBytes=10*1024*1024,
            backupCount=5,
            encoding='utf-8'
        )
        error_handler.setLevel(logging.ERROR)
        # 控制台输出
        console_handler = logging.StreamHandler()
        console_handler.setLevel(logging.INFO)
        # 格式化
        formatter = logging.Formatter(
            '%(asctime)s | %(levelname)-8s | %(name)s | %(message)s',
            datefmt='%Y-%m-%d %H:%M:%S'
        )
        file_handler.setFormatter(formatter)
        error_handler.setFormatter(formatter)
        console_handler.setFormatter(formatter)
        self.logger.addHandler(file_handler)
        self.logger.addHandler(error_handler)
        self.logger.addHandler(console_handler)
        # 统计信息
        self.stats = {
            "start_time": datetime.now(),
            "total_shards": 0,
            "successful": 0,
            "failed": 0,
            "total_records": 0,
            "failed_shards": []
        }
    def log_sync_start(self, shard_id, records_count):
        """记录同步开始"""
        self.stats["total_shards"] += 1
        self.stats["total_records"] += records_count
        self.logger.info(f"▶ 开始同步分片 [{shard_id}] | 记录数: {records_count}")
    def log_sync_progress(self, shard_id, current, total, extra_info=""):
        """记录同步进度"""
        percentage = (current / total) * 100
        self.logger.debug(
            f"分片 [{shard_id}] 进度: {current}/{total} ({percentage:.1f}%) {extra_info}"
        )
    def log_sync_success(self, shard_id, duration, records_synced):
        """记录同步成功"""
        self.stats["successful"] += 1
        speed = records_synced / duration if duration > 0 else 0
        self.logger.info(
            f"✓ 分片 [{shard_id}] 同步成功 | 耗时: {duration:.2f}s | "
            f"记录数: {records_synced} | 速度: {speed:.0f}条/秒"
        )
    def log_sync_failure(self, shard_id, error, duration):
        """记录同步失败"""
        self.stats["failed"] += 1
        self.stats["failed_shards"].append({
            "shard_id": shard_id,
            "error": str(error),
            "duration": duration
        })
        self.logger.error(
            f"✗ 分片 [{shard_id}] 同步失败 | 耗时: {duration:.2f}s | "
            f"错误: {str(error)}",
            exc_info=True
        )
    def log_summary(self):
        """输出同步总结"""
        duration = (datetime.now() - self.stats["start_time"]).total_seconds()
        summary = (
            f"\n{'='*50}\n"
            f"📊 同步总结报告\n"
            f"{'='*50}\n"
            f"开始时间: {self.stats['start_time']}\n"
            f"结束时间: {datetime.now()}\n"
            f"总耗时: {duration:.2f}s\n"
            f"总分片数: {self.stats['total_shards']}\n"
            f"成功: {self.stats['successful']}\n"
            f"失败: {self.stats['failed']}\n"
            f"总记录数: {self.stats['total_records']}\n"
            f"成功率: {(self.stats['successful']/self.stats['total_shards'])*100:.1f}%\n"
        )
        if self.stats["failed_shards"]:
            summary += f"失败分片: {self.stats['failed_shards']}\n"
        self.logger.info(summary)
        # 保存JSON格式的统计
        with open(f'logs/sync_stats_{datetime.now().strftime("%Y%m%d_%H%M%S")}.json', 'w') as f:
            json.dump(self.stats, f, default=str, indent=2)
# 使用示例
def sync_shard_with_logging(shard, sync_logger):
    """带日志的分片同步"""
    shard_id = shard["id"]
    data = shard["data"]
    sync_logger.log_sync_start(shard_id, len(data))
    start_time = time.time()
    try:
        # 模拟同步过程
        for i, record in enumerate(data):
            # 模拟处理每条记录
            time.sleep(0.01)
            # 记录进度(每10%记录一次)
            if i % max(1, len(data)//10) == 0:
                sync_logger.log_sync_progress(shard_id, i+1, len(data))
        duration = time.time() - start_time
        sync_logger.log_sync_success(shard_id, duration, len(data))
        return True
    except Exception as e:
        duration = time.time() - start_time
        sync_logger.log_sync_failure(shard_id, e, duration)
        return False
# 主程序
if __name__ == "__main__":
    sync_logger = SyncLogger()
    shards = [
        {"id": "shard_001", "data": list(range(100))},
        {"id": "shard_002", "data": list(range(200))},
        {"id": "shard_003", "data": list(range(50))},
        {"id": "shard_004", "data": list(range(300))}
    ]
    for shard in shards:
        sync_shard_with_logging(shard, sync_logger)
    sync_logger.log_summary()

结构化JSON日志(适合日志分析)

import logging
import json
from datetime import datetime
import socket
class StructuredLogger:
    """结构化JSON日志"""
    def __init__(self, service_name="shard-sync", log_file="structured_sync.log"):
        self.service_name = service_name
        self.hostname = socket.gethostname()
        logging.basicConfig(
            level=logging.INFO,
            format='%(message)s',  # 只输出消息部分
            handlers=[
                logging.FileHandler(log_file),
                logging.StreamHandler()
            ]
        )
        self.logger = logging.getLogger('structured')
    def _create_log_entry(self, level, message, **kwargs):
        """创建结构化日志条目"""
        entry = {
            "timestamp": datetime.now().isoformat(),
            "service": self.service_name,
            "host": self.hostname,
            "level": level,
            "message": message,
            **kwargs
        }
        return json.dumps(entry, ensure_ascii=False)
    def info(self, message, **kwargs):
        entry = self._create_log_entry("INFO", message, **kwargs)
        self.logger.info(entry)
    def error(self, message, **kwargs):
        entry = self._create_log_entry("ERROR", message, **kwargs)
        self.logger.error(entry)
    def warning(self, message, **kwargs):
        entry = self._create_log_entry("WARNING", message, **kwargs)
        self.logger.warning(entry)
# 使用示例
structured_logger = StructuredLogger()
# 记录分片同步日志
structured_logger.info(
    "开始同步分片",
    shard_id="shard_001",
    records_count=100,
    source="database_a",
    target="database_b"
)
structured_logger.info(
    "分片同步成功",
    shard_id="shard_001",
    duration_seconds=5.2,
    records_synced=100,
    speed_per_second=19.2
)
structured_logger.error(
    "分片同步失败",
    shard_id="shard_003",
    duration_seconds=3.1,
    error_type="ConnectionTimeout",
    error_message="连接数据库超时",
    retry_count=2
)

快速启动模板

如果你需要一个简单可用的方案:

import logging
from datetime import datetime
# 配置日志
logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s [%(levelname)s] %(message)s',
    datefmt='%Y-%m-%d %H:%M:%S',
    handlers=[
        logging.FileHandler(f'shard_sync_{datetime.now().strftime("%Y%m%d")}.log'),
        logging.StreamHandler()
    ]
)
logger = logging.getLogger(__name__)
def sync_shard(shard_info):
    start = datetime.now()
    logger.info(f"开始同步 {shard_info['name']}")
    try:
        # 你的同步逻辑
        # ...
        logger.info(f"完成同步 {shard_info['name']}")
        return True
    except Exception as e:
        logger.error(f"同步失败 {shard_info['name']}: {e}")
        return False
    finally:
        elapsed = (datetime.now() - start).total_seconds()
        logger.info(f"耗时 {elapsed:.2f} 秒")
# 主循环
shards = [...]  # 你的分片列表
for shard in shards:
    sync_shard(shard)

推荐实践

  1. 使用日志级别: DEBUG(调试), INFO(信息), WARNING(警告), ERROR(错误)
  2. 记录关键信息: 时间戳、分片ID、同步状态、耗时、错误详情
  3. 日志轮转: 使用RotatingFileHandler管理日志文件大小
  4. 分类存储: 普通日志和错误日志分开存储
  5. 结构化日志: 对于复杂系统,使用JSON格式便于后续分析

选择哪种方案取决于你的需求复杂度,建议从基础方案开始,逐步完善。

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