パイプライン イベント システム テーブルのリファレンス
ベータ版
このシステムテーブルはベータ版です。
この記事は、アカウント内のパイプラインのLakeflow pipelinesイベントログエントリを記録するpipeline_eventsシステムテーブルのリファレンスです。各行はパイプラインイベントLogからの不変のイベントであり、ライフサイクルの移行、フローの進行状況、データ品質メトリクス、エラー、クラスターリソース、および地域内のすべてのパイプラインとワークスペースにおけるその他の運用データをキャプチャします。
要件
- このシステムテーブルにアクセスするには、ユーザーは次のいずれかを実行する必要があります。
- メタストア管理者とアカウント管理者の両方である、または
- システムスキーマに
USE権限とSELECT権限が必要です。システムテーブルへのアクセス権限の付与を参照してください。
利用可能なパイプライン イベント テーブル
パイプラインイベントのシステムテーブルは、ベータ版期間中は lakeflow_pipeline_events_preview スキーマに存在し、一般提供開始時に lakeflow スキーマに移動します:
テーブル | 説明 | ストリーミングをサポート | 無料保持期間 | グローバルデータまたはリージョンデータが含まれます。 |
|---|---|---|---|---|
pipeline_events(ベータ版) | パイプラインのランによって出力されたパイプラインイベントLogsエントリを記録します。 | はい | 13ヶ月 | リージョン |
スキーマはベータ版期間中は lakeflow_pipeline_events_preview です。一般提供開始時に、テーブルは lakeflow スキーマに移動します(最終的なテーブルパスは system.lakeflow.pipeline_events になります)。ベータ版スキーマに対して記述されたクエリーは、テーブルの移動時に更新する必要があります。
詳細スキーマ参照
パイプライン イベント テーブル スキーマ
パイプラインイベントテーブルは追記専用です。各行は、パイプラインの更新によって発行された単一のイベントを、発行時点のデータとして記録し、行がその場で変更または削除されることはありません。
どのフィールドが行で入力されるかは、イベントタイプによって異なります。error、update_id、および多数の origin.* サブフィールドは、適用されるイベントでのみ設定され、details フィールドの構造も event_type によって異なります。
このテーブルを使用して、過去のパイプラインアクティビティを照会し、パイプラインの失敗に関するアラートを作成し、パイプラインの動作を他のLakeFlowシステムテーブルと関連付けることができます。
テーブルパス : system.lakeflow_pipeline_events_preview.pipeline_events
プライマリーキー : (account_id, pipeline_event_id)
列名 | データ型 | 説明 | 注 |
|---|---|---|---|
| string | このパイプラインイベントが属するアカウントのID | |
| string | このパイプラインイベントが属するワークスペースのID | |
| string | イベントを出力したパイプラインのID | |
| string | イベントを発行したパイプライン更新のID | |
| string | イベントのグローバル一意識別子 | |
| string | イベントのタイプ(例えば、 | 値の全セットについては、イベントタイプの値を参照してください。 |
| struct | イベントの発生元に関するコンテキストメタデータ(クラウドプロバイダー、リージョン、パイプラインタイプ、テーブル名またはフロー名、その他の識別子など) | Origin struct フィールドを参照してください。 |
| string | 人間にとってわかりやすいイベントの説明 | 一部のイベントでは空である可能性があります。 |
| string | イベントの重大度レベル |
|
| string | イベントスキーマの安定性 |
|
| struct | エラーの詳細。エラー情報を持つイベントにのみ設定されます。 | エラー構造体フィールドを参照してください。 |
| バリアント | イベント固有のペイロードです。含まれるフィールドは、〜によって異なります。 | 詳細フィールドを参照してください。 |
| timestamp | パイプラインから出力されたイベントの時刻 | タイムゾーンは |
元の構造体フィールド
サブフィールド | データ型 | 説明 |
|---|---|---|
| string | クラウド プロバイダー(例えば、 |
| string | クラウドプロバイダーリージョン |
| bigint | ワークスペース組織ID |
| string | パイプラインのタイプ |
| string | ユーザーが指定したパイプライン名。 |
| string | パイプライン更新用のコンピュートクラスターID |
| string | イベントがメンテナンス実行によるものである場合、メンテナンス更新のID |
| string | イベントが参照するデータセット(テーブルまたはビュー)の名前 |
| string | イベントが参照するシンクの名前 |
| string | Unity Catalogカタログ名 |
| string | Unity Catalogスキーマ名 |
| string | イベントが参照するフローのID |
| string | イベントが参照するフロー名 |
| bigint | ストリーミングフローのマイクロバッチ ID です。 |
| string | アクションを開始したリクエストID |
| string | マテリアライゼーション名 |
| string | オペレーションのID |
| string | データソースの名前 |
| string | Unity CatalogテーブルID |
| string | 取り込みソースのタイプ(例えば、 |
| string | 取り込みソースの接続名 |
| string | アップストリームシステムでのソースカタログ名 |
| string | 上流システムのソーススキーマ名 |
| string | アップストリームシステムにあるソーステーブル名 |
| string | ソーステーブルのバージョン、該当する場合。 |
エラー構造体フィールド
サブフィールド | データ型 | 説明 |
|---|---|---|
| boolean | エラーにより更新が終了したかどうか |
| array<struct> | エラーに関連する例外の連鎖(最後の根本原因) |
| string | SQLSTATE コード(利用可能な場合) |
| string | Databricks エラークラス (利用可能な場合) |
詳細フィールド
details列はVARIANTであり、含まれるフィールドはevent_typeによって異なります。各イベント タイプで利用可能なフィールドについては、「パイプライン イベント ログ スキーマ」を参照してください。variant_get 関数またはドット構文を使用して、ネストされた値を読み取ります。一般的なアクセスパターンについては、以下のクエリ例を参照してください。
event_type キーがペイロードをラップします。たとえば、flow_progress イベントのメトリクスは $.metrics ではなく $.flow_progress.metrics にあります。すべてのパスにイベントタイプキーを含めてください。
-- Using variant_get (lets you cast to a specific type)
SELECT
pipeline_id,
event_time,
variant_get(details, '$.flow_progress.status', 'STRING') AS flow_status,
variant_get(details, '$.flow_progress.metrics.num_output_rows', 'BIGINT') AS rows_written,
variant_get(details, '$.flow_progress.metrics.backlog_bytes', 'BIGINT') AS backlog_bytes
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 HOUR
-- Using dot syntax (returns VARIANT, cast when needed)
SELECT
pipeline_id,
event_time,
details:flow_progress.status::STRING AS flow_status,
details:flow_progress.metrics.num_output_rows::BIGINT AS rows_written
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 HOUR
クエリの例
-- Flow throughput for a specific pipeline
SELECT
origin.flow_name,
date_trunc('HOUR', event_time) AS hour,
SUM(variant_get(details, '$.flow_progress.metrics.num_output_rows', 'BIGINT')) AS rows_written
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
pipeline_id = '<your-pipeline-id>'
AND event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 7 DAYS
GROUP BY
origin.flow_name,
date_trunc('HOUR', event_time)
ORDER BY
hour DESC,
rows_written DESC
-- The latest error for each pipeline that has errored in the last 7 days, with the outermost exception.
-- The exception chain is ordered with the root cause last, so read element -1 for the root cause.
-- On many errors only the first element carries error_class and sql_state.
SELECT
workspace_id,
pipeline_id,
event_time,
event_type,
message,
error.exceptions[0].error_class AS exception_error_class,
error.exceptions[0].sql_state AS exception_sql_state
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
level = 'ERROR'
AND event_time >= current_timestamp() - INTERVAL 7 DAYS
QUALIFY
ROW_NUMBER() OVER (PARTITION BY workspace_id, pipeline_id ORDER BY event_time DESC) = 1
ORDER BY
event_time DESC
-- Data quality: failed expectations by dataset, per update, in the last 1 day
SELECT
pipeline_id,
update_id,
origin.dataset_name,
expectation.name AS expectation_name,
SUM(expectation.failed_records) AS failed_records
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
LATERAL VIEW explode(variant_get(details, '$.flow_progress.data_quality.expectations', 'ARRAY<STRUCT<name:STRING,dataset:STRING,passed_records:BIGINT,failed_records:BIGINT>>')) AS expectation
WHERE
event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 DAY
GROUP BY
pipeline_id,
update_id,
origin.dataset_name,
expectation.name
HAVING
SUM(expectation.failed_records) > 0
ORDER BY
failed_records DESC
一般的な結合パターン
パイプライン名でフィルタするため、pipelinesテーブルと結合します。
pipelinesテーブルは、ゆっくり変化するディメンション(SCD2)です。結合する前に、各パイプラインの最新バージョンを使用してください。
WITH latest_pipelines AS (
SELECT *
FROM system.lakeflow.pipelines
QUALIFY ROW_NUMBER() OVER (PARTITION BY workspace_id, pipeline_id ORDER BY change_time DESC) = 1
)
SELECT
p.name AS pipeline_name,
e.event_time,
e.event_type,
e.level,
e.message
FROM
system.lakeflow_pipeline_events_preview.pipeline_events e
JOIN
latest_pipelines p
ON e.workspace_id = p.workspace_id
AND e.pipeline_id = p.pipeline_id
WHERE
e.level = 'ERROR'
AND e.event_time >= current_timestamp() - INTERVAL 24 HOURS
ORDER BY
e.event_time DESC
pipeline_update_timeline での結合 update_id
SELECT
u.period_start_time AS update_start,
u.period_end_time AS update_end,
e.event_time,
e.event_type,
e.level,
e.message
FROM
system.lakeflow.pipeline_update_timeline u
JOIN
system.lakeflow_pipeline_events_preview.pipeline_events e
ON e.update_id = u.update_id
WHERE
u.pipeline_id = '<your-pipeline-id>'
AND u.period_start_time >= current_timestamp() - INTERVAL 7 DAYS
ORDER BY
u.period_start_time DESC,
e.event_time ASC
アラートのセットアップ
pipeline_events上で、Databricks SQLのアラートを使用してアラートを構築できます。pipeline_eventsに対してSQLクエリを記述し(オプションで他のLakeFlowシステムテーブルと結合)、それをSQLウェアハウスでスケジュールし、通知送信先(Eメール、Slack、webhook、PagerDutyなど)を設定します。
いくつかの便利な開始点:
過去N分間、パイプラインにイベントが到着しなかった場合にアラートします。
これを使用して、停止している、またはサイレントに失敗しているパイプラインを検出します。
-- Returns one row per pipeline that has not emitted any event in the last 30 minutes.
-- The alert can trigger when this query returns any rows.
SELECT
workspace_id,
pipeline_id,
MAX(event_time) AS last_event_time
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
event_time >= current_timestamp() - INTERVAL 24 HOURS
GROUP BY
workspace_id,
pipeline_id
HAVING
MAX(event_time) < current_timestamp() - INTERVAL 30 MINUTES
特定のフローのバックログが高すぎる場合にアラートします。
バックログはflow_progressイベントでbacklog_bytesとして報告され、ファイルソースの場合はbacklog_filesとしても報告されます。最新の読み取り値がthresholdを超えたときにTriggerします(例:100 MBの未処理作業)。すべてのソースがすべてのメトリクスを報告するわけではないため、使用しているソースが生成するメトリクスでフィルタリングしてください。
-- Returns the most recent backlog reading per flow for a given pipeline.
-- The alert can trigger when backlog_bytes exceeds the threshold for any flow.
WITH latest_flow_progress AS (
SELECT
workspace_id,
pipeline_id,
origin.flow_name,
event_time,
CASE
WHEN variant_get(details, '$.flow_progress.status', 'STRING') = 'COMPLETED' THEN 0
ELSE variant_get(details, '$.flow_progress.metrics.backlog_bytes', 'BIGINT')
END AS backlog_bytes
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
pipeline_id = '<your-pipeline-id>'
AND event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 HOUR
AND (
variant_get(details, '$.flow_progress.metrics.backlog_bytes', 'BIGINT') IS NOT NULL
OR variant_get(details, '$.flow_progress.status', 'STRING') = 'COMPLETED'
)
QUALIFY
ROW_NUMBER() OVER (PARTITION BY workspace_id, pipeline_id, origin.flow_name ORDER BY event_time DESC) = 1
)
SELECT *
FROM latest_flow_progress
WHERE backlog_bytes > 100000000
どのソースが遅れているかを見るために、 $.flow_progress.metrics.source_metrics は各ソースごとの読み取り値の配列で、それぞれのソースの読み取り値の横に source_name が並び、その backlog_bytes、 backlog_records 、 backlog_filesが付けられています。
パイプラインのデータ品質低下に関するアラート
各 flow_progress イベントは、EXPECT … DROP 個のエクスペクテーションによって削除された行数を報告します。更新ウィンドウ全体でデータセットごとにこれらを合計し、合計がthresholdを超えた場合にアラートを通知します。
-- Returns datasets where more than 100 rows were dropped by expectations, per update, in the last hour.
SELECT
pipeline_id,
update_id,
origin.dataset_name,
SUM(variant_get(details, '$.flow_progress.data_quality.dropped_records', 'BIGINT')) AS dropped_records
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 HOUR
GROUP BY
pipeline_id,
update_id,
origin.dataset_name
HAVING
SUM(variant_get(details, '$.flow_progress.data_quality.dropped_records', 'BIGINT')) > 100
フローが処理する行数が少なすぎる場合にアラートを出す
flow_progress イベントはmetrics.num_output_rowsをマイクロバッチごとのカウントとしてレポートするため、ウィンドウ内のイベントを合計すると、そのウィンドウ全体で書き込まれた行数が得られます。throughputが期待される下限を下回った場合にアラートを作成します。たとえば、通常は1時間あたり数千行を書き込むフローが、ほぼゼロしか生成しない場合は、ソースの設定が誤っている可能性があります。
このクエリーは、ウィンドウ内で行数を含む flow_progress イベントを出力したフローのみを報告します。完全に停止したフローはイベントを出力しないため、このアラートを上記のアラート(missing-events)と組み合わせて使用してください。
-- Returns flows that wrote fewer than 100 rows in the last hour.
-- The alert can trigger when this query returns any rows.
SELECT
pipeline_id,
origin.flow_name,
SUM(variant_get(details, '$.flow_progress.metrics.num_output_rows', 'BIGINT')) AS rows_written_last_hour
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
pipeline_id = '<your-pipeline-id>'
AND event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 HOUR
AND variant_get(details, '$.flow_progress.metrics.num_output_rows', 'BIGINT') IS NOT NULL
GROUP BY
pipeline_id,
origin.flow_name
HAVING
SUM(variant_get(details, '$.flow_progress.metrics.num_output_rows', 'BIGINT')) < 100
新規データのレイテンシーが高すぎる場合にアラートを出す
ストリーミングフローの場合、flow_progressイベントはstreaming_metricsでレイテンシーを報告します。stream_latency_msは、データがアップストリームに到達してから、マイクロバッチがDeltaテーブルにコミットされるまでの時間です。最新の読み取り値がthresholdを超えたときに発動するTriggerを設定できます(例:5分)。
タグ付けされたイベント時間を持つストリーミングフローのみが stream_latency_ms をレポートし、それはSDP時間メトリクスが有効な場合に限られます。他のフローはイベントごとにNULLを返し、それらのフローに対してこのアラートが起動することはありません。
-- Returns the most recent new-data latency per flow for a given pipeline.
-- The alert can trigger when stream_latency_ms > 300000 (5 minutes) for any flow.
WITH latest_latency AS (
SELECT
workspace_id,
pipeline_id,
origin.flow_name,
event_time,
CASE
WHEN variant_get(details, '$.flow_progress.status', 'STRING') = 'COMPLETED' THEN 0
ELSE variant_get(details, '$.flow_progress.streaming_metrics.stream_latency_ms', 'BIGINT')
END AS stream_latency_ms
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
pipeline_id = '<your-pipeline-id>'
AND event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 HOUR
AND (
variant_get(details, '$.flow_progress.streaming_metrics.stream_latency_ms', 'BIGINT') IS NOT NULL
OR variant_get(details, '$.flow_progress.status', 'STRING') = 'COMPLETED'
)
QUALIFY
ROW_NUMBER() OVER (PARTITION BY workspace_id, pipeline_id, origin.flow_name ORDER BY event_time DESC) = 1
)
SELECT *
FROM latest_latency
WHERE stream_latency_ms > 300000
本番運用アラートのヒント
- アカウント全体ではなく、各アラートが特定のスコープを対象とするように、
pipeline_idでフィルタリングします(ワークスペースごとにアラートを管理している場合はworkspace_idも対象とします)。 - アラートの感度に合った評価のケイデンスを設定します。高速フェイル信号には短い間隔(例:5分ごと)を使用し、バックログやデータ品質の傾向には長い間隔(例:1時間ごと)を使用します。「クエリが0より多くの行を返すときにトリガーする」という条件は、ほとんどの場合に機能します。
リファレンス
レベル値
Value | 説明 |
|---|---|
| 通常のパイプラインアクティビティ(フローの進捗、更新ライフサイクルの遷移、構成の変更) |
| パイプラインが復旧したものの、致命的ではない問題、または対処が必要な問題 |
| フローまたは更新におけるパイプラインの進行を妨げた失敗 |
| 実行中に発行される定量的な測定値(行数、スループット、レイテンシ) |
成熟度レベルの値
Value | 説明 |
|---|---|
| イベントスキーマは安定しています。互換性に影響する変更は想定されません。本番運用クエリーやアラートの構築に安全にご利用いただけます。 |
| イベントスキーマは今後のリリースで変更される場合があります。注意して使用してください。 |
| イベントタイプまたはスキーマは非推奨であり、今後のリリースで削除される予定です。それから移行してください。 |
イベントタイプの値
event_typeフィールドは列挙です。値のフルセット:
Value | 説明 |
|---|---|
| 新しいパイプライン更新が要求されました。 |
| パイプラインの更新がライフサイクル状態を遷移しました。 |
| 更新内のフロー(データセット)が状態を遷移しました。 |
| フローに関する静的メタデータ。 |
| データセットに関する静的メタデータ。 |
| 出力シンクに関する静的メタデータ。 |
| 非推奨の機能がパイプラインによって使用されました。 |
| クラスターオートスケール決定。 |
| 現在の構成では、操作がサポートされていません。 |
| バッキング コンピュートのタスク スロットおよびオートスケールのメトリクス |
| 更新に関する計画段階の情報 |
| ドライバーまたはエグゼキューターにおけるガベージコレクション負荷。 |
| 更新が異常に終了しました。 |
| クラスター上のディスク容量のひっ迫。 |
| パイプラインフックのライフサイクルの進行状況。 |
| データセットのライフサイクルイベントです。 |
| バックグラウンド操作が状態に遷移しました。 |
| パイプラインが外部 API 呼び出しを行いました。 |
| 一般的な操作の進行状況。 |
| フローを支えるストリーミングクエリの進捗 |
| パイプラインの復元操作の概要 |
| エンジンからのアドバイザリーメッセージです。 |
| ランタイムの設定の詳細。 |
| リソース情報(クラスター、インスタンスタイプなど)。 |
| ファイル通知設定状況(クラウドファイルソース向け)。 |
| Spark Connect における動作変更の通知 |
| ユーザーによるパイプラインに対するアクションです。 |
| イベントに関連付けられたユーザーコードに関するコンテキスト。 |