メインコンテンツまでスキップ

イベントの監視とクエリー

すべての Lakeflow pipeline は、監査 logs、データ品質チェック、パイプラインの進行状況、データリネージをキャプチャするイベントを発行します。次の 2 つのソースからこれらのイベントをクエリーできます。

  • 領域内のワークスペース全体のすべてのパイプラインのイベントが pipeline_eventsシステムテーブルに格納され、イベントをクエリーするための推奨される方法となっています。ベータ版です。
  • パイプラインごとのイベントログは、単一のパイプラインのイベントを含むDeltaテーブルです。

標準のSQLを使用してイベントをクエリーします。入門用として、このページではサンプルダッシュボードとシステムテーブルの一般的なクエリーも提供しています。

要件

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

ダッシュボードの例

このサンプルダッシュボードでは、pipeline_eventsシステムテーブルを読み取って、リージョン内のすべてのパイプラインにわたるパイプラインの更新、フローのthroughput、バックログ、データ品質、およびエラーを追跡します。パイプライン、テーブル、タグ、時間範囲ごとにすべてのページをフィルター処理します。

エラーおよび警告の数を含むパイプラインステータスの内訳

パイプラインごとの1時間の変更行数

エクスペクテーションの失敗:失敗した行数と、経時的な最大失敗率

ターゲットテーブルごとの1時間のバックログバイト数

ダッシュボードをインポート

  1. ダッシュボードの JSON ファイルをdownload。
  2. ダッシュボードをワークスペースにインポートします。手順については、ダッシュボードファイルのインポートを参照してください。

クエリーのモニタリング

ダッシュボードからの次のクエリーは、一般的なパイプラインのモニタリングの使用例を示しています。

パイプラインごとの最新のエラー

このクエリーは、過去7日間にエラーが発生した各パイプラインの最も新しいエラーを、最外殻のエクセプションとともに返します。

SQL
SELECT
workspace_id,
pipeline_id,
event_time,
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

1時間あたりのエラー率

このクエリーでは、スパイクを検出できるように、パイプラインごとの1時間あたりのエラー数をカウントします。

SQL
SELECT
pipeline_id,
date_trunc('HOUR', event_time) AS hour,
count(*) AS num_errors
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
level = 'ERROR'
AND event_time >= current_timestamp() - INTERVAL 7 DAYS
GROUP BY
pipeline_id,
date_trunc('HOUR', event_time)
ORDER BY
hour DESC

フローあたりの変更行数

このクエリーは、1時間あたりにフローによって追加、アップサート、および削除された行を合算します。各メトリクスはマイクロバッチごとのカウントであるため、ウィンドウ全体で合計すると throughput が得られます。

SQL
SELECT
origin.flow_name,
date_trunc('HOUR', event_time) AS hour,
SUM(
ifnull(variant_get(details, '$.flow_progress.metrics.num_output_rows', 'BIGINT'), 0)
+ ifnull(variant_get(details, '$.flow_progress.metrics.num_upserted_rows', 'BIGINT'), 0)
+ ifnull(variant_get(details, '$.flow_progress.metrics.num_deleted_rows', 'BIGINT'), 0)
) AS rows_changed
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
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_changed DESC

フローごとのバックログ

このクエリーは、フローごとの最新のバックログ読み取り値を返します。COMPLETED フローは最新状態であるため、そのバックログは 0 として報告されます。

SQL
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
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 > 0
ORDER BY backlog_bytes DESC

データセットごとの失敗した期待値

このクエリーは、過去1日における、データセットごと、更新ごとの失敗したレコードの期待値を返します。イベントの origin.dataset_name が設定されていない場合でも設定される、エクスペクテーション自身の dataset をキーにします。

SQL
SELECT
pipeline_id,
update_id,
coalesce(expectation.dataset, origin.dataset_name) AS 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,
coalesce(expectation.dataset, origin.dataset_name),
expectation.name
HAVING
SUM(expectation.failed_records) > 0
ORDER BY
failed_records DESC