使用AWS Glue正确分区:解决分区后dt_venda列内容为空问题
问题描述
我有一个用于销售文件ETL的AWS Glue作业,目标是按销售日期(dt_venda)分区,每个分区目录以dt_venda=YYYY-MM-DD格式命名。但输出的CSV文件里,dt_venda列内容为空,我需要保留所有列与记录不变,同时让文件内的dt_venda列填充对应分区的日期值,分区目录命名保持正常。
当前使用的代码片段:
from pyspark.sql.functions import col from pyspark.sql import functions as F # 按销售日期分区DataFrame partitionedDataFrame = coalescedDataFrame.repartition("dt_venda") # 为每个分区生成CSV文件 partitionedDataFrame.write.partitionBy("dt_venda")\ .option("sep", ";")\ .option("header", "true")\ .mode("append")\ .csv("s3://s3-aws-araujo-sf-cx360-cdp-dev/03-cleansed/venda/") job.commit()
解决方案
这是Spark partitionBy 写入的默认行为——被用作分区字段的列会从输出文件中移除,仅保留在分区目录名里。要同时保留该列在文件中且填充对应值,有两种可行方法:
方法1:复制分区列,用副本做分区
直接给dt_venda列创建一个副本,用副本作为分区字段,原列保留在DataFrame中,这样输出文件里会保留dt_venda列的内容,分区目录命名不受影响。
修改后的代码:
from pyspark.sql.functions import col from pyspark.sql import functions as F # 复制dt_venda列作为分区专用字段 df_with_partition_col = coalescedDataFrame.withColumn("partition_dt_venda", col("dt_venda")) # 按副本列重分区并写入 partitionedDataFrame = df_with_partition_col.repartition("partition_dt_venda") partitionedDataFrame.write.partitionBy("partition_dt_venda")\ .option("sep", ";")\ .option("header", "true")\ .mode("append")\ .csv("s3://s3-aws-araujo-sf-cx360-cdp-dev/03-cleansed/venda/") job.commit()
输出的分区目录为partition_dt_venda=YYYY-MM-DD,同时CSV文件内的dt_venda列会保留原始日期值。
方法2:关闭Spark分区列移除特性(Spark 3.1+支持)
如果你的AWS Glue作业使用Spark 3.1及以上版本(比如Glue 3.0/4.0),可以通过配置参数直接让分区列保留在输出文件中,无需修改DataFrame结构。
修改后的代码:
from pyspark.sql.functions import col from pyspark.sql import functions as F # 配置Spark参数,禁用分区列类型推断并保留分区列 spark.conf.set("spark.sql.sources.partitionColumnTypeInference.enabled", "false") partitionedDataFrame = coalescedDataFrame.repartition("dt_venda") partitionedDataFrame.write.partitionBy("dt_venda")\ .option("sep", ";")\ .option("header", "true")\ .option("keepPartitionColumns", "true") # 关键选项:保留分区列在文件内 .mode("append")\ .csv("s3://s3-aws-araujo-sf-cx360-cdp-dev/03-cleansed/venda/") job.commit()
这种方法下,分区目录依然是dt_venda=YYYY-MM-DD,CSV文件内的dt_venda列也会填充对应日期值。
注意事项
- 方法2仅适用于Spark 3.1及更高版本,若使用Glue 2.0(对应Spark 2.4),只能采用方法1。
- 需确认
coalescedDataFrame中的dt_venda列本身有有效值,若原列值为空,即使保留列也会输出空值——你当前的问题是Spark默认移除了分区列,而非原列值为空。
内容的提问来源于stack exchange,提问作者Alexandre
相关产品推荐
相关产品推荐

