本文目录导读:

方法1:基础时间统计
import time
import threading
from typing import List, Dict
class ShardSyncStats:
def __init__(self):
self.shard_times: Dict[int, float] = {}
self.lock = threading.Lock()
def record_shard_start(self, shard_id: int):
"""记录分片开始同步时间"""
with self.lock:
self.shard_times[shard_id] = {'start': time.time(), 'end': None}
def record_shard_end(self, shard_id: int):
"""记录分片结束同步时间"""
with self.lock:
if shard_id in self.shard_times:
self.shard_times[shard_id]['end'] = time.time()
def get_shard_duration(self, shard_id: int) -> float:
"""获取单个分片耗时"""
if shard_id in self.shard_times:
data = self.shard_times[shard_id]
if data['end']:
return data['end'] - data['start']
return 0.0
def get_total_stats(self) -> Dict:
"""获取总体统计"""
durations = [
data['end'] - data['start']
for data in self.shard_times.values()
if data['end']
]
if not durations:
return {}
return {
'total_shards': len(durations),
'total_time': sum(durations),
'avg_time': sum(durations) / len(durations),
'max_time': max(durations),
'min_time': min(durations),
'std_dev': self._std_dev(durations)
}
def _std_dev(self, values: List[float]) -> float:
"""计算标准差"""
avg = sum(values) / len(values)
variance = sum((x - avg) ** 2 for x in values) / len(values)
return variance ** 0.5
# 使用示例
stats = ShardSyncStats()
def sync_shard(shard_id: int, stats: ShardSyncStats):
"""模拟分片同步"""
stats.record_shard_start(shard_id)
# 模拟同步操作
time.sleep(0.5) # 实际替换为真正的同步逻辑
stats.record_shard_end(shard_id)
print(f"分片 {shard_id} 同步完成,耗时: {stats.get_shard_duration(shard_id):.2f}秒")
# 并发执行
threads = []
for i in range(5):
t = threading.Thread(target=sync_shard, args=(i, stats))
threads.append(t)
t.start()
for t in threads:
t.join()
# 输出统计结果
print("\n=== 同步统计结果 ===")
for shard_id in range(5):
print(f"分片 {shard_id}: {stats.get_shard_duration(shard_id):.2f}秒")
print(f"\n总体统计: {stats.get_total_stats()}")
方法2:使用装饰器自动统计
import time
import functools
from collections import defaultdict
def sync_stats_logger(func):
"""分片同步统计装饰器"""
stats_data = defaultdict(list)
@functools.wraps(func)
def wrapper(shard_id, *args, **kwargs):
start_time = time.time()
result = func(shard_id, *args, **kwargs)
end_time = time.time()
duration = end_time - start_time
stats_data[func.__name__].append({
'shard_id': shard_id,
'duration': duration,
'timestamp': time.strftime('%Y-%m-%d %H:%M:%S')
})
print(f"[{time.strftime('%H:%M:%S')}] 分片 {shard_id} 耗时: {duration:.3f}秒")
return result
# 添加统计方法
wrapper.get_stats = lambda: dict(stats_data)
wrapper.print_summary = lambda: print_stats_summary(stats_data)
return wrapper
def print_stats_summary(stats_data):
"""打印统计摘要"""
print("\n=== 同步耗时统计摘要 ===")
for func_name, records in stats_data.items():
durations = [r['duration'] for r in records]
if not durations:
continue
print(f"\n函数: {func_name}")
print(f"总同步分片数: {len(durations)}")
print(f"总耗时: {sum(durations):.3f}秒")
print(f"平均耗时: {sum(durations)/len(durations):.3f}秒")
print(f"最大耗时: {max(durations):.3f}秒")
print(f"最小耗时: {min(durations):.3f}秒")
# 耗时分布
slow_count = sum(1 for d in durations if d > 1.0)
normal_count = len(durations) - slow_count
print(f"耗时>1秒: {slow_count}个")
print(f"耗时<=1秒: {normal_count}个")
# 使用示例
@sync_stats_logger
def sync_shard_data(shard_id: int):
"""实际的分片同步函数"""
# 模拟不同分片的同步时间
time.sleep(0.3 + shard_id * 0.1)
return f"分片 {shard_id} 数据同步完成"
# 执行同步
for i in range(5):
sync_shard_data(i)
# 打印统计摘要
sync_shard_data.print_summary()
方法3:带进度和图表的数据统计
import time
import matplotlib.pyplot as plt
from datetime import datetime
import numpy as np
class ShardSyncAnalyzer:
"""分片同步分析器"""
def __init__(self):
self.records = []
self.start_time = None
def start_sync(self):
"""开始整体同步"""
self.start_time = time.time()
def record_shard(self, shard_id: int, data_size: int = 0):
"""记录分片同步"""
return ShardSyncTimer(self, shard_id, data_size)
def add_record(self, shard_id: int, duration: float,
data_size: int = 0, status: str = "success"):
"""添加记录"""
self.records.append({
'shard_id': shard_id,
'duration': duration,
'data_size': data_size,
'status': status,
'timestamp': datetime.now()
})
def generate_report(self):
"""生成统计报告"""
if not self.records:
return
durations = [r['duration'] for r in self.records
if r['status'] == 'success']
total_time = time.time() - self.start_time if self.start_time else sum(durations)
print("=" * 50)
print("分片同步耗时统计报告")
print("=" * 50)
print(f"\n同步时间: {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}")
print(f"总分片数: {len(self.records)}")
print(f"成功分片: {len(durations)}")
print(f"失败分片: {len(self.records) - len(durations)}")
print(f"\n--- 耗时统计 ---")
print(f"总耗时: {total_time:.2f}秒")
print(f"平均耗时: {np.mean(durations):.3f}秒")
print(f"中位数耗时: {np.median(durations):.3f}秒")
print(f"标准差: {np.std(durations):.3f}秒")
print(f"最大耗时: {max(durations):.3f}秒")
print(f"最小耗时: {min(durations):.3f}秒")
# 百分位数
print(f"\n--- 百分位数 ---")
for p in [50, 75, 90, 95, 99]:
val = np.percentile(durations, p)
print(f"P{p}: {val:.3f}秒")
# 绘制图表
self._plot_stats(durations)
def _plot_stats(self, durations):
"""绘制统计图表"""
fig, axes = plt.subplots(2, 2, figsize=(12, 8))
# 1. 耗时直方图
axes[0, 0].hist(durations, bins=20, edgecolor='black')
axes[0, 0].set_title('分片同步耗时分布')
axes[0, 0].set_xlabel('耗时(秒)')
axes[0, 0].set_ylabel('分片数量')
# 2. 耗时折线图
axes[0, 1].plot(range(len(durations)), durations, 'b-', alpha=0.7)
axes[0, 1].axhline(y=np.mean(durations), color='r',
linestyle='--', label='平均值')
axes[0, 1].set_title('分片同步耗时趋势')
axes[0, 1].set_xlabel('分片序号')
axes[0, 1].set_ylabel('耗时(秒)')
axes[0, 1].legend()
# 3. 箱线图
axes[1, 0].boxplot(durations)
axes[1, 0].set_title('耗时箱线图')
axes[1, 0].set_ylabel('耗时(秒)')
# 4. 累计分布
sorted_durations = np.sort(durations)
cumulative = np.arange(1, len(sorted_durations) + 1) / len(sorted_durations)
axes[1, 1].plot(sorted_durations, cumulative, 'g-')
axes[1, 1].set_title('耗时累计分布')
axes[1, 1].set_xlabel('耗时(秒)')
axes[1, 1].set_ylabel('累计概率')
plt.tight_layout()
plt.savefig('shard_sync_stats.png')
plt.show()
class ShardSyncTimer:
"""分片同步计时器(上下文管理器)"""
def __init__(self, analyzer: ShardSyncAnalyzer,
shard_id: int, data_size: int = 0):
self.analyzer = analyzer
self.shard_id = shard_id
self.data_size = data_size
self.start_time = None
def __enter__(self):
self.start_time = time.time()
return self
def __exit__(self, exc_type, exc_val, exc_tb):
duration = time.time() - self.start_time
status = "failed" if exc_type else "success"
self.analyzer.add_record(
self.shard_id,
duration,
self.data_size,
status
)
# 实时输出
print(f"[分片 {self.shard_id:3d}] " +
f"{'✓' if status == 'success' else '✗'} " +
f"耗时: {duration:.3f}秒")
return exc_type is None
# 使用示例
def main():
analyzer = ShardSyncAnalyzer()
analyzer.start_sync()
# 模拟同步10个分片
for i in range(10):
# 模拟不同大小和耗时的分片
data_size = np.random.randint(100, 1000)
sync_time = np.random.exponential(0.5) + 0.1
with analyzer.record_shard(i, data_size) as timer:
time.sleep(sync_time) # 实际同步操作
# 模拟失败场景(随机)
if np.random.random() < 0.1:
raise Exception("同步失败")
# 生成报告
analyzer.generate_report()
if __name__ == "__main__":
main()
方法4:JSON日志输出
import json
import time
from datetime import datetime
def log_shard_sync(shard_id: int, start_time: float,
end_time: float, status: str = "success"):
"""记录分片同步日志"""
log_entry = {
"timestamp": datetime.now().isoformat(),
"shard_id": shard_id,
"duration": round(end_time - start_time, 3),
"status": status,
"start_time": start_time,
"end_time": end_time
}
# 写入日志文件
with open("shard_sync_log.json", "a") as f:
f.write(json.dumps(log_entry) + "\n")
return log_entry
def analyze_sync_log(log_file: str = "shard_sync_log.json"):
"""分析同步日志"""
records = []
with open(log_file, "r") as f:
for line in f:
if line.strip():
records.append(json.loads(line))
if not records:
print("没有找到同步记录")
return
# 统计分析
durations = [r["duration"] for r in records if r["status"] == "success"]
print(f"总记录数: {len(records)}")
print(f"成功: {len(durations)}, 失败: {len(records) - len(durations)}")
print(f"总耗时: {sum(durations):.3f}秒")
print(f"平均耗时: {sum(durations)/len(durations):.3f}秒")
# 输出耗时较长的分片
slow_shards = [r for r in records if r["duration"] > 1.0]
if slow_shards:
print(f"\n耗时较长的分片 ({len(slow_shards)}个):")
for r in sorted(slow_shards, key=lambda x: x["duration"], reverse=True)[:5]:
print(f" 分片 {r['shard_id']}: {r['duration']:.3f}秒")
# 使用示例
def sync_with_logging():
for i in range(5):
start = time.time()
# 模拟同步
time.sleep(0.5)
end = time.time()
log_entry = log_shard_sync(i, start, end)
print(f"已记录分片 {i} 的同步信息")
if __name__ == "__main__":
sync_with_logging()
analyze_sync_log()
选择建议
- 简单场景:使用方法1,快速集成
- 需要代码侵入低:使用方法2的装饰器
- 需要可视化:使用方法3,带图表展示
- 需要持久化:使用方法4,记录到JSON日志
这些方法可以根据你的具体需求进行组合和调整。