Skip to content

Professional 学習教材 ④ モニタリング・アラート(配点 10%)

Databricks 認定データエンジニア Professional(上位資格)試験のドメイン④「モニタリング・アラート」の日本語自習教材です。 本番運用(プロダクション)観点で、ジョブ/パイプライン/ストリーム/データ品質を「どう監視し」「どう気付き(アラート)」「どう通知するか」を、システムテーブル・データ品質監視(旧 Lakehouse Monitoring)・Structured Streaming メトリクス・SQL アラート/ジョブ通知の 4 本柱で整理しています。 このファイル単体で学習が完結するよう、公式ドキュメントの用語・スキーマ・設定項目・注意点を可能な限りそのまま収録しています。

参照した公式ページ(すべて WebFetch で実読、2026-07 時点):

注記: 指定 URL のうち data-quality-monitoring/lakehouse-monitoring/monitor-output は現在の正規 URL(Unity Catalog 配下)へリダイレクトされました。sql/user/alerts/ は 404 ではなくランディングページとして取得でき、詳細は create サブページで補いました。捏造は行わず、取得できた内容のみを収録しています。


1. このドメインの概要

Professional 試験のこのドメインは、本番ワークロードを「観測(オブザーバビリティ)」し、問題を「検知」し、担当者へ「通知」する一連の運用を問います。Databricks では監視のレイヤーが役割ごとに分かれており、どの目的にどのツールを使うかを正しく選べることが重要です。

監視したいもの主に使うツール出力・確認先
コスト・監査・ジョブ/パイプラインの履歴(アカウント全体)システムテーブルsystem カタログ)SQL でクエリ、ダッシュボード
テーブルのデータ品質(鮮度・完全性・統計・ドリフト)データ品質監視(異常検出+データプロファイル、旧 Lakehouse Monitoring)メトリックテーブル、ダッシュボード
Structured Streaming クエリの健全性(レート・遅延・状態)StreamingQueryListener / Spark UI Streaming タブメトリクス JSON、外部システムへ push
条件に基づく能動的な検知と通知Databricks SQL アラート / ジョブ通知Email / Slack / Teams / PagerDuty / Webhook

学習のコツは 3 つの視点で整理することです。

  1. プル型(履歴分析): システムテーブルとデータ品質監視のメトリックテーブルは「後から SQL でクエリして傾向を見る」もの。リアルタイムではなく 1 日を通して更新される(システムテーブルはリアルタイム監視をサポートしない)。
  2. ストリーミング型(リアルタイム観測): StreamingQueryListener と Spark UI は「今流れているストリームの入力/処理レート・遅延・状態」をイベント単位で観測する。
  3. プッシュ型(能動通知): SQL アラート・ジョブ通知は「条件が満たされたら人に知らせる」。通知先は管理者が定義した通知先(destination)を再利用する。

2. 重要用語集

用語(日本語)English説明
システムテーブルSystem tablessystem カタログにある、アカウントの運用データの Databricks ホスト型分析ストア。コスト監視・監査・コンピュート/ジョブ追跡・データ/AI 監視に使う。Unity Catalog で管理され、読み取り専用。
情報スキーマInformation schemasystem.information_schema。他のシステムテーブルとは動作が異なるメタデータ用スキーマ。
監査ログ(テーブル)Audit logssystem.access.audit。リージョン内ワークスペースの全監査イベント記録(パブリックプレビュー、保持 365 日)。ワークスペースレベルはリージョン、アカウントレベルはグローバルデータ。
課金対象使用状況Billable usagesystem.billing.usage。アカウント全体の課金対象使用量レコード(グローバル、保持 365 日、無料)。
料金テーブルPricing / list pricessystem.billing.list_prices。SKU 価格の履歴ログ。価格変更ごとに行が追加される(グローバル、無料)。usage と結合して USD コストを算出。
ジョブ(LakeFlow)システムテーブルJobs / LakeFlow system tablessystem.lakeflow.*。ジョブ・タスク・実行タイムライン・パイプライン等を追跡(旧称 workflow スキーマ)。
ジョブテーブルjobssystem.lakeflow.jobs。作成された全ジョブを追跡する SCD2(緩やかに変化するディメンション)テーブル。
ジョブ実行タイムラインjob_run_timelinesystem.lakeflow.job_run_timeline。ジョブ実行の開始/終了時刻・結果状態・終了コードを追跡(イミュータブル)。
ジョブタスク実行タイムラインjob_task_run_timelinesystem.lakeflow.job_task_run_timeline。タスク実行ごとの期間・使用コンピュートを追跡。マルチタスク/非レガシジョブのコスト分析に必須。
SCD2 テーブルSlowly changing dimension (type 2)行が変わると新しい行を出力し前行を論理的に置換。jobs/job_tasks/pipelines が該当。各エンティティの最新行は 365 日超でも保持される。
データ品質監視Data quality monitoringUnity Catalog のデータ資産の品質を確保する機能群。異常検出データプロファイルから成る。
Lakehouse Monitoring(旧称)Lakehouse Monitoringデータプロファイル機能の旧称。現在は「データ品質監視 > データプロファイル」に統合。
異常検出Anomaly detectionワンクリックでスキーマ内全テーブルの鮮度完全性を自動監視。インテリジェントスキャンで重要テーブルを優先。
鮮度Freshnessテーブルが最近更新されたか。コミット履歴から次コミット時刻を予測し、異常に遅れると「古い(stale)」とマーク。
完全性Completeness過去 24 時間に書き込まれるべき行数。履歴から予想範囲を予測し、下限を下回ると「不完全」とマーク。
データプロファイルData profilingテーブル/推論ログの概要統計・分布・モデル性能を時系列でキャプチャする定量指標。
プロファイルの分析タイプAnalysis typeSnapshot(スナップショット)/TimeSeries(時系列)/InferenceLog(推論ログ)の 3 種。時間ウィンドウの決まり方が異なる。
メトリックテーブルMetric tablesプロファイルが作成・更新する 2 つのテーブル: プロファイルメトリックテーブルとドリフトメトリックテーブル。
プロファイルメトリックテーブルProfile metrics table{output_schema}.{table_name}_profile_metrics。各列×時間ウィンドウ×スライス×グループの概要統計(count, avg, 分位点等)。
ドリフトメトリックテーブルDrift metrics table{output_schema}.{table_name}_drift_metrics。分布変化を追跡。連続ドリフトとベースラインドリフトを計算。
連続ドリフトConsecutive drift現ウィンドウを直前の時間ウィンドウと比較。集計後に連続ウィンドウが存在する場合のみ計算。
ベースラインドリフトBaseline driftベースラインテーブルの分布と比較。ベースラインテーブルを与えた場合のみ計算。
スライスSliceslicing_exprs で定義するデータの部分集合。slice_key/slice_value 列で識別。NULL/NULL はテーブル全体。
時間ウィンドウTime windowメトリックを計算する時間間隔。granularity(粒度)で決まる struct<start, end>
StreamingQueryListenerStreamingQueryListenerStructured Streaming のメトリクスを外部へ push するリスナー。onQueryStarted/Progress/Idle/Terminated を実装。DBR 11.3 LTS 以降。
StreamingQueryProgressStreamingQueryProgress各マイクロバッチ完了時のメトリクス(numInputRows, レート, durationMs 等)を持つオブジェクト。
入力レートinputRowsPerSecondデータが到着する速度(全ソース集計)。
処理レートprocessedRowsPerSecondSpark がデータを処理する速度(全ソース集計)。
バックプレッシャー/バックログBackpressure / backlog入力レート > 処理レートで未処理データが溜まる状態。Delta/AutoLoader の numBytesOutstanding/numFilesOutstanding、Kafka の *OffsetsBehindLatest で観測。
ウォーターマークWatermark遅延データの状態トリミング基準時刻。eventTime.watermark
ステート演算子State operatorsステートフル処理(集計・重複除去・ストリーム結合等)の状態メトリクス(numRowsTotal 等)。
observe / 監視可能メトリックObservable metricsdf.observe("name", <集計>) で定義する任意の名前付き集計。observedMetrics から取得しデータ品質チェックに使う。
Databricks SQL アラートDatabricks SQL Alertsスケジュールで SQL を実行し、結果が条件を満たしたら通知する機能。状態は OK/TRIGGERED/ERROR
アラート条件Alert condition値(最初の行 or 集計 SUM/AVG 等)+列+演算子(> 等)+しきい値(静的値等)。
通知(Email / Slack / Webhook)Notifications条件成立時の通知先。Email・Slack・Microsoft Teams・PagerDuty・Webhook。
通知先(宛先)Notification destination管理者が作成する再利用可能な通知先。Slack/Teams/PagerDuty/Webhook/Email。ジョブと SQL アラートで利用可。
再通知間隔Notification frequencyOK に戻るまで通知を繰り返す間隔。SQL アラートの通知設定。
ジョブ通知Job notificationsジョブ/タスクの 開始・成功・失敗・期間超過・ストリーミングバックログ超過を通知。
ヘルスルールHealth rulesジョブに定義する健全性ルール(期間・ストリーミングバックログのしきい値等)。jobs.health_rules

