Python脚本如何汇总分片同步执行结果

wen python案例 30

本文目录导读:

Python脚本如何汇总分片同步执行结果

  1. 使用 concurrent.futures(推荐)
  2. 使用 multiprocessing(适合CPU密集型)
  3. 使用 asyncio(异步IO密集型)
  4. 实际应用示例:分片处理文件

在Python中汇总分片同步执行结果,常见的有几种方案,我给你整理最实用的几种:

使用 concurrent.futures(推荐)

from concurrent.futures import ThreadPoolExecutor, as_completed
import time
def process_shard(shard_id, data):
    """处理单个分片"""
    print(f"处理分片 {shard_id} 开始")
    time.sleep(1)  # 模拟耗时操作
    result = f"分片{shard_id}的结果: {len(data)}条数据"
    return shard_id, result
def main():
    # 准备分片数据
    shards = {
        1: ["a", "b", "c"],
        2: ["d", "e"],
        3: ["f", "g", "h", "i"]
    }
    results = {}  # 汇总结果
    # 使用线程池并行执行
    with ThreadPoolExecutor(max_workers=3) as executor:
        # 提交所有任务
        future_to_shard = {
            executor.submit(process_shard, sid, data): sid 
            for sid, data in shards.items()
        }
        # 收集结果
        for future in as_completed(future_to_shard):
            shard_id, result = future.result()
            results[shard_id] = result
            print(f"完成: {result}")
    # 汇总所有结果
    print("\n=== 汇总结果 ===")
    for shard_id in sorted(results.keys()):
        print(f"{shard_id}: {results[shard_id]}")
    # 如果需要,可以合并结果
    total_data_count = sum(len(data) for data in shards.values())
    print(f"总数据量: {total_data_count}条")
if __name__ == "__main__":
    main()

使用 multiprocessing(适合CPU密集型)

from multiprocessing import Pool
import time
def process_shard(shard_info):
    """处理分片(用于多进程)"""
    shard_id, data = shard_info
    print(f"进程处理分片 {shard_id}")
    time.sleep(1)
    return shard_id, sum(data)  # 假设数据是数字
def main():
    shards = [
        (1, [1, 2, 3, 4, 5]),
        (2, [10, 20, 30]),
        (3, [100, 200])
    ]
    # 使用进程池
    with Pool(processes=3) as pool:
        results = pool.map(process_shard, shards)
    # 汇总结果
    print("=== 汇总结果 ===")
    total_sum = 0
    for shard_id, result in sorted(results):
        print(f"分片{shard_id}: 和={result}")
        total_sum += result
    print(f"总和: {total_sum}")
if __name__ == "__main__":
    main()

使用 asyncio(异步IO密集型)

import asyncio
import time
async def process_shard(shard_id, data):
    """异步处理分片"""
    print(f"异步处理分片 {shard_id} 开始")
    await asyncio.sleep(1)  # 模拟IO操作
    return shard_id, f"处理结果_{shard_id}"
async def main():
    shards = {
        1: ["data1"],
        2: ["data2"],
        3: ["data3"]
    }
    # 并发执行所有分片
    tasks = [
        process_shard(sid, data) 
        for sid, data in shards.items()
    ]
    results = await asyncio.gather(*tasks)
    # 汇总结果
    print("\n=== 汇总结果 ===")
    summary = {}
    for shard_id, result in results:
        summary[shard_id] = result
        print(f"分片{shard_id}: {result}")
    return summary
if __name__ == "__main__":
    results = asyncio.run(main())

实际应用示例:分片处理文件

from concurrent.futures import ThreadPoolExecutor, as_completed
import os
import json
class ShardProcessor:
    def __init__(self, max_workers=4):
        self.max_workers = max_workers
        self.results = {}
        self.errors = []
    def process_file_shard(self, file_path, start_line, end_line):
        """处理文件的一个分片"""
        results = []
        try:
            with open(file_path, 'r') as f:
                for i, line in enumerate(f):
                    if start_line <= i < end_line:
                        # 处理每行数据
                        processed = self.process_line(line.strip())
                        results.append(processed)
            return {
                'shard': f"{start_line}-{end_line}",
                'success': True,
                'data': results,
                'count': len(results)
            }
        except Exception as e:
            return {
                'shard': f"{start_line}-{end_line}",
                'success': False,
                'error': str(e)
            }
    def process_line(self, line):
        """处理单行数据(示例)"""
        # 这里实现具体的处理逻辑
        return line.upper()
    def split_file_into_shards(self, file_path, num_shards):
        """将文件分割成多个分片"""
        total_lines = sum(1 for _ in open(file_path))
        shard_size = total_lines // num_shards
        shards = []
        for i in range(num_shards):
            start = i * shard_size
            end = start + shard_size if i < num_shards - 1 else total_lines
            shards.append((file_path, start, end))
        return shards
    def process_all_shards(self, file_path, num_shards=4):
        """并行处理所有分片并汇总结果"""
        shards = self.split_file_into_shards(file_path, num_shards)
        with ThreadPoolExecutor(max_workers=self.max_workers) as executor:
            # 提交任务
            future_to_shard = {
                executor.submit(self.process_file_shard, *shard): shard
                for shard in shards
            }
            # 收集结果
            for future in as_completed(future_to_shard):
                result = future.result()
                shard_info = future_to_shard[future]
                if result['success']:
                    self.results[result['shard']] = result['data']
                else:
                    self.errors.append(result['error'])
        # 汇总和统计
        summary = self.generate_summary()
        return summary
    def generate_summary(self):
        """生成汇总报告"""
        summary = {
            'total_shards': len(self.results),
            'total_records': sum(len(data) for data in self.results.values()),
            'error_count': len(self.errors),
            'shard_details': {
                shard: {
                    'count': len(data),
                    'sample': data[:3]  # 样本数据
                }
                for shard, data in self.results.items()
            }
        }
        if self.errors:
            summary['errors'] = self.errors
        # 合并所有数据
        all_data = []
        for shard_id in sorted(self.results.keys()):
            all_data.extend(self.results[shard_id])
        summary['merged_data'] = all_data
        return summary
# 使用示例
def main():
    processor = ShardProcessor(max_workers=4)
    # 假设有一个大文件
    test_file = "large_data.txt"
    # 创建测试文件
    with open(test_file, 'w') as f:
        for i in range(100):
            f.write(f"line_{i}\n")
    # 处理所有分片并汇总
    summary = processor.process_all_shards(test_file, num_shards=4)
    # 输出汇总结果
    print(json.dumps(summary, indent=2))
    # 清理测试文件
    os.remove(test_file)
if __name__ == "__main__":
    main()
  1. 选择正确的并发方式

    • IO密集型:使用 ThreadPoolExecutorasyncio
    • CPU密集型:使用 multiprocessing.Pool
  2. 结果收集策略

    • 使用 as_completed() 实时获取完成的结果
    • 使用 result() 方法获取每个任务的返回值
  3. 错误处理

    • 每个分片独立处理异常
    • 汇总时分开记录成功和失败的结果
  4. 汇总之类设计

    • 提供汇总统计(总数、成功数、失败数)
    • 可选的数据合并功能
    • 保留分片级别的详细信息

根据你的具体场景选择合适的方案,最常用的是 concurrent.futures 方案。

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