监视和观察自动加载程序

Auto Loader 管道需要主动监控,以便在积压持续增长、架构漂移、数据损坏和流处理停滞等问题影响下游消费者之前发现它们。 本页介绍如何监视关键指标、查询文件级状态、生成可观测性仪表板以及排查常见问题。

有关生产配置详细信息,请参阅 为生产工作负荷配置自动加载程序。 有关配置最佳做法,请参阅 自动加载程序最佳做法

Prerequisites

此页面上的多个监控工作流依赖 cloud_files_state() 来监测每个文件的摄取状态,包括积压查询、延迟计算和架构漂移检测。 cloud_files_state() 是一个表值函数,它返回自动加载程序检查点的文件级引入状态。 默认情况下,并非所有字段都可用。 可用性取决于 Databricks Runtime 版本和配置:

  • Databricks Runtime 18.2 及更高版本discovery_timeprocessed_time并且 commit_time 自动可用。 在 Databricks Runtime 16.4-18.1 上,仅在启用 cloudFiles.cleanSource 时,这些字段才可用。
  • 启用 cloudFiles.cleanSource 的 Databricks Runtime 16.4 及更高版本archive_timearchive_modemove_location可用。

启用 cloudFiles.cleanSource 会产生一些性能开销。 在生产环境中启用它之前,先在预生产环境中针对您的工作负载进行基准测试。

此外:

  • 使用 _metadata 列标注导入的数据。 至少捕获 file_pathfile_modification_time。 请参阅文件元数据列
  • 启用 _rescued_data_corrupt_record 列。

自动加载器关键指标

下表汇总了用于监视自动加载程序管道的最重要指标。 这些指标可从 StreamingQueryListener 进度事件中获取,其中 Auto Loader 特有的值在各源的 metrics 映射中提供。

Metric 它告诉你什么
numFilesOutstanding 等待处理的积压工作中的文件数
numBytesOutstanding 积压文件的大小(以字节为单位)
approximateQueueSize 云队列深度(仅限文件通知模式)
numInputRows 每批次处理的行数
inputRowsPerSecond 数据到达率
processedRowsPerSecond 处理吞吐量
durationMs 故障 每批次的时间消耗分布

需要关注的内容

以下模式表明你的管道可能需要关注。

  • 增长 numFilesOutstanding:积压工作正在增加。 您的数据管道处理速度已跟不上传入的数据。
  • processedRowsPerSecond < inputRowsPerSecond:管道处理数据的速度比到达速度慢。
  • 大型 durationMs.latestOffset:发现文件的速度较慢。 请考虑切换到文件事件机制。
  • durationMs.addBatch:数据处理速度缓慢。 请考虑扩展计算资源或优化转换操作。

有关完整的指标参考,请参阅 Auto Loader 源指标

使用 cloud_files_state 查询文件级状态

cloud_files_state()表值函数提供有关自动加载程序发现的每个文件的详细信息。 以下字段可用。 标记为需要 Databricks Runtime 16.4 及更高版本或 18.2 及更高版本的字段仅在 先决条件中所述的条件下填充。