3. 詳細解説

3-1. システムテーブルによる監視(種類・用途・課金/ジョブ/監査)

システムテーブルとは

  • system カタログにある、アカウント全体の運用データの Databricks ホスト型分析ストア。コスト監視・セキュリティ監査・コンピュート/ジョブ性能追跡・データ/AI 監視に使う。
  • アカウント全体の履歴監視が用途。system カタログはすべての Unity Catalog メタストアに含まれ、accessbillinglakeflowcomputestoragequery などのスキーマを持つ。
  • 読み取り専用(変更不可)。
  • system.information_schema は他のシステムテーブルと動作が異なる(メタデータ用)。

要件と有効化

  • システムテーブルにアクセスするには、ワークスペースが Unity Catalog 有効である必要がある。
  • Unity Catalog によって管理されるため、アカウントに少なくとも 1 つの UC 有効ワークスペースが必要。データはアカウント内全ワークスペース分を含むが、アクセスできるのは UC 有効ワークスペースのみ。
  • メタストアが UC 特権モデル バージョン 1.0 上にある必要がある。

アクセス権限

  • アクセスは Unity Catalog で管理。アカウント管理者かつメタストア管理者は既定でアクセス可能。
  • 他のユーザーに許可するには、管理者が system カタログへの USE CATALOG、システムスキーマへの USE SCHEMA、システムスキーマへの SELECT を付与する。

データの場所・保持・リージョン

  • データはメタストアと同じリージョンの Databricks ホスト型ストレージに格納され、OpenSharing(Delta Sharing) で安全に共有される。
  • 各テーブルに無料の保持期間がある(多くは 365 日、一部 30/90/180 日や「不定」)。
  • 同じクラウドリージョン内の全ワークスペースの運用データを含む。一部テーブルはグローバルデータ(billing.usage, list_prices, audit のアカウントレベル, workspaces 等)。別リージョンのレコードはそのリージョンのワークスペースから見る。

主なシステムテーブル一覧(試験頻出)

テーブルパス用途ストリーミング無料保持データ範囲
system.access.audit監査ログ(全監査イベント)あり365 日WS レベル=リージョン/アカウントレベル=グローバル
system.billing.usage課金対象使用量(コスト監視の中心)あり365 日グローバル
system.billing.list_pricesSKU 価格履歴(USD 換算に使用)なし不定グローバル
system.lakeflow.jobs全ジョブ(SCD2)あり365 日リージョン
system.lakeflow.job_tasks全ジョブタスク(SCD2)あり365 日リージョン
system.lakeflow.job_run_timelineジョブ実行の開始/終了・結果あり365 日リージョン
system.lakeflow.job_task_run_timelineタスク実行の開始/終了・使用コンピュートあり365 日リージョン
system.lakeflow.pipelines / pipeline_update_timelineパイプラインと更新履歴あり365 日リージョン
system.compute.clustersクラスターの構成履歴(SCD2)あり365 日リージョン
system.compute.node_timeline汎用/ジョブコンピュートの使用率メトリックあり90 日リージョン
system.compute.warehouses / warehouse_eventsSQL ウェアハウス構成/イベント(start/stop/scale)あり365 日リージョン
system.access.table_lineage / column_lineageテーブル/列のリネージあり365 日リージョン
system.query.historySQL ウェアハウス/サーバーレスのクエリ履歴なし365 日リージョン
system.data_quality_monitoring.table_resultsデータ品質チェック結果(鮮度・完全性・インシデント)なし不定リージョン

