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

