如何在Databricks中读取指定分区NDJSON并添加date列?
读取指定S3分区路径时自动添加date列的方法
场景描述
S3存储中数据按分区目录组织,路径格式为:
s3://my_bucket/my_prefix/date=2023-05-01/ s3://my_bucket/my_prefix/date=2023-05-02/ ... s3://my_bucket/my_prefix/date=2023-09-01/
每个分区下有NDJSON文件。读取整个存储桶时,Spark会自动识别分区并添加date列,但读取指定路径列表path_list时,默认不会生成该列,以下是几种解决方法:
方法1:用通配符路径替代路径列表(推荐)
直接传递包含分区标识的路径(单个或多个),Spark会自动解析分区字段生成date列:
# 连续日期用通配符匹配 df = spark.read.text("s3://my_bucket/my_prefix/date=2023-05-0[1-3]/") # 不连续日期直接列出目标路径 df = spark.read.text( "s3://my_bucket/my_prefix/date=2023-05-01/", "s3://my_bucket/my_prefix/date=2023-05-03/", "s3://my_bucket/my_prefix/date=2023-09-01/" )
方法2:手动从文件路径提取date字段
如果必须使用预定义的path_list,可以读取后通过路径解析添加列:
from pyspark.sql.functions import input_file_name, regexp_extract # 读取指定路径列表 df = spark.read.text(path_list) # 从文件路径中提取date值 df = df.withColumn( "date", regexp_extract(input_file_name(), r"date=(\d{4}-\d{2}-\d{2})", 1) )
input_file_name()获取每条记录对应的源文件路径,regexp_extract匹配路径中的date=YYYY-MM-DD格式,提取出日期字符串作为date列值。
方法3:读取父目录后过滤目标分区
先读取分区的父目录,让Spark自动识别所有分区,再过滤需要的日期:
# 读取父目录,自动识别date分区列 df = spark.read.text("s3://my_bucket/my_prefix/") # 过滤指定日期 target_dates = ["2023-05-01", "2023-05-03", "2023-09-01"] df = df.filter(df.date.isin(target_dates))
该方法适合分区总数不多的场景,避免手动构造路径列表。
内容的提问来源于stack exchange,提问作者Tristan Tran
相关产品推荐
相关产品推荐