課金使用量テーブルと料金テーブルは無料。パブリックプレビューのテーブルもプレビュー期間中は無料だが将来課金される可能性がある。

運用上の重要な注意点(試験ポイント)

  • リアルタイム監視はサポートされない。データは 1 日を通して更新され、直近イベントが未反映のことがある。
  • スキーマ進化: 既存テーブルに新しい列がいつでも追加され得る(構造体列にも新フィールド)。既存列は変更・削除されない。固定スキーマ依存のクエリは壊れる可能性があるため、別テーブルへ書き出す場合はスキーマの進化を有効化する。
  • 選択性の低いクエリは失敗する: System Table query returned too much data. Please repeat query with more selective predicates. 述語(WHERE の期間絞り込み等)を付ける。
  • ストリーミング利用時: skipChangeCommits=true を設定(削除でストリームが中断しないように)。VACUUM の既定保持は 7 日で、ストリーミングクエリが 7 日以上遅れると中断し得る。追いつくには Trigger.AvailableNow(DBR 18.0 以降)や実行頻度の増加を使う。
  • 非推奨スキーマ: system.operational_datasystem.lineage(空テーブル)。

ジョブ/パイプラインの本番監視(system.lakeflow)

system.lakeflow の 6 テーブル(jobs, job_tasks, job_run_timeline, job_task_run_timeline, pipelines, pipeline_update_timeline)でジョブの成否・所要時間・再試行・コストを分析できる。

  • result_statejob_run_timeline): SUCCEEDED / FAILED / SKIPPED / CANCELLED / TIMED_OUT / ERROR / BLOCKED / NULL(長時間実行の中間スライス行)。実行の終了を表す行にのみ設定される
  • termination_code: SUCCESS / CANCELLED / SKIPPED / DRIVER_ERROR / CLUSTER_ERROR / WORKSPACE_RUN_LIMIT_EXCEEDED / MAX_CONCURRENT_RUNS_EXCEEDED / RESOURCE_NOT_FOUND など。失敗原因の切り分けに使う。
  • run_type: JOB_RUN(標準)/SUBMIT_RUN(runs/submit の一回実行)/WORKFLOW_RUN(ノートブックワークフロー、job_run_timeline にのみ記録され jobs/job_tasks には出ない)。
  • タイムラインのスライスロジック: period_start_timeperiod_end_time は 1 行あたり最大 1 時間。1 時間超の実行は複数行に分割され、result_state/termination_code は最終行のみ設定。実行されなかった場合は period_start_time == period_end_time の行になる(理由は termination_code を確認)。
  • 期間列の注意: run_duration_seconds 等の期間列はレガシ単一タスクジョブのみ設定され、マルチタスク/非レガシは 0。マルチタスクは job_task_run_timeline を使う。
  • コスト帰属の注意: system.billing.usageusage_metadata.job_id はジョブコンピュートまたはサーバーレスで実行したジョブにのみ設定。汎用(All-Purpose)コンピュートで実行したジョブは正確なコスト帰属ができない(リソース共有のため)。正確なコスト追跡には専用ジョブコンピュートまたはサーバーレスを推奨。WORKFLOW_RUN のコストは親ノートブックに帰属する。

3-2. データ品質監視(プロファイル/ドリフトメトリクス、メトリックテーブル、ダッシュボード)

全体像

「データ品質監視(Data quality monitoring)」は Unity Catalog のデータ資産の品質を確保する機能群で、2 つの柱がある。

  1. 異常検出(Anomaly detection): ワンクリックでスケーラブルにスキーマ内全テーブルを監視。インテリジェントスキャンが重要テーブルを優先し影響の低いテーブルをスキップしてコストを抑える。履歴データパターンを解析して各テーブルの鮮度完全性を自動評価する。
    • 鮮度: コミット履歴からテーブルごとのモデルを作り次コミット時刻を予測。異常に遅れると「古い」とマーク。
    • 完全性: 過去 24 時間に書き込まれるべき行数を履歴から予測。予想範囲の下限未満なら「不完全」とマーク。
  2. データプロファイル(Data profiling、旧 Lakehouse Monitoring): テーブルの分布やモデル性能の履歴メトリックをキャプチャする定量指標。推論テーブル(モデルの入力+予測)を監視することで GenAI アプリ・ML モデル・モデルサービングの性能も追跡できる。

アカウント管理者・メタストア管理者はガバナンスハブの Data ページで「正常なテーブルの割合」などデータ品質の健全性を確認できる。

データプロファイルが答えられる問い

  • null/ゼロ値の割合はどう変化しているか(データ整合性)。
  • 数値列の 90 パーセンタイル、カテゴリ列の分布は昨日とどう違うか(統計分布)。
  • 現在データと既知のベースライン、または連続する時間枠の間に誤差(ドリフト)があるか。
  • データのスライスごとの分布・ドリフトは。
  • ML モデルの入力・予測が時間とともにどう変わるか。モデル A と B の性能比較。

観測の時間粒度を制御でき、カスタムメトリックも設定できる。監視対象テーブルは変更されず、それらを作るジョブにオーバーヘッドも追加しない。

コスト

  • データ品質監視はサーバーレスコンピュートで実行され、DATA_QUALITY_MONITORING 課金製品のサーバーレス DBU として課金。
  • コストは監視テーブルの数・サイズ・評価頻度に依存。異常検出のインテリジェントスキャンはコスト制御に寄与。

メトリックテーブル(monitor-output の核心)

