Python脚本如何统计分片同步耗时数据

wen python案例 29

本文目录导读:

Python脚本如何统计分片同步耗时数据

  1. 方法1:基础时间统计
  2. 方法2:使用装饰器自动统计
  3. 方法3:带进度和图表的数据统计
  4. 方法4:JSON日志输出
  5. 选择建议

方法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. 简单场景:使用方法1,快速集成
  2. 需要代码侵入低:使用方法2的装饰器
  3. 需要可视化:使用方法3,带图表展示
  4. 需要持久化:使用方法4,记录到JSON日志

这些方法可以根据你的具体需求进行组合和调整。

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