自動ローダーの監視と観察

自動ローダー パイプラインでは、ダウンストリーム コンシューマーに影響を与える前に、バックログの増加、スキーマの誤差、破損したデータ、ストールしたストリームなどの問題を検出するために、アクティブな監視が必要です。 このページでは、主要なメトリックの監視、ファイル レベルの状態のクエリ、監視ダッシュボードの構築、一般的な問題のトラブルシューティングを行う方法について説明します。

運用構成の詳細については、 実稼働ワークロードの自動ローダーの構成に関するページを参照してください。 構成のベスト プラクティスについては、「 自動ローダーのベスト プラクティス」を参照してください。

前提条件

このページのいくつかの監視ワークフローでは、バックログ クエリ、待機時間の計算、スキーマドリフト検出など、ファイルごとのインジェスト状態を監視するために、 cloud_files_state() に依存しています。 cloud_files_state() は、自動ローダー チェックポイントのファイル レベルのインジェスト状態を返すテーブル値関数です。 既定では、すべてのフィールドを使用できるわけではありません。 可用性は、Databricks ランタイムのバージョンと構成によって異なります。

  • Databricks Runtime 18.2 以降: discovery_timeprocessed_timecommit_time は自動的に使用できます。 Databricks Runtime 16.4– 18.1 では、これらのフィールドは cloudFiles.cleanSource が有効になっている場合にのみ使用できます。
  • cloudFiles.cleanSourceが有効になっている Databricks Runtime 16.4 以降では、archive_timearchive_modemove_locationを使用できます。

cloudFiles.cleanSourceを有効にすると、パフォーマンスのオーバーヘッドが発生します。 運用環境で有効にする前に、実稼働前環境のワークロードに対してベンチマークを実行します。

Additionally:

  • _metadata列を使用して取り込まれたデータに注釈を付けます。 少なくとも file_pathfile_modification_timeをキャプチャします。 「ファイル メタデータ列」を参照してください。
  • _rescued_data列と_corrupt_record列を有効にします。

Auto Loader の主な指標

次の表は、自動ローダー パイプラインを監視する最も重要なメトリックをまとめたものです。 これらのメトリックは、 StreamingQueryListener 進行状況イベントから使用でき、各ソースの metrics マップの下に自動ローダー固有の値が公開されます。

Metric それがあなたに伝えるもの
numFilesOutstanding 処理を待機しているバックログ内のファイルの数
numBytesOutstanding ファイル バックログのサイズ (バイト単位)
approximateQueueSize クラウド キューの深さ (ファイル通知モードのみ)
numInputRows 1バッチあたりの処理行数
inputRowsPerSecond データ到着率
processedRowsPerSecond 処理スループット
durationMs 内訳 各バッチで時間が費やされる場所

何に注目すべきか

次のパターンは、パイプラインに注意が必要な場合があることを示しています。

  • 増加中numFilesOutstanding: バックログがたまりつつあります。 パイプラインが流入するデータに追いついていません。
  • processedRowsPerSecond < inputRowsPerSecond: パイプラインは、データの処理速度が到着よりも遅くなっています。
  • 大きな durationMs.latestOffset: ファイルの検出が遅い。 ファイル イベントに切り替えることを検討してください。
  • 大きな durationMs.addBatch: データ処理が遅い。 コンピューティングのスケーリングまたは変換の最適化を検討してください。

メトリックの完全なリファレンスについては、「 自動ローダー ソース メトリック」を参照してください。

を使用してファイル レベルの状態を照会する cloud_files_state

cloud_files_state()テーブル値関数は、自動ローダーによって検出された各ファイルに関する詳細情報を提供します。 次のフィールドを使用できます。 Databricks Runtime 16.4 以降または 18.2 以降が必要とマークされているフィールドは、「前提条件」で説明されている条件でのみ設定 されます