テーブルにプロファイルを実行すると、2 つのメトリックテーブルが作成・更新される。

  • プロファイルメトリックテーブル {output_schema}.{table_name}_profile_metrics: 各列と、時間ウィンドウ・スライス・グループ化列の各組み合わせの概要統計。InferenceLog ではモデル精度も含む。
  • ドリフトメトリックテーブル {output_schema}.{table_name}_drift_metrics: メトリック分布の変化を追跡。値そのものではなく「変化」を可視化・アラートするのに使う。
    • 連続(CONSECUTIVE)ドリフト: ウィンドウを直前の時間ウィンドウと比較(集計後に連続ウィンドウがある場合のみ)。
    • ベースライン(BASELINE)ドリフト: ベースラインテーブルの分布と比較(ベースラインテーブル提供時のみ)。

{output_schema}output_schema_name で指定したカタログ+スキーマ、{table_name} はプロファイル対象テーブル名。

統計の計算単位(ウィンドウとスライス)

  • メトリックはウィンドウ(時間間隔)ごとに計算。Snapshot 分析ではウィンドウは「更新時刻の単一時点」、TimeSeries/InferenceLog では granularity(粒度)と timestamp_col に基づく。
  • メトリックは常にテーブル全体に対して計算。加えて slicing_exprs を指定するとスライスごとにも計算。
    • slicing_exprs=["col_1", "col_2 > 10"] は「col_2 > 10 に 1 つ」「col_2 <= 10 に 1 つ」「col_1 の一意値ごとに 1 つ」のスライスを生成。テーブル全体は slice_key = NULL, slice_value = NULL
  • InferenceLog ではモデル ID ごとにも計算。

プロファイルメトリックテーブルのスキーマ(抜粋)

グループ化列: window(struct start/end), granularity, model_id_col(InferenceLog のみ), log_type(入力 or ベースライン), slice_key, slice_value, column_name:table はテーブル全体に適用されるメトリックの特別名), data_type, monitor_version

概要統計列(主なもの): count, num_nulls, avg, quantiles(1000 分位点の配列。50% は quantiles[500]), distinct_countapprox_count_distinct のため概数), min, max, stddev, num_zeros, num_nan, min/max/avg_size, min/max/avg_len, frequent_items(上位 100 頻出), median, percent_null, percent_zeros, percent_distinct

モデル精度(InferenceLog かつ label_col と prediction_col 指定時):

  • 分類: accuracy_score, log_lossprediction_proba_col 必要), roc_auc_score, confusion_matrix, precision, recall, f1_score
  • 回帰: mean_squared_error, root_mean_squared_error, mean_average_error, mean_absolute_percentage_error, r2_score
  • 公平性とバイアス(problem_type=classification 時): predictive_parity, predictive_equality, equal_opportunity, statistical_parity

ドリフトメトリックテーブルのスキーマ(抜粋)

追加グループ化列: window_cmp(CONSECUTIVE の比較ウィンドウ), drift_typeBASELINE or CONSECUTIVE)。

ドリフト(差分)列: count_delta, avg_delta, percent_null_delta, percent_zeros_delta, percent_distinct_delta, non_null_columns_delta

分布ドリフト検定(差分は「現ウィンドウ − 比較ウィンドウ」):

  • カテゴリ列のみ: chi_squared_test(カイ二乗検定), tv_distance(全変動距離), l_infinity_distance, js_distance(Jensen–Shannon)。
  • 数値列のみ: ks_test(KS 検定), wasserstein_distance, population_stability_index(PSI)。
    • PSI の解釈: < 0.1 有意な変化なし、< 0.2 中程度、>= 0.2 有意な母集団変化。

時系列/推論プロファイルは作成時刻から30 日を振り返る。この境界により最初の分析ウィンドウは部分的になり得る(最初のウィンドウのみに影響)。

メトリックテーブルへのクエリとダッシュボード

  • メトリックテーブルには直接 SQL でクエリできる(例は §4)。
  • プロファイル実行時にダッシュボードも作成される(data-profiling/monitor-dashboard)。ドリフトテーブルを使えば「特定の値」ではなく「データの変化」を可視化・アラートできる。

3-3. ストリーミングクエリの監視(メトリクス、リスナー、遅延/レート)

組み込み監視: Spark UI

  • Databricks は Spark UI の Streaming タブで Structured Streaming の組み込み監視を提供する。
  • writeStream.queryName("<query-name>") を付けて一意のクエリ名を与えると、Spark UI でどのメトリックがどのストリームのものか区別しやすい。

外部サービスへのメトリクス push: StreamingQueryListener

  • StreamingQueryListener(Apache Spark のストリーミングクエリリスナーインターフェース)でメトリクスを外部サービス(アラート・ダッシュボード用)へ push できる。DBR 11.3 LTS 以降で Python / Scala 対応。
  • 実装するコールバック:
    • onQueryStarted: クエリ開始時(DataStreamWriter.start() と同期。ブロックしないこと)。
    • onQueryProgress: 状態更新時(ストリーミングバッチの最後にのみ配信)。
    • onQueryIdle: ソースにデータがなく新データ待ちの時に配信。
    • onQueryTerminated: 停止時(エラー有無問わず)。
  • 注意: リスナーの処理遅延はクエリの処理速度に大きく影響し得る。処理ロジックを制限し、Kafka のような高速応答システムへの書き込みが推奨。長時間データを処理中は onQueryIdleonQueryProgress も送られないことがあるが、クエリは正常。
  • UC 対応コンピュートの制限: StreamingQueryListener で UC 管理オブジェクトと対話するには DBR 15.1 以降(資格情報使用)または専用アクセスモードが必要。標準アクセスモードの Scala は DBR 16.1 以降が必要。

監視可能メトリック(df.observe)

  • 監視可能メトリックはクエリ(DataFrame)に定義する任意の名前付き集計関数。DataFrame が完了点(バッチ完了 or ストリーミングエポック到達)に達すると、その間に処理したデータのメトリックを含むイベントが生成される。
  • 取得先はモードで異なる:
    • バッチ: QueryExecutionListenerQueryExecution.observedMetrics)。
    • ストリーミング/マイクロバッチ: StreamingQueryListenerStreamingQueryProgress.observedMetrics)。
  • Databricks は**continuous トリガーモードをサポートしない**(マイクロバッチのみ)。
  • 用途例: エラー行の割合を監視し、malformed / cnt > 0.5 でアラートを発火(§4 参照)。

