如何查看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
相关产品推荐
相关产品推荐

