如何进一步加速S3存储桶中500个CSV文件的合并流程
问题背景
- 存储情况:S3 Bucket中存储约500个CSV文件
备注:单个文件大小约1GB,单文件包含约600万行数据。
- 核心需求:将所有CSV文件拼接合并为单个文件,优化全流程处理速度,落地可执行的优化手段。
初始实现代码
最初采用pandas读取解析CSV的方案实现,代码如下:
import trio import boto3 import pandas as pd from functools import partial AWS_ID = 'Hidden' AWS_SECRET = 'Hidden' Bucket_Name = 'Hidden' limiter = trio.CapacityLimiter(10) async def read_object(bucket, object_csv, sender): async with limiter, sender: print(f'Reading {object_csv}') test = bucket.Object(object_csv) test = test.get()['Body'] data = await trio.to_thread.run_sync(partial(pd.read_csv, test, header=None)) await sender.send(data) async def main(): async with trio.open_nursery() as nurse: s3 = boto3.resource( service_name='s3', aws_access_key_id=AWS_ID, aws_secret_access_key=AWS_SECRET, ) bucket = s3.Bucket(Bucket_Name) allfiles = [i.key for i in bucket.objects.all()] sender, receiver = trio.open_memory_channel(0) nurse.start_soon(rec, receiver) async with sender: for f in allfiles: nurse.start_soon(read_object, bucket, f, sender.clone()) async def rec(receiver): alldf = [] async with receiver: async for df in receiver: alldf.append(df) final = pd.concat(alldf, ignore_index=True) print(final) if __name__ == "__main__": try: trio.run(main) except KeyboardInterrupt: exit('Job Cancelled!')
经性能定位,该版本的核心瓶颈为:
data = await trio.to_thread.run_sync(partial(pd.read_csv, test, header=None))
单文件通过pd.read_csv读取解析的耗时约为2分钟,即使开启多线程并行执行,整体处理耗时仍然过高。
第一次优化后的实现
调整实现逻辑,去掉pandas的CSV解析步骤,改为直接读取S3对象的二进制流后直接写入本地输出文件,更新后的代码如下:
limiter = trio.CapacityLimiter(10) async def read_object(bucket, object_csv, sender): async with limiter, sender: print(f'Reading {object_csv}') test = bucket.Object(object_csv) test = test.get()['Body'] data = await trio.to_thread.run_sync(test.read) await sender.send(data) print(f'Done Reading {object_csv}') async def main(): async with trio.open_nursery() as nurse: s3 = boto3.resource( service_name='s3', aws_access_key_id=AWS_ID, aws_secret_access_key=AWS_SECRET, ) bucket = s3.Bucket(Bucket_Name) sender, receiver = trio.open_memory_channel(0) nurse.start_soon(rec, receiver) async with sender: for csv in bucket.objects.all(): nurse.start_soon(read_object, bucket, csv.key, sender.clone()) async def rec(receiver): async with receiver, await trio.open_file('output.csv', 'wb') as f: count = 0 async for df in receiver: count += 1 await f.write(df) await f.write(b"\n") print(f'Collected {count}', flush=True, end='\r') if __name__ == "__main__": try: trio.run(main) except KeyboardInterrupt: exit('Job Cancelled!')
咨询问题
基于当前更新后的二进制流直读直写的代码实现,是否还有进一步提升处理速度的可行方案?
内容的提问来源于stack exchange,提问作者αԋɱҽԃ αмєяιcαη
相关产品推荐
相关产品推荐

