本文目录导读:

为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)
推荐实践
- 使用日志级别: DEBUG(调试), INFO(信息), WARNING(警告), ERROR(错误)
- 记录关键信息: 时间戳、分片ID、同步状态、耗时、错误详情
- 日志轮转: 使用
RotatingFileHandler管理日志文件大小 - 分类存储: 普通日志和错误日志分开存储
- 结构化日志: 对于复杂系统,使用JSON格式便于后续分析
选择哪种方案取决于你的需求复杂度,建议从基础方案开始,逐步完善。