PySpark读取Azure Blob存储JSON时如何获取文件时间戳?
获取Azure Blob存储JSON文件的时间戳并转为DataFrame
首先明确:没有和input_file_name()完全对应的内置函数直接获取文件时间戳,但可以通过以下几种方案在读取数据的同时获取文件的时间戳(通常是最后修改时间):
方案1:使用Spark的binaryFile格式读取(推荐)
Spark的binaryFile数据源会返回文件的完整元数据,包括修改时间、文件路径、文件大小等,同时保留文件内容,你可以解析内容为JSON并合并元数据:
# 读取Blob存储中的JSON文件,获取元数据和二进制内容 df = spark.read.format("binaryFile") \ .option("pathGlobFilter", "*.json") # 只过滤JSON文件 .option("recursiveFileLookup", "true") # 递归读取子目录(可选) .load("/mnt/your-blob-container/target-path/") # 解析二进制内容为JSON字符串,再转为结构化数据 from pyspark.sql.functions import udf, col, from_json from pyspark.sql.types import StructType, StringType, IntegerType import json # 自定义UDF解析二进制内容 def parse_binary_to_json(content): return json.loads(content.decode("utf-8")) parse_json_udf = udf(parse_binary_to_json, StringType()) # 替换为你的JSON实际Schema json_schema = StructType() \ .add("user_id", IntegerType()) \ .add("user_name", StringType()) \ .add("action", StringType()) # 解析JSON并合并元数据字段 final_df = df.withColumn("json_data", from_json(parse_json_udf(col("content")), json_schema)) \ .select( "path", # 文件路径,等价于input_file_name()的结果 "modificationTime", # 文件最后修改时间戳 "json_data.*" # 展开JSON的所有字段 ) final_df.show()
方案2:先获取Blob元数据再关联DataFrame
如果需要单独获取文件元数据,可通过Azure Blob Storage SDK先拉取所有目标文件的时间戳,再和读取的JSON DataFrame关联:
from azure.storage.blob import BlobServiceClient from pyspark.sql.functions import lit # 初始化Blob客户端 connection_string = "your-azure-blob-connection-string" container_name = "your-container-name" blob_service_client = BlobServiceClient.from_connection_string(connection_string) container_client = blob_service_client.get_container_client(container_name) # 收集所有JSON文件的路径和最后修改时间 file_metadata = [] for blob in container_client.list_blobs(name_starts_with="target-path/"): if blob.name.endswith(".json"): # 注意:这里的路径要和input_file_name()返回的格式一致(比如挂载后的路径) full_path = f"/mnt/{container_name}/{blob.name}" file_metadata.append({ "file_path": full_path, "last_modified": blob.last_modified }) # 转为Spark元数据DataFrame metadata_df = spark.createDataFrame(file_metadata) # 读取JSON数据并添加文件路径字段 json_df = spark.read.json("/mnt/your-blob-container/target-path/*.json") \ .withColumn("file_path", input_file_name()) # 关联元数据和JSON数据 final_df = json_df.join(metadata_df, on="file_path", how="inner")
方案3:Databricks环境下用dbutils.fs.ls获取元数据
如果是在Databricks中,可以直接使用dbutils.fs.ls快速获取文件元数据,再关联读取的JSON数据:
# 获取目标目录下的所有文件及元数据 file_list = dbutils.fs.ls("/mnt/your-blob-container/target-path/") file_metadata = [] for file in file_list: if file.name.endswith(".json"): file_metadata.append({ "file_path": file.path, "last_modified": file.modificationTime }) # 转为Spark元数据DataFrame metadata_df = spark.createDataFrame(file_metadata) # 读取JSON并关联 json_df = spark.read.json("/mnt/your-blob-container/target-path/*.json") \ .withColumn("file_path", input_file_name()) final_df = json_df.join(metadata_df, on="file_path", how="inner")
方案对比
- 方案1:无需额外依赖,读取和获取元数据一步完成,性能最优,适合大多数场景。
- 方案2:适合需要批量获取Blob元数据或进行复杂元数据筛选的场景,但需要引入Azure SDK依赖。
- 方案3:Databricks专属,操作简单,但仅适用于Databricks环境。
内容的提问来源于stack exchange,提问作者mytabi
相关产品推荐
相关产品推荐

