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

运行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流会在检查点目录记录已处理文件,直接读取偏移量文件即可追踪:

  1. 定位管道检查点路径:默认在管道配置的存储位置下的checkpoints/<view_name>/目录,格式如下:
    abfss://<container>@<storage_account>.dfs.core.windows.net/<pipeline_storage_path>/checkpoints/lmax_ns_data_view/

  2. 读取并解析检查点中的文件信息:

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 22:41:15