StreamingQueryProgress の主要メトリクス

フィールド意味
id再起動しても保持される一意のクエリ ID。
runId開始/再起動ごとに一意な ID。
nameユーザー指定のクエリ名(未指定なら null)。
batchId処理中バッチの一意 ID。再試行で同 ID が複数回、処理データがなければインクリメントされない。
numInputRowsトリガーで処理したレコード数(全ソース集計)。
inputRowsPerSecond入力レート: データ到着速度(全ソース集計)。
processedRowsPerSecond処理レート: Spark の処理速度(全ソース集計)。
durationMsマイクロバッチ各ステージの所要時間。
eventTimeイベント時刻の統計(avg/max/min/watermark)。
stateOperatorsステートフル処理の状態メトリクス。
sourcesソースごとの情報・メトリクス(バックログ含む)。
sinkシンクの情報。
observedMetricsdf.observe で定義した任意集計。

durationMs の主なステージ: addBatch(マイクロバッチ実行、計画時間を除く), getBatch, latestOffset, queryPlanning, triggerExecution(計画+実行の合計), walCommit(新オフセットのコミット), commitOffsets。処理時間のボトルネック分析に使う。

eventTime: avg/max/min/watermark。ウォーターマークはステートフル集計の状態トリミングに使われる。

バックプレッシャー/遅延(レートとバックログ)

入力レート > 処理レートが続くと未処理データ(バックログ)が溜まる=バックプレッシャー。ソースごとに次で観測する。

  • Delta / Auto Loader ソース: sources.metrics.numBytesOutstanding(未処理ファイルの合計サイズ), numFilesOutstanding(未処理ファイル数)がバックログメトリック。Auto Loader は加えて approximateQueueSize(通知モード時)。
  • Kafka ソース: avgOffsetsBehindLatest / maxOffsetsBehindLatest / minOffsetsBehindLatest(最新オフセットからの遅れ), estimatedTotalBytesBehindLatest(未消費推定バイト)。
    • DBR 17.1 以降はマイクロバッチ完了ごとに最新オフセットを取得するため、常時受信トピックではバックログに小さな 0 以外が出ることがある(正常、遅延ではない)。DBR 17.0 以降はマイクロバッチ開始時に取得し、追いついていれば 0 を返し得る。
  • Kinesis: avgMsBehindLatest / maxMsBehindLatest / minMsBehindLatest
  • PubSub / Pulsar: numRecordsReadyToProcess / sizeOfRecordsReadyToProcess など。

stateOperators(ステートフル処理の観測)

operatorNamestateStoreSave, dedupe, symmetricHashJoin 等)ごとに numRowsTotal, numRowsUpdated, numRowsRemoved, memoryUsedBytes, numRowsDroppedByWatermark(遅すぎて集計に含められなかった行=遅延データの存在を示す), numShufflePartitions, numStateStoreInstances(ストリーム結合はパーティションごとに 4 インスタンス)などを持つ。customMetrics に RocksDB 等のバックエンド固有メトリクスが入る。

テーブル識別子のマッピング(注意点)

  • ストリーミングメトリクスの reservoirId は Delta トランザクションログの一意 ID であり、Unity Catalog の tableId にはマップされない。Delta テーブルの ID は DESCRIBE DETAIL <table-name>id フィールドで確認できる。

3-4. アラートと通知(SQL アラート、ジョブ通知、Email/Slack/Webhook)

Databricks SQL アラート

  • スケジュールで SQL クエリを実行し、結果が定義条件を満たしたら通知する機能。スケジュール実行時に関連クエリが実行され条件が評価される。アラート履歴で過去の評価結果を確認できる。
  • 用途: ビジネス KPI 追跡、データ品質監視(UC データ品質モニター/異常検出と組み合わせ)、コスト傾向監視(サーバーレス課金システムテーブルにアラート)、SQL ウェアハウス/クエリ健全性監視、監査・セキュリティイベント検知、AI エージェント障害検知、LakeFlow ジョブのタスクとして条件チェック実行。
従来(レガシー)アラートとの違い(試験ポイント)
  • クエリの再利用不可: 各アラートは自身のクエリ定義を所有し、既存の保存済み SQL クエリを再利用できない(新エディターで直接作成)。
  • 状態が簡略化: UNKNOWN を廃止。評価は OK / TRIGGERED / ERROR に解決される。
  • 移行中は新旧アラートを並行利用できる。
アラートエディターの構成要素
  1. クエリエディター: アラート対象のクエリを作成・テスト。
  2. コンピューティング: アラートクエリを実行する SQL ウェアハウスを選択。
  3. スケジュール: 定期実行スケジュール(例: 毎時 0 分から 5 分ごと)。Quartz Cron 構文で編集可能。
  4. 共有: 権限(他ユーザーの操作方法)。
  5. 条件: トリガーする値のしきい値。
  6. 通知: 送信先ユーザー/通知先。OK に戻るまで通知を繰り返す頻度を設定可能。
  7. 詳細設定: 特別な値・条件のアラート。
条件の設定
  • チェックする: クエリ結果の列の最初の値、または 1 列全行への集計(SUM / AVERAGE 等)
  • チェックするを選択。
  • 論理演算子(例 >)を選択。
  • しきい値: 既定は静的な値(Static value)(例 4000)。
  • 「テスト条件」でプレビュー可能。
  • 重要: アラートはパラメーターを含むクエリをサポートしない
詳細設定
  • OK 時に通知(Notify on OK): OK 復帰時にも通知。
  • 空の結果状態: クエリが結果を返さない場合に返す特別な状態を設定。
  • テンプレート: 通知テンプレートをカスタマイズ。
通知テンプレート(変数)

