本文目录导读:

在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()
-
选择正确的并发方式:
- IO密集型:使用
ThreadPoolExecutor或asyncio - CPU密集型:使用
multiprocessing.Pool
- IO密集型:使用
-
结果收集策略:
- 使用
as_completed()实时获取完成的结果 - 使用
result()方法获取每个任务的返回值
- 使用
-
错误处理:
- 每个分片独立处理异常
- 汇总时分开记录成功和失败的结果
-
汇总之类设计:
- 提供汇总统计(总数、成功数、失败数)
- 可选的数据合并功能
- 保留分片级别的详细信息
根据你的具体场景选择合适的方案,最常用的是 concurrent.futures 方案。