フィールド タイプ Description
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 MOVEDELETE、または NULL ( 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 パイプラインの適切に構造化された自動ローダー パイプラインには、 {table}_source ビュー (自動ローダー ソース定義)、 {table}_bronze ストリーミング テーブル ( _rescued_data 列と _corrupt_record 列を含む生データ インジェスト)、解析できないデータを含む行を検疫する corrupt_records_sink 、ダウンストリームで使用するための {table} クリーン ビューが含まれます。

  • ブロンズ ストリーミング テーブルに対する期待値を設定して、スキーマの誤差とデータの破損を監視します。 _rescued_data IS NULL は予期しないスキーマ変更を検出し、 _corrupt_record IS NULL は解析不可能なデータを検出します。 Lakeflow パイプラインは、データが到着し、可観測性証跡が生成されると、これらの期待を評価します。 パイプラインに対する警告、行の削除、または失敗に対する期待値を構成できます。

パイプラインの event_log_raw ビューを作成した後、自動ローダー固有のメトリックに対して次のクエリを使用します。

フローごとのインジェスト スループットを監視する:

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 パイプラインの外部で自動ローダーを実行する場合は、次の構造化ストリーミング監視アプローチを使用します。

  • StreamingQueryListenerから読み取ることで、各バッチから自動ローダー固有のメトリックをキャプチャするsource.metricsを実装します。
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())

Note

リスナーでロジックを処理すると、クエリ処理が遅くなる可能性があります。 リスナーコールバック内での計算量を抑え、その中で外部への同期書き込みは避ける。代わりに、軽量なテレメトリを非同期に送出するか、永続化のためにメトリクスを別ジョブに引き渡す。

  • ソースの進行状況から numInputRowsinputRowsPerSecond、および processedRowsPerSecond を使用してスループットを計算します。1 秒あたりのファイル数とバッチごとの 1 秒あたりの行数です。

  • インジェストの待機時間を計算するには、 create_timecommit_timecloud_files_state() からエンド ツー エンドの待機時間と比較します。 待機時間を処理するには、 durationMs の内訳 ( latestOffsetaddBatch、その他の報告されたバッチ フェーズなど) を使用して、ボトルネックであるステージを特定します。

  • df.observe()を使用して、ストリーミング DataFrame でインライン データ品質メトリックを直接定義します。 メトリックは、StreamingQueryListenerobservedMetrics進行状況イベントに表示されます。

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 ストリーミング タブと監視ダッシュボードで自動ローダー ストリームを簡単に区別できます。

完全な構造化ストリーミング監視リファレンスについては、「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 ベースラインより大幅に高い

一般的な問題のトラブルシューティング

次の表では、自動ローダー パイプラインの一般的な問題、その原因、およびそれらを解決するための推奨されるアクションについて説明します。

Issue 考えられる原因 推奨されるアクション
バックログは処理よりも速く成長する 不足している計算リソース、データ スキュー、またはスロットルされたレート制限 コンピューティングをスケーリングし、Spark UI でスキューを確認し、 maxFilesPerTrigger 設定を確認してバッチ サイズを制御する
ファイルが検出されない ファイル イベントが正しく構成されていない、アクセス許可の問題、またはストリームが 7 日以内に実行されない 外部の場所のアクセス許可を確認し、Unity カタログ UI でファイル イベントの設定を確認し、RocksDB 状態の有効期限が切れないようにストリームが少なくとも 7 日ごとに実行されていることを確認します
ストリームの起動に時間がかかりすぎる 大規模なチェックポイント状態のダウンロード (RocksDB) 非同期状態の読み込みのために Databricks Runtime 15.3 以降にアップグレードすると、起動時間が約 90% 短縮されます
重複するファイル処理 過度な cloudFiles.maxFileAge 設定またはチェックポイントの破損 保守的な maxFileAge (最低 90 日以上) を使用し、チェックポイントの整合性を確認し、チェックポイント ストレージのライフサイクル ポリシーを回避する
パイプラインの再起動を引き起こすスキーマの進化 スキーマの頻繁な変更または互換性のない変更 schemaEvolutionModeの確認、型の昇格のaddNewColumnsWithTypeWideningへの切り替え、または非常に動的なスキーマにバリアント型を使用する
シンクに蓄積される破損したデータ ソース データ品質の問題 _corrupt_record検疫シンクでパターンを確認し、ソース データの生成を確認し、アップストリーム検証を追加することを検討します
discovery_timecommit_time が入力されていません cleanSource なしで 18.2 未満の Databricks Runtime で実行中 Databricks Runtime 18.2 以降にアップグレードするか、Databricks Runtime 16.4– 18.1 で cloudFiles.cleanSource を有効にする

その他のトラブルシューティングについては、 自動ローダーに関する FAQ を参照してください。