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

如何在Spark/HDFS中获取读取文件的文件名/路径以实现分区输出?

获取Spark读取文件的日期键并用于输出分区

这里有几种实用的方法可以获取文件名中的日期部分,适配你的每日单文件处理场景:

方法1:Spark DataFrame API 提取日期(推荐)

利用Spark内置的input_file_name()函数获取每条记录对应的源文件路径,再通过正则提取日期键:

from pyspark.sql.functions import input_file_name, regexp_extract

# 读取HDFS上的CSV文件
raw_df = spark.read.csv("hdfs:///raw/*.csv", header=True, inferSchema=True)

# 添加文件名列并提取日期部分
df_with_metadata = raw_df.withColumn("source_file", input_file_name()) \
                         .withColumn("date_key", regexp_extract("source_file", r"(\d{4}-\d{2}-\d{2})", 1))

# 由于每日仅上传一个文件,直接取唯一的日期键
date_key = df_with_metadata.select("date_key").distinct().collect()[0][0]

# 按日期分区写入S3
raw_df.write.mode("overwrite") \
     .option("header", "true") \
     .csv(f"s3://your-target-bucket/processed-data/date={date_key}/")

说明:

  • regexp_extract中的正则(\d{4}-\d{2}-\d{2})精准匹配文件名中的YYYY-MM-DD格式日期
  • 因为每日仅处理单个文件,distinct().collect()[0][0]可以安全获取唯一的日期值

方法2:RDD API 直接获取文件名

如果使用RDD处理数据,可以通过wholeTextFiles一次性拿到文件名和文件内容:

# 读取文件,返回(文件名, 文件内容)的RDD
file_rdd = spark.sparkContext.wholeTextFiles("hdfs:///raw/*.csv")

# 从文件名中提取日期键
date_key = file_rdd.map(lambda x: x[0].split("/")[-1].replace("-stats.csv", "")).distinct().collect()[0]

# 后续处理RDD并写入S3(示例)
processed_rdd = file_rdd.flatMap(lambda x: x[1].split("\n")[1:])  # 跳过表头
processed_rdd.saveAsTextFile(f"s3://your-target-bucket/processed-data/date={date_key}/")

方法3:外部传递日期键(更高效)

如果在上传文件到HDFS的环节已经知道日期,可以直接把日期作为参数传给Spark脚本,避免在Spark中解析文件名:

  1. 上传脚本中提取日期(比如从文件名中获取),然后启动Spark作业时传入参数:
# 假设上传的文件名是2022-07-27-stats.csv
DATE_KEY=$(basename /path/to/uploaded/file | cut -d'-' -f1-3)
spark-submit --conf spark.app.date_key=$DATE_KEY your-spark-script.py
  1. 在Spark脚本中读取该参数:
from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()
date_key = spark.sparkContext.getConf().get("spark.app.date_key")

# 直接用日期键写入S3
spark.read.csv("hdfs:///raw/*.csv", header=True) \
     .write.mode("overwrite") \
     .option("header", "true") \
     .csv(f"s3://your-target-bucket/processed-data/date={date_key}/")

这种方法不需要在Spark中处理文件元数据,性能更优,适合固定每日单文件的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 05:18:19