カスタマイズしない限り既定テンプレートを使用。標準エディターでは {{VARIABLE_NAME}} で変数参照: ALERT_STATUS, ALERT_CONDITION, ALERT_THRESHOLD, ALERT_COLUMN, ALERT_NAME, ALERT_URL, QUERY_RESULT_TABLE(先頭 100 行の HTML テーブル、HTML 表示は Email 宛先のみ), QUERY_RESULT_VALUE, QUERY_RESULT_ROWS, QUERY_RESULT_COLS。 標準エディターは限定的な HTML タグをサポート(Email のみレンダリング)。Markdown エディターでは @VARIABLE_NAME で参照し、ALERT_NAME/STATUS/CONDITION/THRESHOLD/COLUMN/URLQUERY_RESULT_TABLE をサポート。

アクティブ実行制限(本番運用の注意)
  • アラート実行は他のジョブ種(ノートブックジョブ、パイプライン)と同じワークスペースのアクティブ実行クォータ(既定 250 同時実行)を共有する。
  • 全スロット使用中にスケジュールアラートが発生すると、そのアラート実行はスキップされ評価されない。
  • 対策: スケジュールをずらす01:30, 01:32, 01:34 …)、恒常的に上限に達するならサポートにクォータ増加を依頼。

ジョブ通知(LakeFlow ジョブ)

ジョブ実行および個々のタスクに対し、次のイベントで通知を設定できる:

  • 開始 / 正常完了 / 失敗 / 期間がしきい値超過 / ストリーミングバックログメトリックがしきい値超過

通知先は 1 つ以上のメールアドレスまたはサードパーティ宛先(Slack, Microsoft Teams, PagerDuty, Webhook)。ジョブ/タスクごと、通知イベント種別ごとに最大 3 つのシステム宛先

ジョブレベルとタスクレベル、フィルタ
  • 失敗タスクが再試行されてもジョブレベル通知は送られない。失敗タスクごとに通知したいならタスク通知を使う。
  • 「失敗を伴う成功(Success with failures)」で完了したジョブは成功扱い。通知を受けたいなら「成功」を選ぶ。
  • 期間超過通知を受けるには予想期間(タイムアウト)を設定する必要がある。
  • 通知の追加/編集にはジョブへの CAN MANAGE または IS OWNER 権限が必要。
  • 通知の抑制(フィルタ): 「スキップされた実行の通知をミュート」「キャンセルされた実行の通知をミュート」「最後の再試行まで通知をミュート(タスク、既定 3 回再試行)」。ただしジョブレベルのミュートはタスク通知には効かない。
ストリーミングバックログ通知の挙動
  • 10 分間の平均バックログが定義しきい値を超えると送信。
  • 過剰メッセージ防止のため、送信後 30 分待って次を判断。バックログが高いままなら 30 分間隔で更新を受け取る。
  • ジョブサービスがアクティブなストリーミングクエリを追跡できる必要があるため、ジョブ内で awaitTermination() を使わない。
HTTP Webhook のイベントコード

jobs.on_start / jobs.on_success(成功または「失敗を伴う成功」)/ jobs.on_failure / jobs.on_duration_warning_threshold_exceeded。ペイロードは event_type, workspace_id, run.run_id, job.job_id/name(タスクは task.task_key, run.parent_run_id)を含む JSON。

通知先(Notification destinations)

  • システム通知はワークフローの実行イベント(開始・成功・失敗)を知らせるメッセージ。既定はユーザーの Email だが、管理者が Webhook で代替宛先を構成でき、イベントドリブン統合を作れる。
  • 通知先の管理にはワークスペース管理者である必要がある。構成すると全ユーザーが利用可能。
  • サポートされる宛先: Email / Slack / Webhook / Microsoft Teams / PagerDuty
  • 作成手順: 設定 > ワークスペース管理者 > 通知 > 管理 > +ターゲットの追加 > 種類選択 > 構成 > 作成。
  • セキュリティ: 宛先構成は暗号化保存。宛先ごとに異なる資格情報を推奨(Slack/Teams=URL、PagerDuty=統合キー、Webhook=HTTP 基本認証のユーザー名/パスワード)。侵害時に個別失効できる。
  • ネットワーク要件: HTTPS 必須(信頼された CA の SSL 証明書)。コントロールプレーンとデータプレーン送信 IP の両方を許可リストに追加。送信 IP は 30 日に 1 回更新され得るため ip-ranges.json を定期確認。
  • 重要(依存回避): Slack / Microsoft Teams メッセージの内容・書式は将来変更され得るため、特定スキーマに依存する処理を作らない。厳密なスキーマが必要ならユーザー定義 Webhook を使う。
  • 制限: 通知先を構成できるのは Databricks SQL とジョブのみ。Email 宛先は 1,300 文字制限。Email 以外(Slack/Teams 等)はカスタム本文で HTML 非対応(一部は Markdown 可)。ジョブの Email 宛先はジョブ設定で手動入力が必要。

4. 構文・コード例

4-1. システムテーブルへのアクセス権付与

sql
GRANT USE CATALOG ON CATALOG system TO `data-engineers`;
GRANT USE SCHEMA  ON SCHEMA  system.lakeflow TO `data-engineers`;
GRANT SELECT      ON SCHEMA  system.lakeflow TO `data-engineers`;

4-2. コスト監視(billing.usage × list_prices で USD 換算)

sql
-- SKU 別の日次コスト(USD)。述語で期間を必ず絞る(選択性の低いクエリは失敗する)
SELECT
  u.usage_date,
  u.sku_name,
  SUM(u.usage_quantity)               AS qty,
  SUM(u.usage_quantity * p.pricing.default) AS usd
FROM system.billing.usage u
LEFT JOIN system.billing.list_prices p
  ON u.sku_name = p.sku_name
  AND p.price_start_time <= u.usage_start_time
  AND (p.price_end_time  >= u.usage_start_time OR p.price_end_time IS NULL)
  AND p.currency_code = 'USD'
WHERE u.usage_date >= CURRENT_DATE() - INTERVAL 30 DAYS
GROUP BY ALL
ORDER BY usd DESC;

4-3. ジョブの成否・失敗ジョブの抽出(job_run_timeline)

sql
-- 直近 7 日の日次ジョブ数を結果状態別に集計
SELECT
  workspace_id,
  to_date(period_start_time) AS date,
  result_state,
  COUNT(DISTINCT run_id)     AS job_count
