如何从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))
关键修改点
- 用
response['Body'].iter_lines(decode_unicode=True)替代data.splitlines(),直接逐行读取S3文件流,避免一次性加载大文件到内存。 - 将裸
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())
关键修改点
- 构造S3文件路径:
f"s3://{s3_bucket_name}/{obj['Key']}",直接传给dd.read_csv,让Dask通过s3fs库直接读取S3文件。 - 提前安装
s3fs依赖:pip install s3fs,这是Dask访问S3的必要库。 - 调整
ARCID_no_hit的路径分割逻辑,适配S3路径的格式。 - 通过环境变量设置AWS凭证,更符合Dask的使用习惯(也可保留boto3配置方式,Dask会自动读取)。
内容的提问来源于stack exchange,提问作者Astro_raf
相关产品推荐
相关产品推荐

