パイプライン イベントシステムテーブルリファレンス

Important

このシステム テーブルは ベータ版です

この記事は、アカウント内のpipeline_events イベント ログ エントリを記録する、 システム テーブルのリファレンスです。 各行は、パイプライン イベント ログからの不変イベントであり、ライフサイクルの遷移、フローの進行状況、データ品質メトリック、エラー、クラスター リソース、およびリージョン内のすべてのパイプラインとワークスペース全体のその他の運用データをキャプチャします。

必要条件

  • このシステム テーブルにアクセスするには、ユーザーは次のいずれかを行う必要があります。

使用可能なパイプライン イベント テーブル

パイプラインイベントシステムのテーブルはベータ版中は lakeflow_pipeline_events_preview スキーマに存在し、一般公開時に lakeflow スキーマに移動します。

テーブル Description ストリーミングをサポート 無料の保持期間 グローバルデータまたは地域データを含む
pipeline_events (ベータ版) パイプラインの実行によって出力されたパイプライン イベント ログ エントリを記録します はい 13 か月 地域

Note

このスキーマはベータ版で lakeflow_pipeline_events_preview されています。 一般公開時にテーブルは lakeflow スキーマに移動します(最終的なテーブルパスは system.lakeflow.pipeline_eventsとなります)。 ベータスキーマに対して書かれたクエリは、テーブルが移動する際に更新されなければなりません。

詳細なスキーマ リファレンス

パイプライン イベント テーブル スキーマ

パイプライン イベント テーブルは追加専用です。 各行は、パイプラインの更新によって生成された 1 つのイベントを出力時に記録し、行が変更されたり削除されたりすることはありません。

行に入力されるフィールドは、イベントの種類によって異なります。 errorupdate_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_progressupdate_progresscreate_updateなど) イベント タイプの値 の全セットについては、こちらを参照してください。
origin 構造体 クラウド プロバイダー、リージョン、パイプラインの種類、テーブルまたはフロー名、その他の識別子など、イベントの発生元に関するコンテキスト メタデータ Origin struct fieldsを参照してください。
message 文字列 イベントの人間が判読できる説明 一部のイベントでは空の場合があります。
level 文字列 イベントの重大度レベル INFOWARNERRORMETRICSのいずれか。 レベル を参照してください。
maturity_level 文字列 イベント スキーマの安定性 STABLEEVOLVINGDEPRECATED。 詳細は 成熟度レベルの数値を参照してください。
error 構造体 エラーの詳細。 エラー情報を含むイベントに対してのみ設定されます エラー 構造体フィールドを参照してください。
details バリアント イベント固有のペイロード。 含まれるフィールドは、〘 event_type 詳細 を参照してください。
event_time timestamp パイプラインによってイベントが生成された時刻 +00:00 (UTC) として記録されたタイムゾーン。

起源構造フィールド

サブフィールド データの種類 Description
cloud 文字列 クラウド プロバイダー ( AWSAZUREGCPなど)
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_SERVERSALESFORCEなど)
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_namebacklog_bytesbacklog_recordsbacklog_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 イベントに関連するユーザーコードの背景。