FROM system.lakeflow.job_run_timeline
WHERE period_start_time > CURRENT_TIMESTAMP() - INTERVAL 7 DAYS
  AND result_state IS NOT NULL          -- 中間スライス行(NULL)を除外
GROUP BY ALL;

-- 失敗ジョブと終了コード(原因切り分け)
SELECT job_id, run_id, result_state, termination_code, period_end_time
FROM system.lakeflow.job_run_timeline
WHERE result_state IN ('FAILED','ERROR','TIMED_OUT')
  AND period_start_time > CURRENT_TIMESTAMP() - INTERVAL 7 DAYS;

4-4. ジョブ実行あたりのコスト(job_run_timeline × billing.usage × list_prices)

sql
WITH jobs_usage AS (
  SELECT *, usage_metadata.job_id AS job_id, usage_metadata.job_run_id AS run_id
  FROM system.billing.usage
  WHERE billing_origin_product = 'JOBS'
),
usage_usd AS (
  SELECT j.*, j.usage_quantity * p.default AS usage_usd
  FROM jobs_usage j
  LEFT JOIN system.billing.list_prices p
    ON j.sku_name = p.sku_name
    AND p.price_start_time <= j.usage_start_time
    AND (p.price_end_time >= j.usage_start_time OR p.price_end_time IS NULL)
    AND p.currency_code = 'USD'
)
SELECT u.workspace_id, u.job_id, u.run_id,
       SUM(u.usage_usd) AS usd,
       FIRST(t.result_state, TRUE) AS result_state
FROM usage_usd u
LEFT JOIN system.lakeflow.job_run_timeline t USING (workspace_id, job_id, run_id)
GROUP BY ALL
ORDER BY usd DESC
LIMIT 100;

4-5. データ品質: メトリックテーブルへのクエリ

sql
-- プロファイルメトリック(InferenceLog 例): バージョン 1 の集計メトリックを時系列で
SELECT window.start, column_name, count, num_nulls, distinct_count, frequent_items
FROM census_monitor_db.adult_census_profile_metrics
WHERE model_id = 1
  AND slice_key IS NULL           -- テーブル全体(スライスなし)
  AND column_name = 'income_predicted'
ORDER BY window.start;

-- ドリフトメトリック: PSI が有意(>=0.2)な数値列を検出
SELECT window.start, column_name, drift_type, population_stability_index
FROM my_schema.my_table_drift_metrics
WHERE population_stability_index >= 0.2
  AND drift_type = 'CONSECUTIVE'
ORDER BY window.start DESC;

4-6. StreamingQueryListener(Python)

python
from pyspark.sql.streaming import StreamingQueryListener

class MyListener(StreamingQueryListener):
    def onQueryStarted(self, event):
        print(f"'{event.name}' [{event.id}] started")

    def onQueryProgress(self, event):
        p = event.progress
        # 入力レート > 処理レート が続けばバックプレッシャー
        print(f"batch={p.batchId} in/s={p.inputRowsPerSecond} proc/s={p.processedRowsPerSecond}")
        for s in p.sources:
            m = s.metrics or {}
            # Delta/AutoLoader のバックログ / Kafka の遅れ
            if "numFilesOutstanding" in m:
                print("backlog files:", m["numFilesOutstanding"])
            if "maxOffsetsBehindLatest" in m:
                print("kafka lag:", m["maxOffsetsBehindLatest"])

    def onQueryIdle(self, event):
        pass

    def onQueryTerminated(self, event):
        print(f"{event.id} terminated")

spark.streams.addListener(MyListener())

4-7. observe によるデータ品質チェック(Python)

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

# 名前付き集計を DataFrame に付与
observed_df = df.observe(
    "metric",
    count(lit(1)).alias("cnt"),
    count(col("error")).alias("malformed"),
)
observed_df.writeStream.format("...").queryName("ingest_bronze").start()

class QualityListener(StreamingQueryListener):
    def onQueryProgress(self, event):
        row = event.progress.observedMetrics.get("metric")
        if row and row.cnt and row.malformed / row.cnt > 0.5:
            # ここでアラート発火(Webhook 等へ)
            print(f"ALERT: malformed {row.malformed}/{row.cnt}")
    def onQueryStarted(self, event): pass
    def onQueryTerminated(self, event): pass

spark.streams.addListener(QualityListener())

4-8. システムテーブルの増分ストリーミング読み取り

python
# 削除でストリームが壊れないよう skipChangeCommits=true を設定
(spark.readStream
      .option("skipChangeCommits", "true")
      .table("system.billing.usage"))

4-9. Delta テーブルの reservoirId 確認

sql
-- ストリーミングメトリクスの reservoirId は DESCRIBE DETAIL の id にマップされる
DESCRIBE DETAIL my_catalog.my_schema.my_table;

4-10. SQL アラート用クエリと Cron スケジュール例

sql
-- しきい値監視の例(金額 SUM > 4000 でトリガー)
SELECT to_date(tpep_pickup_datetime) AS date, SUM(fare_amount) AS amount
FROM `samples`.`nyctaxi`.`trips`
GROUP BY ALL
ORDER BY 1 DESC;
-- 条件: 値=SUM, 列=amount, 演算子=>, しきい値=Static value 4000
-- スケジュール: Quartz Cron(例: 毎時 0 分から 5 分ごと)

