You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何进一步加速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αη

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.29 10:51:21