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

如何查看Databricks Auto Loader处理的文件名及验证读取情况

排查Auto Loader文件处理及时间戳问题

核心问题分析

你当前的Auto Loader代码未设置checkpointLocation,导致流处理状态无法持久化,无法使用官方文档的cloud_files_state查询已处理文件。同时Delta表无2023年记录,需分步骤验证:


1. 临时验证已处理文件(无需重启流)

在流处理的DataFrame中添加源文件元数据列,写入Delta表后即可查询哪些文件被处理:

  • 在withColumn转换环节加入以下代码,把文件名和路径写入Delta表:
    from pyspark.sql.functions import col, input_file_name, input_file_path
    
    # 添加源文件元数据列
    df_stream_in = df_stream_in.withColumn("source_file_name", input_file_name()) \
                               .withColumn("source_file_path", input_file_path())
    
  • 写入Delta表后,执行以下查询:
    -- 查询已处理的文件名列表
    SELECT DISTINCT source_file_name FROM your_delta_table;
    
    -- 查询最大时间戳
    SELECT MAX(`UTC Time Stamp`) FROM your_delta_table;
    

2. 持久化流状态(必须配置,避免重复处理/丢失状态)

修改流写入代码,添加checkpointLocation(指定DBFS或云存储路径),这样Auto Loader会持久化已处理文件的状态:

df_stream_in.writeStream \
            .format("delta") \
            .option("checkpointLocation", "/dbfs/mnt/your-checkpoint-path")  # 替换为实际路径
            .table("your_delta_table")

配置完成后,即可用官方方法查询Auto Loader状态:

SELECT * FROM cloud_files_state('/dbfs/mnt/your-checkpoint-path');

结果中的processedFiles字段会列出所有已处理的文件路径及处理时间。

3. 验证源文件是否可正常读取

先通过静态读取确认存储账户中的2023年文件确实包含有效数据:

# 读取最新的2023年文件(可根据路径过滤)
latest_2023_files = spark.read.format("csv") \
                          .option("header", "true") \
                          .schema(dataset_schema) \
                          .load(f"{file_location}/**/2023/**")  # 替换为实际路径过滤规则

# 查看数据及最大时间戳
display(latest_2023_files)
latest_2023_files.agg(max(col("UTC Time Stamp"))).show()

如果静态读取能获取到2023年数据,说明问题出在流处理状态或配置上;如果静态读也没有数据,需检查文件路径、权限或文件格式。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 21:33:14