如何用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
相关产品推荐
相关产品推荐

