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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 05:25:29