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

如何从AWS S3存储桶读取大体积CSV文件并合并为DataFrame?

解决S3存储桶CSV文件合并为DataFrame的两个问题

一、原代码大文件被跳过的修复

你的原代码中,大文件被跳过并非csv.reader的限制——csv.reader本身没有文件大小限制,它支持逐行处理流数据。问题出在response['Body'].read().decode('utf-8')这一步:一次性将整个大文件读取到内存,可能触发内存不足或解码异常,而宽泛的except语句直接跳过了出错的文件,没有给出任何报错信息。

修改后的原代码

import boto3
import csv
import pandas as pd

# 设置S3存储桶和CSV文件目录路径
aws_access_key_id ='XXXXXXXXXX'
aws_secret_access_key='XXXXXXXXXXXXXX'
s3_bucket_name = 'arcodp'
folder_name = 'lab_data/'

# 获取存储桶目录下所有CSV文件列表
s3 = boto3.client('s3', aws_access_key_id=aws_access_key_id, aws_secret_access_key=aws_secret_access_key)

paginator = s3.get_paginator('list_objects_v2')
pages = paginator.paginate(Bucket=s3_bucket_name, Prefix=folder_name)

csv_files = [obj['Key'] for page in pages for obj in page['Contents'] if obj['Key'].endswith('.csv')]

# 初始化存储数据的列表
df_list = []
ARCID_lst = []

# 逐个读取CSV文件并合并数据
for file in csv_files:
    try: 
        response = s3.get_object(Bucket=s3_bucket_name, Key=file)
        # 直接使用S3文件流逐行读取,避免一次性加载大文件到内存
        csv_reader = csv.reader(response['Body'].iter_lines(decode_unicode=True), delimiter='|', quoting=csv.QUOTE_NONE)
        rows_list = list(csv_reader)
        df_list.extend(rows_list)
    except Exception as e:
        # 捕获具体异常并打印,方便排查问题
        print(f"处理文件{file}时出错: {str(e)}")
        ARCID_no_hit = file.split('/')[1].split('_')[0]
        ARCID_lst.append(ARCID_no_hit)

# 转换为Pandas DataFrame
df_par = pd.DataFrame(df_list)

# 打印前10行数据
print(df_par.head(10))

关键修改点

  1. 用response['Body'].iter_lines(decode_unicode=True)替代data.splitlines(),直接逐行读取S3文件流,避免一次性加载大文件到内存。
  2. 将裸except改为except Exception as e,并打印错误信息,能明确知道大文件被跳过的具体原因(比如编码错误、内存不足等)。

二、Dask代码的修复

你的Dask代码报错“找不到本地文件”,是因为dd.read_csv需要接收文件路径(本地路径或S3路径),但你传入的是已经读取到内存的字符串data,它会把这个字符串当成本地文件名,自然找不到。

修改后的Dask代码

import boto3
from dask import delayed
import dask.dataframe as dd
import csv
import os

# 设置AWS凭证(也可通过环境变量、~/.aws/credentials文件配置)
os.environ['AWS_ACCESS_KEY_ID'] = 'XXXXXXXXXXXXX'
os.environ['AWS_SECRET_ACCESS_KEY'] = 'XXXXXXXXXX'

s3_bucket_name = 'arcodp'
folder_name = 'lab_data/'

# 获取存储桶目录下所有CSV文件的S3路径
s3 = boto3.client('s3', aws_access_key_id=os.environ['AWS_ACCESS_KEY_ID'], aws_secret_access_key=os.environ['AWS_SECRET_ACCESS_KEY'])

paginator = s3.get_paginator('list_objects_v2')
pages = paginator.paginate(Bucket=s3_bucket_name, Prefix=folder_name)

csv_files = [f"s3://{s3_bucket_name}/{obj['Key']}" for page in pages for obj in page['Contents'] if obj['Key'].endswith('.csv')]

df_list = []
ARCID_lst = []

# 注意:需先安装s3fs库才能让Dask访问S3,执行命令:pip install s3fs
for file in csv_files:
    try:
        # 直接传入S3路径,让Dask自行读取文件
        df = delayed(dd.read_csv)(file, sep='|', header=None, quoting=csv.QUOTE_NONE, engine='c')
        df_list.append(df)
    except Exception as e:
        print(f"处理文件{file}时出错: {str(e)}")
        ARCID_no_hit = file.split('/')[3].split('_')[0]  # S3路径格式为s3://bucket/folder/file,索引对应调整
        ARCID_lst.append(ARCID_no_hit)

# 合并所有延迟加载的Dask DataFrame
df_combined = dd.from_delayed(df_list)

# 计算得到最终的Pandas DataFrame
df_par = df_combined.compute()

# 打印前5行数据
print(df_par.head())

关键修改点

  1. 构造S3文件路径:f"s3://{s3_bucket_name}/{obj['Key']}",直接传给dd.read_csv,让Dask通过s3fs库直接读取S3文件。
  2. 提前安装s3fs依赖:pip install s3fs,这是Dask访问S3的必要库。
  3. 调整ARCID_no_hit的路径分割逻辑,适配S3路径的格式。
  4. 通过环境变量设置AWS凭证,更符合Dask的使用习惯(也可保留boto3配置方式,Dask会自动读取)。

内容的提问来源于stack exchange,提问作者Astro_raf

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 05:08:10