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

PySpark实现从HDFS文件夹名提取日期并合并至月度文件夹(无需新增列)

问题

当前HDFS路径/dev/data/下存在按日组织的文件夹结构,示例如下:

2024.03.30
   part-00001.avro
   part-00002.avro
2024.03.31
   part-00001.avro
   part-00002.avro
2024.04.01
   part-00001.avro
   part-00002.avro
2024.04.02
   part-00001.avro
   part-00002.avro

每个日文件夹下均包含Avro文件。现需将这些文件迁移至新路径,合并为按月份组织的文件夹结构,示例如下:

2024.03
2024.04
2024.05

所有Avro文件需放置在对应月度文件夹下。

已编写如下代码,通过现有exit_date列生成month列并分区写入:

day_df = spark.read.format("avro").load("path/to/dev/data")
month_df = day_df.withColumn("month", F.date_format(F.col("exit_date"),'yyyy-MM'))
month_df.write.partitionBy("month").format("com.databricks.spark.avro").save("path/to/dest")

现寻求无需创建month列即可实现该需求的方案。

解决方案

你可以直接在partitionBy中使用日期格式化表达式,无需额外创建month列,Spark会自动基于表达式计算分区值:

from pyspark.sql import functions as F

day_df = spark.read.format("avro").load("path/to/dev/data")
day_df.write.partitionBy(F.date_format(F.col("exit_date"), 'yyyy-MM')) \
           .format("com.databricks.spark.avro") \
           .save("path/to/dest")

如果需要和示例一致的yyyy.MM格式分区文件夹名,只需修改格式化字符串:

day_df.write.partitionBy(F.date_format(F.col("exit_date"), 'yyyy.MM')) \
           .format("com.databricks.spark.avro") \
           .save("path/to/dest")

要是不想依赖数据中的exit_date列,也可以从文件路径里提取原日文件夹的日期信息来生成月度分区:

from pyspark.sql import functions as F

day_df = spark.read.format("avro").load("path/to/dev/data")
day_df.write.partitionBy(
    F.date_format(
        F.to_date(F.regexp_extract(F.input_file_name(), r'(\d{4}\.\d{2}\.\d{2})', 1), 'yyyy.MM.dd'),
        'yyyy.MM'
    )
).format("com.databricks.spark.avro").save("path/to/dest")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 09:40:12