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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 17:25:36