5. 試験で問われるポイント

  • ツール選択: 「アカウント全体のコスト/ジョブ履歴を後から分析」→システムテーブル、「テーブルの鮮度・完全性・分布の変化を自動監視」→データ品質監視、「流れているストリームのレート・遅延・状態」→StreamingQueryListener/Spark UI、「条件成立で通知」→SQL アラート/ジョブ通知。目的とツールの対応を確実に。
  • システムテーブルの前提: system カタログ、Unity Catalog 必須、アカウント管理者+メタストア管理者が既定アクセス、他ユーザーは USE CATALOG+USE SCHEMA+SELECT、読み取り専用、リージョン単位(一部グローバル)、無料保持期間あり。リアルタイム監視は不可(1 日を通して更新)。
  • 主要テーブルパスの暗記: system.billing.usage(コスト)、system.billing.list_prices(価格、USD 換算)、system.access.audit(監査)、system.lakeflow.job_run_timeline(ジョブ実行)、system.lakeflow.job_task_run_timeline(タスク/マルチタスクのコスト)。
  • ジョブコスト帰属: usage_metadata.job_id はジョブコンピュート/サーバーレスのみ設定。汎用(All-Purpose)コンピュートのジョブは正確なコスト帰属ができない→専用ジョブコンピュートかサーバーレスを推奨。WORKFLOW_RUN は親ノートブックに帰属。
  • result_state / termination_code: 失敗切り分け。中間スライス行では NULL(最終行のみ設定)。
  • データ品質の 2 本柱: 異常検出(鮮度=コミット遅延、完全性=24h の行数)とデータプロファイル(旧 Lakehouse Monitoring)。サーバーレスで課金(DATA_QUALITY_MONITORING)。監視対象テーブルは変更されない。
  • メトリックテーブル 2 種: _profile_metrics(統計)と _drift_metrics(変化)。ドリフトは連続(直前ウィンドウ)とベースライン(ベースラインテーブル)。カテゴリ列=カイ二乗/JS/TV、数値列=KS/Wasserstein/PSI。PSI >= 0.2 で有意な変化
  • 分析タイプでウィンドウが変わる: Snapshot=単一時点、TimeSeries/InferenceLog=粒度ベース。推論プロファイルは 30 日を振り返る。
  • ストリーミングの健全性指標: inputRowsPerSecond(入力)vs processedRowsPerSecond(処理)でバックプレッシャー判定。バックログは Delta/AutoLoader=numFilesOutstanding/numBytesOutstanding、Kafka=*OffsetsBehindLatestcontinuous トリガーは非サポート。
  • リスナーの制約: 処理を軽くし高速シンク(Kafka)へ書く。UC 対応コンピュートでは DBR 15.1/16.1 以降などの条件。reservoirId ≠ UC tableId
  • SQL アラート: 状態は OK/TRIGGERED/ERRORUNKNOWN 廃止)。パラメーター付きクエリ非対応。クエリは再利用不可(各アラートが所有)。アクティブ実行クォータ(既定 250)を共有しスキップされ得る→スケジュールをずらす。
  • ジョブ通知: 開始/成功/失敗/期間超過/ストリーミングバックログ。「失敗を伴う成功」は成功扱い。再試行中のタスクはジョブレベル通知が飛ばない→タスク通知を使う。イベント種別ごと最大 3 システム宛先。ストリーミングバックログは 10 分平均で判定、30 分クールダウン、awaitTermination() を使わない。
  • 通知先: 管理者が作成、Email/Slack/Teams/PagerDuty/Webhook、HTTPS 必須、SQL とジョブのみ、Slack/Teams はスキーマ非依存にするなら Webhook 推奨。

6. 理解度チェックリスト

  • [ ] システムテーブルが system カタログにあり Unity Catalog 必須・読み取り専用・リージョン単位(一部グローバル)であることを説明できる
  • [ ] システムテーブルへのアクセスに必要な権限(USE CATALOG+USE SCHEMA+SELECT、既定はアカウント管理者かつメタストア管理者)を挙げられる
  • [ ] system.billing.usagesystem.billing.list_prices を結合して USD コストを算出するクエリを書ける
  • [ ] system.access.audit(監査)、system.lakeflow.*(ジョブ)など主要テーブルパスと用途を対応づけられる
  • [ ] job_run_timelineresult_state / termination_code で失敗ジョブを切り分けられ、中間スライス行が NULL になる理由を説明できる
  • [ ] 汎用コンピュートのジョブはコスト帰属が不正確で、専用ジョブ/サーバーレスが推奨される理由を説明できる
  • [ ] システムテーブルがリアルタイム監視非対応で、スキーマ進化・選択性・ストリーミング時の skipChangeCommits に注意が要ることを理解している
  • [ ] データ品質監視の 2 本柱(異常検出=鮮度・完全性/データプロファイル)と旧称 Lakehouse Monitoring を説明できる
  • [ ] 鮮度と完全性がそれぞれ何を(コミット遅延/24h の行数)どう予測するか説明できる
  • [ ] プロファイルメトリックテーブルとドリフトメトリックテーブルの命名(_profile_metrics/_drift_metrics)と役割を区別できる
  • [ ] 連続ドリフトとベースラインドリフトの違い、カテゴリ列(カイ二乗/JS/TV)と数値列(KS/Wasserstein/PSI)の指標を挙げられる
  • [ ] PSI の閾値(<0.1 / <0.2 / >=0.2)の意味を説明できる
  • [ ] 分析タイプ(Snapshot/TimeSeries/InferenceLog)でウィンドウの決まり方が変わることを理解している
  • [ ] データ品質監視がサーバーレスで課金され、監視対象テーブルを変更しないことを理解している
  • [ ] StreamingQueryListener の 4 コールバックと、入力レート/処理レート/バックログでバックプレッシャーを判定する方法を説明できる
  • [ ] Delta/AutoLoader(numFilesOutstanding 等)と Kafka(*OffsetsBehindLatest)のバックログメトリックを区別できる
  • [ ] df.observeobservedMetrics を使ったストリーミングのデータ品質チェックを書ける
  • [ ] reservoirId が UC の tableId と異なり DESCRIBE DETAILid に対応することを知っている
  • [ ] SQL アラートの状態(OK/TRIGGERED/ERROR)、条件(値・列・演算子・しきい値)、Cron スケジュール、パラメーター非対応を説明できる
  • [ ] SQL アラートがアクティブ実行クォータを共有しスキップされ得ること、対策(スケジュールずらし)を説明できる
  • [ ] ジョブ通知のイベント(開始/成功/失敗/期間超過/ストリーミングバックログ)と「失敗を伴う成功」「タスク通知 vs ジョブ通知」を区別できる
  • [ ] 通知先の種類(Email/Slack/Teams/PagerDuty/Webhook)、管理者が作成すること、Slack/Teams はスキーマ非依存に Webhook 推奨であることを説明できる