如何在AWS Glue中使用PySpark处理多个CSV文件?
AWS Glue动态帧处理多CSV文件实用指南
一、读取S3路径下多个/全部CSV并自动合并
你当前的代码其实已经支持读取指定S3文件夹下的所有CSV文件——当paths参数指向S3文件夹而非单个文件时,Glue会自动加载该文件夹内所有CSV文件并合并成一个动态帧。如果需要读取子文件夹内的文件,只需在connection_options里添加recurse=True参数即可。
补充带表头的读取示例(避免把表头当成数据行):
t = glueC.create_dynamic_frame_from_options( connection_type="s3", connection_options={ "paths": ["s3://somebucket/inputfolder/"], "recurse": True # 可选,开启后读取子文件夹内的文件 }, format="csv", format_options={"withHeader": True} # 声明CSV包含表头 )
读取完成后,多个文件的数据已经自动合并在动态帧t中,无需额外合并操作。
二、动态帧处理数据的行数/数据量上限
PySpark和Glue动态帧本身没有硬编码的行数或数据量上限,它的处理能力完全取决于你配置的Glue集群资源(节点数量、CPU、内存、存储容量)以及数据的分区策略。只要集群资源足够,PB级别的数据也能处理。
实际操作中需要注意:合理设置数据分区,避免单个分区数据过大(导致内存溢出)或小文件过多(影响读写性能),这是实践中的优化点,而非技术上限。
三、合并处理后写入不同CSV文件
根据你的需求,分两种常见场景给出实现方式:
1. 按条件拆分输出到独立路径
如果需要按业务逻辑(比如某字段值)拆分数据到不同文件,可通过Filter转换拆分动态帧后分别写入:
from awsglue.transforms import Filter # 假设动态帧t中有category字段,按该字段值拆分 frame_a = Filter.apply(frame=t, f=lambda x: x["category"] == "A") frame_b = Filter.apply(frame=t, f=lambda x: x["category"] == "B") # 写入不同S3路径 glueC.write_dynamic_frame.from_options( frame=frame_a, connection_type="s3", connection_options={"path": "s3://somebucket/output/category_a/"}, format="csv", format_options={"withHeader": True} ) glueC.write_dynamic_frame.from_options( frame=frame_b, connection_type="s3", connection_options={"path": "s3://somebucket/output/category_b/"}, format="csv", format_options={"withHeader": True} )
2. 按字段自动分区输出
如果希望按某个字段(比如日期)自动创建子文件夹存储对应数据,可使用partitionKeys参数:
glueC.write_dynamic_frame.from_options( frame=t, connection_type="s3", connection_options={ "path": "s3://somebucket/output/", "partitionKeys": ["date"] # 按date字段分区 }, format="csv", format_options={"withHeader": True} )
执行后会在输出路径下生成类似date=2024-05-20/的子文件夹,每个文件夹下存储对应日期的数据文件。
内容的提问来源于stack exchange,提问作者nohardfeelings
相关产品推荐
相关产品推荐

