运行Delta Live Tables工作流时,如何查看管道已处理的文件?
解决方案
方法1:在数据中嵌入源文件路径,直接查询已处理文件
修改DLT视图代码,添加源文件路径列,让最终的Bronze表保留文件追踪信息:
@dlt.view( name="lmax_ns_data_view", comment="Raw one nanosecond 3-month snapshot from LMAX" ) def name_of_view(): import pyspark.sql.functions as F from pyspark.sql.types import StructType, StructField, IntegerType, LongType, StringType, DoubleType # Set the storage path storage_path = "<storage_path>" # Set the schema string_columns = ["str1", "st2"] float_columns = ["fl1", "fl1"] int_32_columns = ["int32_col"] int_64_columns = ["col1", "col2", "col3"] data_schema = StructType( [StructField(col, IntegerType(), True) for col in int_32_columns] + [StructField(col, LongType(), True) for col in int_64_columns] + [StructField(col, StringType(), True) for col in string_columns] + [StructField(col, DoubleType(), True) for col in float_columns] ) df = spark.readStream.format("cloudFiles") \ .option("cloudFiles.format", "parquet") \ .option("cloudFiles.includeExistingFiles", "true") # 确保首次运行处理所有现有文件 .schema(data_schema) \ .load(storage_path) \ # 添加源文件路径列 df = df.withColumn("source_file", F.input_file_name()) return df
修改完成后,可通过以下SQL查询Bronze表获取处理进度:
-- 查看所有已处理的唯一文件路径 SELECT DISTINCT source_file FROM <your_catalog>.<your_schema>.name_of_streaming_table; -- 统计已处理文件数量 SELECT COUNT(DISTINCT source_file) AS processed_file_count FROM <your_catalog>.<your_schema>.name_of_streaming_table;
方法2:解析DLT流的检查点文件
DLT的Cloud Files流会在检查点目录记录已处理文件,直接读取偏移量文件即可追踪:
定位管道检查点路径:默认在管道配置的存储位置下的
checkpoints/<view_name>/目录,格式如下:abfss://<container>@<storage_account>.dfs.core.windows.net/<pipeline_storage_path>/checkpoints/lmax_ns_data_view/读取并解析检查点中的文件信息:
checkpoint_offset_path = "<your_checkpoint_path>/source/offsets" offset_df = spark.read.json(checkpoint_offset_path) display(offset_df.select(F.explode("files").alias("processed_file_path")))
方法3:查询DLT事件日志的Cloud Files处理事件
过滤事件日志中的特定事件类型,获取批次处理的文件明细:
SELECT timestamp, details.files AS processed_files, details.numFiles AS processed_file_count FROM table(dlt.event_log('<your_pipeline_id>')) WHERE event_type = 'cloud_files_processed' ORDER BY timestamp DESC
该查询会按时间倒序返回每次批次处理的文件列表和数量,实时追踪处理进度。
内容的提问来源于stack exchange,提问作者Sorath
相关产品推荐
相关产品推荐

