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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 16:40:13