Important
このシステム テーブルは ベータ版です。
この記事は、アカウント内のpipeline_events イベント ログ エントリを記録する、 システム テーブルのリファレンスです。 各行は、パイプライン イベント ログからの不変イベントであり、ライフサイクルの遷移、フローの進行状況、データ品質メトリック、エラー、クラスター リソース、およびリージョン内のすべてのパイプラインとワークスペース全体のその他の運用データをキャプチャします。
必要条件
- このシステム テーブルにアクセスするには、ユーザーは次のいずれかを行う必要があります。
- メタストア管理者とアカウント管理者の両方であるか、または
- システム スキーマに対する
USE権限とSELECT権限を持つ。 システム テーブルへのアクセス権の付与を参照してください。
使用可能なパイプライン イベント テーブル
パイプラインイベントシステムのテーブルはベータ版中は lakeflow_pipeline_events_preview スキーマに存在し、一般公開時に lakeflow スキーマに移動します。
| テーブル | Description | ストリーミングをサポート | 無料の保持期間 | グローバルデータまたは地域データを含む |
|---|---|---|---|---|
| pipeline_events (ベータ版) | パイプラインの実行によって出力されたパイプライン イベント ログ エントリを記録します | はい | 13 か月 | 地域 |
Note
このスキーマはベータ版で lakeflow_pipeline_events_preview されています。 一般公開時にテーブルは lakeflow スキーマに移動します(最終的なテーブルパスは system.lakeflow.pipeline_eventsとなります)。 ベータスキーマに対して書かれたクエリは、テーブルが移動する際に更新されなければなりません。
詳細なスキーマ リファレンス
パイプライン イベント テーブル スキーマ
パイプライン イベント テーブルは追加専用です。 各行は、パイプラインの更新によって生成された 1 つのイベントを出力時に記録し、行が変更されたり削除されたりすることはありません。
行に入力されるフィールドは、イベントの種類によって異なります。
error、 update_id、および多くの origin.* サブフィールドは、それらが適用されるイベントにのみ設定され、 details フィールドの構造も event_typeによって異なります。
このテーブルを使用して、履歴パイプライン アクティビティのクエリを実行し、パイプラインエラーに関するアラートを作成し、パイプラインの動作を他の Lakeflow システム テーブルと関連付けます。
テーブル パス: system.lakeflow_pipeline_events_preview.pipeline_events
主キー: (account_id, pipeline_event_id)
| 列名 | データの種類 | Description | Notes |
|---|---|---|---|
account_id |
文字列 | このパイプライン イベントが属するアカウントの ID | |
workspace_id |
文字列 | このパイプライン イベントが属するワークスペースの ID | |
pipeline_id |
文字列 | イベントを生成したパイプラインの ID | |
update_id |
文字列 | イベントを生成したパイプライン更新の ID | |
pipeline_event_id |
文字列 | イベントのグローバル一意識別子 | |
event_type |
文字列 | イベントの種類 ( flow_progress、 update_progress、 create_updateなど) |
イベント タイプの値 の全セットについては、こちらを参照してください。 |
origin |
構造体 | クラウド プロバイダー、リージョン、パイプラインの種類、テーブルまたはフロー名、その他の識別子など、イベントの発生元に関するコンテキスト メタデータ | Origin struct fieldsを参照してください。 |
message |
文字列 | イベントの人間が判読できる説明 | 一部のイベントでは空の場合があります。 |
level |
文字列 | イベントの重大度レベル |
INFO、WARN、ERROR、METRICSのいずれか。 レベル 値を参照してください。 |
maturity_level |
文字列 | イベント スキーマの安定性 |
STABLE、EVOLVING、DEPRECATED。 詳細は 成熟度レベルの数値を参照してください。 |
error |
構造体 | エラーの詳細。 エラー情報を含むイベントに対してのみ設定されます | エラー 構造体フィールドを参照してください。 |
details |
バリアント | イベント固有のペイロード。 含まれるフィールドは、〘 event_type |
詳細 欄を参照してください。 |
event_time |
timestamp | パイプラインによってイベントが生成された時刻 |
+00:00 (UTC) として記録されたタイムゾーン。 |
起源構造フィールド
| サブフィールド | データの種類 | Description |
|---|---|---|
cloud |
文字列 | クラウド プロバイダー ( AWS、 AZURE、 GCPなど) |
region |
文字列 | クラウド プロバイダーリージョン |
org_id |
bigint | ワークスペースの組織 ID |
pipeline_type |
文字列 | パイプラインの種類 |
pipeline_name |
文字列 | パイプラインのユーザー指定の名前 |
cluster_id |
文字列 | パイプラインの更新をサポートするコンピューティング クラスター ID |
maintenance_id |
文字列 | イベントがメンテナンス実行からの場合は、メンテナンス更新プログラムの ID |
dataset_name |
文字列 | イベントが参照するデータセット (テーブルまたはビュー) の名前 |
sink_name |
文字列 | イベントが参照するシンクの名前 |
catalog_name |
文字列 | Unity カタログ カタログ名 |
schema_name |
文字列 | Unity カタログ スキーマ名 |
flow_id |
文字列 | イベントが参照するフローの ID |
flow_name |
文字列 | イベントが参照するフローの名前 |
batch_id |
bigint | ストリーミング フローのマイクロバッチ ID。
bigint 互換性のために |
request_id |
文字列 | アクションを開始した要求 ID |
materialization_name |
文字列 | 具体化の名前 |
operation_id |
文字列 | 操作の ID |
source_name |
文字列 | データ ソースの名前 |
uc_table_id |
文字列 | Unity カタログ テーブル ID |
ingestion_source_type |
文字列 | インジェスト ソースの種類 ( SQL_SERVER、 SALESFORCEなど) |
ingestion_source_connection_name |
文字列 | インジェスト ソースの接続名 |
ingestion_source_catalog_name |
文字列 | アップストリーム システムのソース カタログ名 |
ingestion_source_schema_name |
文字列 | アップストリーム システムのソース スキーマ名 |
ingestion_source_table_name |
文字列 | アップストリーム システムのソース テーブル名 |
ingestion_source_table_version |
文字列 | ソース テーブルのバージョン (該当する場合) |
エラー構造フィールド
| サブフィールド | データの種類 | Description |
|---|---|---|
fatal |
boolean | エラーが原因で更新プログラムが終了したかどうか |
exceptions |
array<struct> | エラーに関連付けられている例外のチェーン (根本原因最後) |
exceptions[].sql_state |
文字列 | SQLSTATE コード (使用可能な場合) |
exceptions[].error_class |
文字列 | Databricks エラー クラス (使用可能な場合) |
詳細分野
details列はVARIANTであり、含まれるフィールドはevent_typeによって異なります。 各イベントの種類で使用できるフィールドについては、「 パイプライン イベント ログ スキーマ」を参照してください。
variant_get関数またはドット構文を使用して、入れ子になった値を読み取ります。 一般的なアクセス パターンについては、以下のクエリ例を参照してください。
event_typeキーがペイロードを巻きつけます。 例えば、 flow_progress イベントの指標は $.flow_progress.metricsであり、 $.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 アラートを使用して、にアラートを作成できます。
pipeline_events (必要に応じて他の Lakeflow システム テーブルと結合) に対して SQL クエリを作成し、SQL ウェアハウスでスケジュールを設定し、通知先 (電子メール、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_progressbacklog_bytesイベントごとに報告され、ファイルソースについてはbacklog_filesとして報告されます。 直近の読み取り値が閾値を超えたとき(例えば、未処理の作業が100MB分)発生したときにトリガーされます。 すべての情報源がすべての指標を報告しているわけではないので、情報源が入力されている指標でフィルタリングしてください。
-- 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は各ソースごとの読み取り値の配列で、それぞれのソースの測定値の横にそのソースの1source_name、backlog_bytes、backlog_recordsがbacklog_files付きです。
パイプラインでのデータ品質の低下に関するアラート
各 flow_progress イベントは、予想される EXPECT … DROP 削除された行の数を報告します。 これらを更新ウィンドウでデータセットごとに合計し、合計が閾値を超えたらアラートします。
-- 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 マイクロバッチごとのカウントとして報告されるため、ウィンドウ内のイベントを合計すると、そのウィンドウに書き込まれた行が得られます。 スループットが予想される下限を下回ったときにアラートを作成してください。 例えば、通常は1時間に何千行も書き込むフローがほぼゼロになる場合、ソースの設定が誤っていることを示しています。
このクエリは、ウィンドウ内に行数のある flow_progress イベントを発生させたフローのみを報告します。 完全に停止したフローはイベントを発生させないため、このアラートと上記の欠損イベントアラートを組み合わせてください。
-- 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 はデータが上流に到達してからマイクロバッチがデルタテーブルにコミットするまでの時間です。 直近の読み取り値が閾値を超えたタイミング(例えば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 分ごとなど) を、バックログとデータ品質の傾向には長い間隔 (時間単位など) を使用します。 "クエリが 0 行を超える行を返すときにトリガーする" 条件は、ほとんどの場合に機能します。
参照値
レベル値
| Value | Description |
|---|---|
INFO |
通常のパイプライン アクティビティ (フローの進行状況、更新ライフサイクルの遷移、構成の変更)。 |
WARN |
パイプラインが復旧した致命的でない問題、または注意が必要な場合があります。 |
ERROR |
パイプラインがフローまたは更新を進行できなかったエラー。 |
METRICS |
実行中に出力される定量的な測定値 (行数、スループット、待機時間)。 |
満期レベルの価値
| Value | Description |
|---|---|
STABLE |
イベント スキーマは安定しています。 破壊的変更は想定されていません。 運用環境のクエリとアラートを安全に構築できます。 |
EVOLVING |
イベント スキーマは、今後のリリースで変更される可能性があります。 慎重に使用してください。 |
DEPRECATED |
イベントの種類またはスキーマは非推奨となり、今後のリリースで削除される予定です。 移行します。 |
イベントタイプの値
event_type フィールドは列挙型です。 値の完全なセット:
| Value | Description |
|---|---|
create_update |
新しいパイプラインの更新が要求されました。 |
update_progress |
パイプラインの更新がライフサイクル状態を経て遷移しました。 |
flow_progress |
更新プログラム内のフロー (データセット) が状態に遷移しました。 |
flow_definition |
フローに関する静的メタデータ。 |
dataset_definition |
データセットに関する静的メタデータ。 |
sink_definition |
出力シンクに関する静的メタデータ。 |
deprecation |
非推奨の機能がパイプラインによって使用されました。 |
autoscale |
クラスターの自動スケーリングの決定。 |
unsupported_operation |
現在の構成でサポートされていない操作。 |
cluster_resources |
バッキング コンピューティングのタスク スロットと自動スケール メトリック。 |
planning_information |
更新プログラムの計画フェーズ情報。 |
gc_pressure |
ドライバーまたは Executor に対するガベージ コレクションの負荷。 |
abnormal_termination |
更新プログラムが異常終了しました。 |
disk_space |
クラスターのディスク領域の負荷。 |
hook_progress |
パイプライン フックのライフサイクルの進行状況。 |
dataset_life_cycle |
データセットのライフサイクル イベント。 |
background_operation |
状態に遷移したバックグラウンド操作。 |
remote_api_usage |
パイプラインが送信 API 呼び出しを行いました。 |
operation_progress |
汎用操作の進行状況。 |
stream_progress |
フローをバックアップするストリーミング クエリの進行状況。 |
rewind_summary |
パイプラインの巻き戻し操作の概要。 |
advisory |
エンジンからのアドバイザリ メッセージ。 |
runtime_details |
詳細なランタイム構成。 |
resource_info |
リソース情報 (クラスター、インスタンスの種類など)。 |
file_notification_set_up |
ファイル通知のセットアップ状態 (クラウド ファイル ソースの場合)。 |
behavior_change_in_spark_connect |
Spark Connect の動作変更通知。 |
user_action |
パイプラインに対してユーザーが開始したアクション。 |
user_code_context |
イベントに関連するユーザーコードの背景。 |