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

如何用AWS Glue/Spark将S3分区CSV转为按日分区拆分的Parquet

我之前在处理PB级数据迁移的时候,碰到过几乎一模一样的场景——带分区的CSV源表有分区错误、每日几百GB增量,还要转成指定文件数的Parquet分区表。给你梳理一套落地性很强的解决方案,覆盖所有痛点:

解决方案:从分区CSV到按日分区Parquet(含错误处理、文件拆分、资源配置)

一、先搞定源CSV的分区错误问题

分区错误是最容易踩的坑,比如S3路径标了dt=2024-05-01,但数据里的真实日期是2024-05-02,不先修正的话,转出来的Parquet表分区完全没用。

1. 读取时同时获取路径分区和真实数据日期

不要完全依赖Glue数据目录的元数据,直接读S3路径,把路径里的分区字段和数据里的真实日期拉出来对比:

from pyspark.sql.functions import input_file_name, regexp_extract

# 直接读取S3路径的CSV,跳过Glue目录的分区元数据
df = spark.read.csv("s3://your-source-bucket/raw-path/", header=True, inferSchema=True)

# 从文件路径中提取分区日期(适配S3分区路径格式:dt=YYYY-MM-DD)
df = df.withColumn("path_partition_dt", regexp_extract(input_file_name(), r'dt=(\d{4}-\d{2}-\d{2})', 1))

# 假设数据里的真实日期字段是`event_date`,对比路径分区和真实日期
df = df.withColumn("is_partition_error", df.path_partition_dt != df.event_date)

2. 隔离或修正错误数据

  • 要是想快速推进流程,先把错误数据隔离到单独路径,后续人工排查:
# 拆分正确/错误数据集
valid_df = df.filter(df.is_partition_error == False)
invalid_df = df.filter(df.is_partition_error == True)

# 错误数据写入隔离桶
invalid_df.write.mode("append").csv("s3://your-error-bucket/partition-mismatch/")
  • 要是想自动修正,就把错误数据按真实日期重新分区写入临时路径,再用S3 CLI/Glue API把文件移到正确的源分区目录里(适合规则明确的错误场景)。

二、转换为按日分区、指定文件数的Parquet

这一步的核心是控制每个日期分区的文件数量,同时保证分区逻辑正确。

1. 精准控制每个日期的文件数

用repartition(目标文件数, 分区字段)的组合,既能按日期分区,又能让每个日期分区刚好生成n个文件:

# 设定每日要拆分的文件数,比如n=10(根据你的单日数据量调整,几百GB的话10-20个比较合适)
target_file_count_per_day = 10

# 先按真实日期分区,再在每个分区内拆分指定数量的文件
processed_df = valid_df.repartition(target_file_count_per_day, "event_date")

# 写入Parquet,按event_date做分区
processed_df.write.mode("append") \
  .partitionBy("event_date") \
  .parquet("s3://your-target-bucket/parquet-data/")

为什么不用coalesce?coalesce是减少分区数,适合小数据量;repartition是重新洗牌分区,能保证每个日期分区的文件数均匀,更适合几百GB的大场景。

2. 更新Glue数据目录(可选但推荐)

如果要在Glue里直接查询Parquet表,用Spark SQL创建外部表并刷新分区:

# 创建临时视图
processed_df.createOrReplaceTempView("daily_parquet_view")

# 创建外部表(替换成你的字段和库名)
spark.sql("""
CREATE EXTERNAL TABLE IF NOT EXISTS your_glue_db.daily_parquet_table (
  user_id STRING,
  event_type STRING,
  event_timestamp TIMESTAMP
)
PARTITIONED BY (event_date STRING)
STORED AS PARQUET
LOCATION 's3://your-target-bucket/parquet-data/'
""")

# 刷新分区(增量写入后必须执行)
spark.sql("MSCK REPAIR TABLE your_glue_db.daily_parquet_table")

三、资源配置优化(针对每日数百GB数据)

几百GB的数据,资源配不好要么跑不动,要么浪费钱,给你几个实战参数:

1. Glue作业节点配置

  • 节点类型:优先选G.2X(8vCPU、32GB内存),CPU密集型场景足够;如果有复杂UDF或大字段处理,升级到G.4X
  • 节点数量:按经验,每100GB数据配5-10个G.2X节点。比如300GB数据,配15-20个节点(如果源文件是大量小文件,节点数可以再加5-10个,提升并行处理能力)
  • 作业参数:在Glue作业的「Job parameters」里设置:
    • --spark.driver.cores: 4
    • --spark.driver.memory: 16g
    • --spark.executor.cores: 4
    • --spark.executor.memory: 16g
    • --spark.executor.instances: 节点数-1(留一个节点给driver)

2. Spark性能优化参数

在代码开头加这些配置,提升大文件处理和分区效率:

# 优化Parquet写入速度和压缩比
spark.conf.set("spark.sql.parquet.compression.codec", "snappy")
spark.conf.set("spark.sql.parquet.writeLegacyFormat", "false")

# 增量写入时只覆盖目标分区,避免全表重写
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")

# 调整洗牌分区数,和executor数量匹配(一般是executor数*4)
spark.conf.set("spark.sql.shuffle.partitions", 80)

四、每日增量自动化处理

既然是每日新增分区,用Glue触发器实现自动化:

  • 创建定时触发器,每天凌晨等前一天数据写完后触发(比如凌晨2点)
  • 代码里通过日期参数读取前一天的源数据,避免全量扫描:
from datetime import datetime, timedelta

# 计算前一天的日期
yesterday = (datetime.today() - timedelta(days=1)).strftime("%Y-%m-%d")

# 只读取前一天的源分区
df = spark.read.csv(f"s3://your-source-bucket/raw-path/dt={yesterday}/", header=True, inferSchema=True)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 04:35:28