同一S3桶内跨文件夹拷贝大量小JSON文件的优化方案咨询
S3百万级小文件迁移/合并优化方案
一、保留单个文件迁移的优化方案
原代码性能问题原因
- 所有copy操作串行执行,S3 API单次请求有网络延迟,百万级文件串行处理耗时会指数级上升
- 未使用分页器遍历对象,对象量超过1000时会出现漏遍历的问题
优化实现(多线程并行处理)
IO密集型场景用多线程可以大幅提升处理效率,示例代码如下:
import boto3 from concurrent.futures import ThreadPoolExecutor, as_completed s3 = boto3.resource('s3') bucket_name = 'dev-bucket' source_prefix = 'data/' target_prefix = 'final/' bucket = s3.Bucket(bucket_name) MAX_WORKERS = 50 # 可根据实际网络情况调整,S3单桶QPS上限很高,50-200都可尝试 def copy_s3_obj(obj): # 跳过前缀本身的目录对象 if obj.key.endswith('/'): return 0 # 提取原文件名,拼到目标前缀下 file_name = obj.key.split('/')[-1] target_key = f"{target_prefix}{file_name}" # 注意如果有重名文件会被覆盖,需要的话可以加重名逻辑 bucket.Object(target_key).copy_from(CopySource={'Bucket': bucket_name, 'Key': obj.key}) return 1 if __name__ == "__main__": file_count = 0 # 用分页器遍历所有源对象 all_objs = list(bucket.objects.filter(Prefix=source_prefix)) try: with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor: futures = [executor.submit(copy_s3_obj, obj) for obj in all_objs] for future in as_completed(futures): file_count += future.result() print(f"迁移完成,共处理{file_count}个文件") except Exception as e: print(f"迁移失败,错误信息:{e}") exit(1)
如果文件量实在太大,也可以直接用AWS原生的S3 Batch Operations,直接在AWS侧执行批量复制,不需要本地跑代码。
二、合并为单个JSON文件后迁移的实现方案
原代码问题
response = response.append(text)错误:list.append()是原地修改方法,返回值为None,执行后response会被赋值为None,后续操作直接报错- 所有get_object操作串行执行,读取效率极低
- 没有处理JSON格式合法性,直接拼接字符串会导致最终文件不是合法JSON
优化实现
import boto3 import json from concurrent.futures import ThreadPoolExecutor, as_completed s3_client = boto3.client('s3') bucket_name = 'dev-bucket' source_prefix = 'data/' target_key = 'final/merged.json' MAX_WORKERS = 50 def read_s3_json(obj_key): try: resp = s3_client.get_object(Bucket=bucket_name, Key=obj_key) content = resp['Body'].read().decode('utf-8') return json.loads(content) except Exception as e: print(f"读取文件{obj_key}失败:{e}") return None if __name__ == "__main__": # 先遍历所有源文件key all_keys = [] paginator = s3_client.get_paginator('list_objects_v2') for page in paginator.paginate(Bucket=bucket_name, Prefix=source_prefix): if 'Contents' not in page: continue for obj in page['Contents']: key = obj['Key'] if not key.endswith('.json'): continue all_keys.append(key) # 多线程并行读取所有JSON merged_data = [] with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor: futures = [executor.submit(read_s3_json, key) for key in all_keys] for future in as_completed(futures): data = future.result() if data: # 如果每个文件是单条对象直接加,如果是数组就extend if isinstance(data, list): merged_data.extend(data) else: merged_data.append(data) # 写入合并后的JSON到S3 s3_client.put_object( Bucket=bucket_name, Key=target_key, Body=json.dumps(merged_data, ensure_ascii=False).encode('utf-8') ) print(f"合并完成,共写入{len(merged_data)}条记录")
如果总数据量过大内存装不下,可以改成流式写入JSON Lines格式,边读边写,不需要把所有数据都存在内存里。
内容的提问来源于stack exchange,提问作者user16443603
相关产品推荐
相关产品推荐

