Auto Loader 管道需要主动监控,以便在积压持续增长、架构漂移、数据损坏和流处理停滞等问题影响下游消费者之前发现它们。 本页介绍如何监视关键指标、查询文件级状态、生成可观测性仪表板以及排查常见问题。
有关生产配置详细信息,请参阅 为生产工作负荷配置自动加载程序。 有关配置最佳做法,请参阅 自动加载程序最佳做法。
Prerequisites
此页面上的多个监控工作流依赖 cloud_files_state() 来监测每个文件的摄取状态,包括积压查询、延迟计算和架构漂移检测。
cloud_files_state() 是一个表值函数,它返回自动加载程序检查点的文件级引入状态。 默认情况下,并非所有字段都可用。 可用性取决于 Databricks Runtime 版本和配置:
-
Databricks Runtime 18.2 及更高版本:
discovery_time,processed_time并且commit_time自动可用。 在 Databricks Runtime 16.4-18.1 上,仅在启用cloudFiles.cleanSource时,这些字段才可用。 -
启用
cloudFiles.cleanSource的 Databricks Runtime 16.4 及更高版本:archive_time、archive_mode和move_location可用。
启用 cloudFiles.cleanSource 会产生一些性能开销。 在生产环境中启用它之前,先在预生产环境中针对您的工作负载进行基准测试。
此外:
- 使用
_metadata列标注导入的数据。 至少捕获file_path和file_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 |
MOVE、 DELETE或 NULL (需要 cloudFiles.cleanSource) |
move_location |
STRING |
当 cloudFiles.cleanSource 为 MOVE 时的目标路径 |
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())
注释
侦听器中的处理逻辑可能会降低查询处理速度。 限制监听器回调中的计算,并避免在其中进行同步的外部写入;而应异步发送轻量级遥测数据,或将指标交由单独的作业持久化。
使用源进度中的
numInputRows、inputRowsPerSecond和processedRowsPerSecond计算吞吐量,即每个批次的每秒文件数和每秒行数。要计算摄取延迟,请比较
create_time中的commit_time和cloud_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_time 和 commit_time 未填充 |
在低于 18.2 的 Databricks Runtime 上运行且未使用 cleanSource |
升级到 Databricks Runtime 18.2 及更高版本或在 Databricks Runtime 16.4-18.1 上启用cloudFiles.cleanSource |
有关其他故障排除,请参阅 自动加载程序常见问题解答。