领域 类型 说明
path STRING 文件的路径
size BIGINT 文件的大小(以字节为单位)
create_time TIMESTAMP 文件创建时间
discovery_time TIMESTAMP 自动加载程序发现文件时(Databricks Runtime 16.4 及更高版本)
processed_time TIMESTAMP 自动加载程序处理文件时(Databricks Runtime 16.4 及更高版本)
commit_time TIMESTAMP 文件提交到检查点时(Databricks Runtime 16.4 及更高版本)
archive_time TIMESTAMP 文件归档的时间(需要 cloudFiles.cleanSource
archive_mode STRING MOVEDELETENULL (需要 cloudFiles.cleanSource
move_location STRING cloudFiles.cleanSourceMOVE 时的目标路径
ingestion_state STRING 当前文件引入状态

调查文件引入状态

以下查询涵盖常见的诊断方案。

查找所有未处理的文件(当前积压工作):

SELECT * FROM cloud_files_state('path/to/checkpoint')
WHERE ingestion_state != 'COMMITTED';

计算平均引入延迟(从创建文件到提交的时间):

SELECT avg(unix_timestamp(commit_time) - unix_timestamp(create_time)) AS avg_latency_seconds
FROM cloud_files_state('path/to/checkpoint')
WHERE commit_time IS NOT NULL AND create_time IS NOT NULL;

查找损坏或跳过的文件:

SELECT path, ingestion_state, size, create_time
FROM cloud_files_state('path/to/checkpoint')
WHERE ingestion_state LIKE 'SKIPPED%';

跟踪存档进度(需要 cloudFiles.cleanSource):

SELECT archive_mode, count(*) AS file_count
FROM cloud_files_state('path/to/checkpoint')
GROUP BY archive_mode;

查找从发现到提交延迟较高的文件,以识别瓶颈:

SELECT
  path,
  size,
  unix_timestamp(commit_time) - unix_timestamp(discovery_time) AS processing_latency_seconds,
  unix_timestamp(commit_time) - unix_timestamp(create_time) AS end_to_end_latency_seconds
FROM cloud_files_state('path/to/checkpoint')
WHERE commit_time IS NOT NULL
ORDER BY end_to_end_latency_seconds DESC
LIMIT 20;

有关完整的 SQL 参考,请参阅 cloud_files_state 表值函数

在 Lakeflow 管道中监视自动加载程序

Databricks 建议将 Lakeflow 管道用于生产自动加载程序管道。 要利用其内置监控功能:

  • 将 Lakeflow 管道事件日志存储在 Delta 表中,以便查询它以获取可观测性数据。 通过流水线的高级设置或 API 进行此项配置。 有关详细信息,请参阅 管道事件日志

  • 为实现可观测性而设计你的流水线。 Lakeflow 管道中结构良好的 Auto Loader 管道包括一个 {table}_source 视图(Auto Loader 源定义)、一个 {table}_bronze 流式表(使用 _rescued_data_corrupt_record 列摄取原始数据)、一个 corrupt_records_sink(用于隔离包含无法解析数据的行),以及一个供下游使用的 {table} 干净视图。

  • 在你的青铜流式表上设置预期,以监控架构漂移和数据损坏。 _rescued_data IS NULL 检测意外的架构更改并 _corrupt_record IS NULL 检测不可分析的数据。 Lakeflow 管道会在数据到达时评估这些预期,并生成可观测性轨迹。 你可以将预期配置为告警、丢弃行或使管道失败。

为您的管道创建 event_log_raw 视图后,使用以下查询来查看 Auto Loader 特定指标。

监视每个流的引入吞吐量:

SELECT
  origin.flow_name,
  origin.update_id,
  timestamp,
  TRY_CAST(details:flow_progress.metrics.num_output_rows AS BIGINT) AS rows_written
FROM event_log_raw
WHERE event_type = 'flow_progress'
ORDER BY timestamp DESC;

按流监控数据积压:

SELECT
  origin.flow_name,
  timestamp,
  DOUBLE(details:flow_progress.metrics.backlog_bytes) AS backlog_bytes
FROM event_log_raw
WHERE event_type = 'flow_progress'
  AND details:flow_progress.metrics.backlog_bytes IS NOT NULL
ORDER BY timestamp DESC;

汇总预期冲突以检测架构偏移和损坏数据:

SELECT
  origin.flow_name,
  explode(from_json(
    details:flow_progress.data_quality.expectations,
    'array<struct<name:string, dataset:string, passed_records:bigint, failed_records:bigint>>'
  )) AS expectation
FROM event_log_raw
WHERE event_type = 'flow_progress'
  AND details:flow_progress.data_quality.expectations IS NOT NULL;

有关 Lakeflow 管道总体监控指南,请参阅 监控管道管道事件日志

使用结构化流式处理监控自动加载程序

在 Lakeflow 管道之外运行 Auto Loader 时,请使用以下 Structured Streaming 监控方法。

  • 实现一个 StreamingQueryListener,通过读取 source.metrics 捕获每个批次中的 Auto Loader 特有指标。
from pyspark.sql.streaming import StreamingQueryListener

class AutoLoaderMonitor(StreamingQueryListener):
    def onQueryStarted(self, event):
        pass

    def onQueryProgress(self, event):
        for source in event.progress.sources:
            if "CloudFilesSource" in source.description:
                metrics = source.metrics
                files_outstanding = metrics.get("numFilesOutstanding", "0")
                bytes_outstanding = metrics.get("numBytesOutstanding", "0")
                rows_per_sec = source.processedRowsPerSecond
                # Push metrics to your monitoring system (for example, write to a Delta table)

    def onQueryIdle(self, event):
        pass

    def onQueryTerminated(self, event):
        pass

spark.streams.addListener(AutoLoaderMonitor())

注释

侦听器中的处理逻辑可能会降低查询处理速度。 限制监听器回调中的计算,并避免在其中进行同步的外部写入;而应异步发送轻量级遥测数据,或将指标交由单独的作业持久化。

  • 使用源进度中的 numInputRowsinputRowsPerSecondprocessedRowsPerSecond 计算吞吐量,即每个批次的每秒文件数和每秒行数。

  • 要计算摄取延迟,请比较 create_time 中的 commit_timecloud_files_state(),以计算端到端延迟。 对于处理延迟,请使用durationMs细分(例如,latestOffsetaddBatch和其他报告的批处理阶段)来确定哪个阶段是瓶颈。

  • 使用 df.observe() 直接在流式 DataFrame 上定义内联数据质量指标。 在 StreamingQueryListener 下的 observedMetrics 进度事件中可以看到指标。

from pyspark.sql.functions import count, lit, col

observed_df = df.observe(
    "auto_loader_quality",
    count(lit(1)).alias("total_rows"),
    count(col("_rescued_data")).alias("rescued_rows"),
    count(col("_corrupt_record")).alias("corrupt_rows")
)
  • 使用 .queryName() 为每个流分配唯一名称,以便更轻松地区分 Spark UI 的 Streaming 选项卡和监控仪表板中的 Auto Loader 流。

有关结构化流处理监控的完整参考信息,请参阅在 Azure Databricks 上监控结构化流处理查询

生成可观测性仪表板

合并来自多个源的数据,为自动加载程序管道构建全面的可观测性仪表板。 此表显示可用于构建可观测性仪表板的一些建议源。

数据源 可观测性数据
cloud_files_state() 文件级引入状态:每个文件的发现、处理、提交和存档时间戳
Lakeflow 管道事件日志 管道运行历史记录、每批流指标和数据质量预期结果
管道输出表 每个引入表的行数和写入数据量

然后,可以将可观测性数据聚合到专用表中,作为仪表板和警报的基础:

  • 汇总一段时间内的管道运行状态(成功或失败),这些状态派生自 event_type = 'update_progress' 事件。
  • 文件引入聚合指标(积压大小、吞吐量、每批次延迟),源自 cloud_files_state()event_type = 'flow_progress' 事件。
  • 使用从事件日志中的 num_output_rows 派生的每个表的行数和数据量来生成表统计信息。
  • 从每次更新的详细错误日志和预期违规记录中收集调试信息,这些信息源自填充了 event_type = 'flow_progress'data_quality 事件。

这些聚合表可以为 AI/BI 仪表板和 SQL 警报提供支持。 建议的仪表板面板包括管道运行状态时间线、引入积压工作趋势、吞吐量趋势、引入延迟分布、数据质量指标、架构演变事件和文件存档状态。

监视架构演变事件

使用以下方法在架构变更发生时检测它们。

  • 预期冲突计数中的 _rescued_data 非 NULL 值表示架构偏移。 针对 failed_records > 0 预期查询事件日志中的 no rescued data
  • 对已配置的 _schemas 中的 cloudFiles.schemaLocation 目录所做的更改(或者,当未单独设置架构位置时,仅对检查点中的该目录所做的更改)表明已发生架构演变。 可以通过单独的监控作业轮询此目录。
  • 不要将同一流名称中先出现 onQueryTerminated 事件、随后出现 onQueryStarted 这一情况,本身视为模式演进的充分证据。 流会因多种原因而重启(例如集群重启、代码部署和临时性存储错误)。 在断定发生了架构演进之前,先将重启与独立信号——_schemas 目录更改或 _rescued_data 预期违背——结合起来分析。
  • 使用 _metadata.file_path 来识别哪些文件引入了架构更改。 将cloud_files_state()此项与path字段联接,以将架构更改与特定文件和批处理相关联。

使用此示例查询通过预期冲突检测最近的架构偏移:

SELECT
  timestamp,
  origin.flow_name,
  exp.name AS expectation_name,
  exp.failed_records
FROM (
  SELECT
    timestamp,
    origin,
    explode(from_json(
      details:flow_progress.data_quality.expectations,
      'array<struct<name:string, dataset:string, passed_records:bigint, failed_records:bigint>>'
    )) AS exp
  FROM event_log_raw
  WHERE event_type = 'flow_progress'
    AND details:flow_progress.data_quality.expectations IS NOT NULL
)
WHERE exp.name = '<rescued-data expectation name>'
  AND exp.failed_records > 0
ORDER BY timestamp DESC;

为常见问题设置警报

使用 Databricks SQL 警报或管道通知在影响下游使用者之前检测问题。

以下 SQL 检测到日益增多的积压工作,并可用作 Databricks SQL 警报的基础。 计划它定期运行(例如每 5 分钟运行一次),并在结果为非空时发出警报。

-- Alert when backlog exceeds threshold or trends upward across recent batches
WITH recent_backlog AS (
  SELECT
    origin.flow_name,
    timestamp,
    DOUBLE(details:flow_progress.metrics.backlog_bytes) AS backlog_bytes,
    ROW_NUMBER() OVER (PARTITION BY origin.flow_name ORDER BY timestamp DESC) AS rn
  FROM event_log_raw
  WHERE event_type = 'flow_progress'
    AND details:flow_progress.metrics.backlog_bytes IS NOT NULL
)
SELECT flow_name, backlog_bytes, timestamp
FROM recent_backlog
WHERE rn = 1
  AND backlog_bytes > 1073741824  -- alert when backlog exceeds 1 GB

下表汇总了建议的警报条件:

要检测的内容 如何检测 何时发出警报
积压工作日益增多 numFilesOutstanding 呈上升趋势 在多个批次中持续增加
流停滞 无进度事件 N 分钟无事件(基于预期的触发器间隔)
高摄取延迟 commit_time - create_time 超出 SLA 阈值
数据质量下降 预期失败率 未满足预期的行占比增加
架构演变事件 _rescued_data IS NOT NULL 预期违例计数中的任何非 NULL 值
文件发现速度缓慢 durationMs.latestOffset 明显高于基线

排查常见问题

下表描述了常见的自动加载程序管道问题、其可能原因以及解决这些问题的推荐操作。

问题 可能的原因 建议的操作
积压工作的速度比处理更快 计算资源不足、数据倾斜或速率限制被节流 扩容计算资源,通过 Spark UI 检查倾斜情况,并复查 maxFilesPerTrigger 设置以控制批次大小
未发现文件 文件事件配置错误、权限问题或流未在 7 天内运行 验证外部位置权限,检查 Unity 目录 UI 中的文件事件设置,并确保流至少每 7 天运行一次,以避免 RocksDB 状态过期
流启动时间过长 大型检查点状态下载(RocksDB) 升级到 Databricks Runtime 15.3 及更高版本以执行异步状态加载,从而将启动时间缩短约 90%
重复文件处理 激进的 cloudFiles.maxFileAge 设置或检查点损坏 采用保守的 maxFileAge(至少 90 天),验证检查点完整性,并避免对检查点存储使用生命周期策略
导致管道重启的架构演变 频繁或不兼容的架构更改 查看 schemaEvolutionMode,如需进行类型提升,请切换到 addNewColumnsWithTypeWidening,或使用变体类型来支持高度动态的架构
接收器中累积损坏数据 源数据质量问题 检查 _corrupt_record 隔离接收器中的模式,审查源数据生成过程,并考虑添加上游验证
discovery_timecommit_time 未填充 在低于 18.2 的 Databricks Runtime 上运行且未使用 cleanSource 升级到 Databricks Runtime 18.2 及更高版本或在 Databricks Runtime 16.4-18.1 上启用cloudFiles.cleanSource

有关其他故障排除,请参阅 自动加载程序常见问题解答