1. はじめに
この Codelab では、データ使用者向けの技術的なブループリントを提供します。このドキュメントでは、データ ガバナンスに対する「コードファースト」のアプローチについて説明し、堅牢な品質とメタデータ管理を開発ライフサイクルに直接組み込む方法を示します。Knowledge Catalog は、インテリジェントなデータ ファブリックとして機能し、組織がデータレイクからデータ ウェアハウスまで、データ資産全体にわたってデータを一元的に管理、モニタリング、統制できるようにします。
この Codelab では、Knowledge Catalog、BigQuery、Antigravity CLI を活用して、複雑なデータをフラット化し、プログラムでプロファイリングし、インテリジェントなデータ品質ルール候補を生成し、自動化された品質スキャンをデプロイする方法を紹介します。主な目的は、エラーが発生しやすくスケーリングが難しい手動の UI 主導のプロセスから脱却し、堅牢でバージョン管理可能な「ポリシー コード」フレームワークを確立することです。
学習内容
- マテリアライズド ビューを使用してネストされた BigQuery データをフラット化し、包括的なプロファイリングを有効にする方法。
- Knowledge Catalog Python クライアント ライブラリを使用して、Knowledge Catalog プロファイル スキャンをプログラムでトリガーして管理する方法。
- プロファイル データをエクスポートして、生成 AI モデルの入力として構造化する方法。
- Antigravity CLI のプロンプトを設計して、プロファイル データを分析し、Knowledge Catalog 準拠の YAML ルールファイルを生成する方法。
- AI によって生成された構成を検証するためのインタラクティブな人間参加型(HITL)プロセスの重要性。
- 生成されたルールを自動データ品質スキャンとしてデプロイする方法。
必要なもの
- Google Cloud アカウントと Google Cloud プロジェクト
- ウェブブラウザ(Chrome など)
主なコンセプト: Knowledge Catalog のデータ品質の柱
効果的なデータ品質戦略を構築するには、Knowledge Catalog のコア コンポーネントを理解することが不可欠です。
- データ プロファイル スキャン: データを分析し、null の割合、個別の値の数、値の分布などの統計メタデータを生成する Knowledge Catalog ジョブ。これは、プログラマティックな「検出」フェーズとして機能します。
- データ品質ルール: データが満たす必要のある条件を定義する宣言型ステートメント(
NonNullExpectation、SetExpectation、RangeExpectationなど)。 - ルール候補の生成 AI: 大規模言語モデル(Gemini など)を使用してデータ プロファイルを分析し、関連するデータ品質ルールを提案します。これにより、ベースライン品質フレームワークの定義プロセスが加速します。
- データ品質スキャン: 事前定義されたルールまたはカスタムルールのセットに対してデータを検証する Knowledge Catalog ジョブ。
- プログラムによるガバナンス: ガバナンス コントロール(品質ルールなど)をコード(YAML ファイルや Python スクリプトなど)として管理することを中心としたテーマ。これにより、自動化、バージョン管理、CI/CD パイプラインへの統合が可能になります。
- 人間参加型(HITL): 人間の専門知識と監督を自動化されたワークフローに統合するための重要な制御ポイント。AI によって生成された構成の場合、HITL は、デプロイ前に提案の正確性、ビジネス関連性、安全性を検証するために不可欠です。
2. 設定と要件
Cloud Shell の起動
Google Cloud はノートパソコンからリモートで操作できますが、この Codelab では、Google Cloud Shell(Cloud 上で動作するコマンドライン環境)を使用します。
Google Cloud コンソールで、右上のツールバーにある Cloud Shell アイコンをクリックします。

プロビジョニングと環境への接続にはそれほど時間はかかりません。完了すると、次のように表示されます。

この仮想マシンには、必要な開発ツールがすべて用意されています。永続的なホーム ディレクトリが 5 GB 用意されており、Google Cloud で稼働します。そのため、ネットワークのパフォーマンスと認証機能が大幅に向上しています。この Codelab での作業はすべて、ブラウザ内から実行できます。インストールは不要です。
必要な API を有効にして環境を構成する
Cloud Shell で、プロジェクト ID が設定されていることを確認します。
export PROJECT_ID=$(gcloud config get-value project)
gcloud config set project $PROJECT_ID
export LOCATION="us-central1"
export BQ_LOCATION="us"
export DATASET_ID="kc_dq_codelab"
export TABLE_ID="ga4_transactions"
gcloud services enable dataplex.googleapis.com \
bigquery.googleapis.com \
serviceusage.googleapis.com \
aiplatform.googleapis.com
この例では、使用する一般公開サンプルデータも us(マルチリージョン)にあるため、ロケーションとして us(マルチリージョン)を使用しています。BigQuery では、クエリのソースデータと宛先テーブルが同じロケーションに存在する必要があります。
専用の BigQuery データセットを作成する
サンプルデータと結果を格納する新しい BigQuery データセットを作成します。
bq --location=us mk --dataset $PROJECT_ID:$DATASET_ID
サンプルデータを準備する
この Codelab では、Google Merchandise Store の難読化された e コマースデータを含む一般公開データセットを使用します。一般公開データセットは読み取り専用であるため、独自のデータセットに可変コピーを作成する必要があります。次の bq コマンドは、kc_dq_codelab データセットに新しいテーブル ga4_transactions を作成します。スキャンを迅速に実行するため、1 日分のデータ(2021-01-31)をコピーします。
bq query \
--use_legacy_sql=false \
--destination_table=$PROJECT_ID:$DATASET_ID.$TABLE_ID \
--replace=true \
'SELECT * FROM `bigquery-public-data.ga4_obfuscated_sample_ecommerce.events_20210131`'
デモ ディレクトリを設定する
まず、この Codelab に必要なフォルダ構造とサポート ファイルを含む GitHub リポジトリのクローンを作成します。
# Perform a shallow clone to get only the latest repository structure without the full history
git clone --depth 1 --filter=blob:none --sparse https://github.com/GoogleCloudPlatform/devrel-demos.git
cd devrel-demos
# Specify and download only the folder we need for this lab
git sparse-checkout set data-analytics/programmatic-dq
cd data-analytics/programmatic-dq
このディレクトリがアクティブな作業領域になります。以降のファイルはすべてここに作成されます。
3. Knowledge Catalog プロファイリングによる自動データ検出
Knowledge Catalog のデータ プロファイリングは、null の割合、一意性、値の分布など、データに関する統計情報を自動的に検出するための強力なツールです。このプロセスは、データの構造と品質を理解するために不可欠です。ただし、Knowledge Catalog のプロファイリングには、テーブル内のネストされたフィールドや繰り返されるフィールド(RECORD 型や ARRAY 型など)を完全に検査できないという既知の制限事項があります。列が複合型であることは識別できますが、ネストされた構造内の個々のフィールドをプロファイリングすることはできません。
この問題を解決するために、データを専用のマテリアライズド ビューにフラット化します。この戦略により、すべてのフィールドが最上位の列になり、Knowledge Catalog で各フィールドを個別にプロファイリングできます。
ネストされたスキーマについて
まず、ソーステーブルのスキーマを確認します。Google アナリティクス 4(GA4)データセットには、ネストされた列と繰り返し列が複数含まれています。すべてのネストされた構造を含む完全なスキーマをプログラムで取得するには、bq show コマンドを使用して、出力を JSON ファイルとして保存します。
bq show --schema --format=json $PROJECT_ID:$DATASET_ID.$TABLE_ID > bq_schema.json
bq_schema.json ファイルを調べると、デバイス、地域、e コマース、繰り返しレコード項目などの複雑な構造が明らかになります。これらは、効果的なプロファイリングのためにフラット化が必要な構造です。
マテリアライズド ビューによるデータの平坦化
このネストされたデータの課題に対する最も効果的で実用的なソリューションは、マテリアライズド ビュー(MV)を作成することです。フラット化された結果を事前に計算することで、MV はクエリのパフォーマンスとコストの面で大きなメリットをもたらし、アナリストやプロファイリング ツールにシンプルなリレーショナル構造を提供します。
最初に思いつくのは、すべてを 1 つの巨大なビューにフラット化することかもしれません。ただし、この直感的なアプローチには、重大なデータ破損につながる危険な落とし穴が潜んでいます。これが重大な間違いである理由を説明します。
mv_ga4_user_session_flat.sql
CREATE OR REPLACE MATERIALIZED VIEW `$PROJECT_ID.$DATASET_ID.mv_ga4_user_session_flat`
OPTIONS (
enable_refresh = true,
refresh_interval_minutes = 30
) AS
SELECT
event_date, event_timestamp, event_name, user_pseudo_id, user_id, stream_id, platform,
device.category AS device_category,
device.operating_system AS device_os,
device.operating_system_version AS device_os_version,
device.language AS device_language,
device.web_info.browser AS device_browser,
geo.continent AS geo_continent,
geo.country AS geo_country,
geo.region AS geo_region,
geo.city AS geo_city,
traffic_source.name AS traffic_source_name,
traffic_source.medium AS traffic_source_medium,
traffic_source.source AS traffic_source_source
FROM
`$PROJECT_ID.$DATASET_ID.ga4_transactions`;
mv_ga4_ecommerce_transactions.sql
CREATE OR REPLACE MATERIALIZED VIEW `$PROJECT_ID.$DATASET_ID.mv_ga4_ecommerce_transactions`
OPTIONS (
enable_refresh = true,
refresh_interval_minutes = 30
) AS
SELECT
event_date, event_timestamp, user_pseudo_id, ecommerce.transaction_id,
ecommerce.total_item_quantity,
ecommerce.purchase_revenue_in_usd,
ecommerce.purchase_revenue,
ecommerce.refund_value_in_usd,
ecommerce.refund_value,
ecommerce.shipping_value_in_usd,
ecommerce.shipping_value,
ecommerce.tax_value_in_usd,
ecommerce.tax_value,
ecommerce.unique_items
FROM
`$PROJECT_ID.$DATASET_ID.ga4_transactions`
WHERE
ecommerce.transaction_id IS NOT NULL;
mv_ga4_ecommerce_items.sql
CREATE OR REPLACE MATERIALIZED VIEW `$PROJECT_ID.$DATASET_ID.mv_ga4_ecommerce_items`
OPTIONS (
enable_refresh = true,
refresh_interval_minutes = 30
) AS
SELECT
event_date, event_timestamp, event_name, user_pseudo_id, ecommerce.transaction_id,
item.item_id,
item.item_name,
item.item_brand,
item.item_variant,
item.item_category,
item.item_category2,
item.item_category3,
item.item_category4,
item.item_category5,
item.price_in_usd,
item.price,
item.quantity,
item.item_revenue_in_usd,
item.item_revenue,
item.coupon,
item.affiliation,
item.item_list_name,
item.promotion_name
FROM
`$PROJECT_ID.$DATASET_ID.ga4_transactions`,
UNNEST(items) AS item
WHERE
ecommerce.transaction_id IS NOT NULL;
次に、bq コマンドライン ツールを使用してこれらのテンプレートを実行します。envsubst コマンドは各ファイルを読み取り、$PROJECT_ID や $DATASET_ID などの変数をシェル環境の値に置き換え、最終的な有効な SQL を bq query コマンドにパイプします。
envsubst < mv_ga4_user_session_flat.sql | bq query --use_legacy_sql=false
envsubst < mv_ga4_ecommerce_transactions.sql | bq query --use_legacy_sql=false
envsubst < mv_ga4_ecommerce_items.sql | bq query --use_legacy_sql=false
Python クライアントを使用してプロファイル スキャンを実行する
フラット化されたプロファイリング可能なビューができたので、各ビューに対して Knowledge Catalog データ プロファイル スキャンをプログラムで作成して実行できます。次の Python スクリプトは、google-cloud-dataplex クライアント ライブラリを使用してこのプロセスを自動化します。
スクリプトを実行する前に、プロジェクト ディレクトリ内に隔離された Python 環境を作成することが重要なベスト プラクティスです。これにより、プロジェクトの依存関係が個別に管理され、Cloud Shell 環境内の他のパッケージとの競合を防ぐことができます。
# Create the virtual environment
python3 -m venv dq_venv
# Activate the environment
source dq_venv/bin/activate
次に、新しく有効にした環境内に Knowledge Catalog クライアント ライブラリをインストールします。
# Install the Knowledge Catalog client library
pip install google-cloud-dataplex
環境をセットアップしてライブラリをインストールしたら、オーケストレーション スクリプトを作成する準備が整います。
Cloud Shell ツールバーで [エディタを開く] をクリックします。1_run_scan.py という名前で新しいファイルを作成し、次の Python コードを貼り付けます。GitHub リポジトリのクローンを作成すると、このファイルはすでにフォルダにあります。
このスクリプトは、各マテリアライズド ビューのスキャンを作成し(まだ存在しない場合)、スキャンを実行してから、すべてのスキャンジョブが完了するまでポーリングします。
import os
import sys
import time
from google.cloud import dataplex_v1
from google.api_core.exceptions import AlreadyExists
def create_and_run_scan(
client: dataplex_v1.DataScanServiceClient,
project_id: str,
location: str,
data_scan_id: str,
target_resource: str,
) -> dataplex_v1.DataScanJob | None:
"""
Creates and runs a single data profile scan.
Returns the executed Job object without waiting for completion.
"""
parent = client.data_scan_path(project_id, location, data_scan_id).rsplit('/', 2)[0]
scan_path = client.data_scan_path(project_id, location, data_scan_id)
# 1. Create Data Scan (skips if it already exists)
try:
data_scan = dataplex_v1.DataScan()
data_scan.data.resource = target_resource
data_scan.data_profile_spec = dataplex_v1.DataProfileSpec()
print(f"[INFO] Creating data scan '{data_scan_id}'...")
client.create_data_scan(
parent=parent,
data_scan=data_scan,
data_scan_id=data_scan_id
).result() # Wait for creation to complete
print(f"[SUCCESS] Data scan '{data_scan_id}' created.")
except AlreadyExists:
print(f"[INFO] Data scan '{data_scan_id}' already exists. Skipping creation.")
except Exception as e:
print(f"[ERROR] Error creating data scan '{data_scan_id}': {e}")
return None
# 2. Run Data Scan
try:
print(f"[INFO] Running data scan '{data_scan_id}'...")
run_response = client.run_data_scan(name=scan_path)
print(f"[SUCCESS] Job started for '{data_scan_id}'. Job ID: {run_response.job.name.split('/')[-1]}")
return run_response.job
except Exception as e:
print(f"[ERROR] Error running data scan '{data_scan_id}': {e}")
return None
def main():
"""Main execution function"""
# --- Load configuration from environment variables ---
PROJECT_ID = os.environ.get("PROJECT_ID")
LOCATION = os.environ.get("LOCATION")
DATASET_ID = os.environ.get("DATASET_ID")
if not all([PROJECT_ID, LOCATION, DATASET_ID]):
print("[ERROR] One or more required environment variables are not set.")
print("Please ensure PROJECT_ID, LOCATION, and DATASET_ID are exported in your shell.")
sys.exit(1)
print(f"[INFO] Using Project: {PROJECT_ID}, Location: {LOCATION}, Dataset: {DATASET_ID}")
# List of Materialized Views to profile
TARGET_VIEWS = [
"mv_ga4_user_session_flat",
"mv_ga4_ecommerce_transactions",
"mv_ga4_ecommerce_items"
]
# ----------------------------------------------------
client = dataplex_v1.DataScanServiceClient()
running_jobs = []
# 1. Create and run jobs for all target views
print("\n--- Starting Data Profiling Job Creation and Execution ---")
for view_name in TARGET_VIEWS:
data_scan_id = f"profile-scan-{view_name.replace('_', '-')}"
target_resource = f"//bigquery.googleapis.com/projects/{PROJECT_ID}/datasets/{DATASET_ID}/tables/{view_name}"
job = create_and_run_scan(client, PROJECT_ID, LOCATION, data_scan_id, target_resource)
if job:
running_jobs.append(job)
print("-------------------------------------------------------\n")
if not running_jobs:
print("[ERROR] No jobs were started. Exiting.")
return
# 2. Poll for all jobs to complete
print("--- Monitoring job completion status (checking every 30 seconds) ---")
completed_jobs = {}
while running_jobs:
jobs_to_poll_next = []
print(f"\n[STATUS] Checking status for {len(running_jobs)} running jobs...")
for job in running_jobs:
job_id_short = job.name.split('/')[-1][:13]
try:
updated_job = client.get_data_scan_job(name=job.name)
state = updated_job.state
if state in (dataplex_v1.DataScanJob.State.RUNNING, dataplex_v1.DataScanJob.State.PENDING, dataplex_v1.DataScanJob.State.CANCELING):
print(f" - Job {job_id_short}... Status: {state.name}")
jobs_to_poll_next.append(updated_job)
else:
print(f" - Job {job_id_short}... Status: {state.name} (Complete)")
completed_jobs[job.name] = updated_job
except Exception as e:
print(f"[ERROR] Could not check status for job {job_id_short}: {e}")
running_jobs = jobs_to_poll_next
if running_jobs:
time.sleep(30)
# 3. Print final results
print("\n--------------------------------------------------")
print("[SUCCESS] All data profiling jobs have completed.")
print("\nFinal Job Status Summary:")
for job_name, job in completed_jobs.items():
job_id_short = job_name.split('/')[-1][:13]
print(f" - Job {job_id_short}: {job.state.name}")
if job.state == dataplex_v1.DataScanJob.State.FAILED:
print(f" - Failure Message: {job.message}")
print("\nNext step: Analyze the profile results and generate quality rules.")
if __name__ == "__main__":
main()
Cloud Shell ターミナルからスクリプトを実行します。
python 1_run_scan.py
スクリプトは 3 つのマテリアライズド ビューのプロファイリングをオーケストレートし、ステータスをリアルタイムで更新します。完了すると、各ビューの機械可読性の高い統計プロファイルが作成され、ワークフローの次の段階である AI を活用したデータ品質ルールの生成の準備が整います。
完了したプロファイル スキャンは、Google Cloud コンソールで確認できます。
- ナビゲーション メニューで、[統制] セクションの [Knowledge Catalog] と [データのプロファイリングと品質] に移動します。
- 3 つのプロファイル スキャンが最新のジョブ ステータスとともに表示されます。スキャンをクリックすると、詳細な結果が表示されます。
BigQuery プロファイルから AI 対応の入力へ
Knowledge Catalog プロファイル スキャンが正常に実行されました。結果は Knowledge Catalog API で利用できますが、生成 AI モデルの入力として使用するには、構造化されたローカル ファイルに抽出する必要があります。
次の Python スクリプト 2_dq_profile_save.py は、mv_ga4_user_session_flat ビューの最新の成功したプロファイル スキャン ジョブをプログラムで検索します。次に、完全な詳細プロファイルの結果を取得し、dq_profile_results.json という名前のローカル JSON ファイルとして保存します。このファイルは、次のステップで AI 分析の直接入力として使用されます。
Cloud Shell エディタで、2_dq_profile_save.py という名前の新しいファイルを作成し、次のコードを貼り付けます。前の手順と同様に、リポジトリのクローンを作成した場合は、ファイルの作成をスキップできます。
import os
import sys
import json
from google.cloud import dataplex_v1
from google.api_core.exceptions import NotFound
from google.protobuf.json_format import MessageToDict
# --- Configuration ---
# The Materialized View to analyze is fixed for this step.
TARGET_VIEW = "mv_ga4_user_session_flat"
OUTPUT_FILENAME = "dq_profile_results.json"
def save_to_json_file(content: dict, filename: str):
"""Saves the given dictionary content to a JSON file."""
try:
with open(filename, "w", encoding="utf-8") as f:
# Use indent=2 for a readable, "pretty-printed" JSON file.
json.dump(content, f, indent=2, ensure_ascii=False)
print(f"\n[SUCCESS] Profile results were saved to '{filename}'.")
except (IOError, TypeError) as e:
print(f"[ERROR] An error occurred while saving the file: {e}")
def get_latest_successful_job(
client: dataplex_v1.DataScanServiceClient,
project_id: str,
location: str,
data_scan_id: str
) -> dataplex_v1.DataScanJob | None:
"""Finds and returns the most recently succeeded job for a given data scan."""
scan_path = client.data_scan_path(project_id, location, data_scan_id)
print(f"\n[INFO] Looking for the latest successful job for scan '{data_scan_id}'...")
try:
# List all jobs for the specified scan, which are ordered most-recent first.
jobs_pager = client.list_data_scan_jobs(parent=scan_path)
# Iterate through jobs to find the first one that succeeded.
for job in jobs_pager:
if job.state == dataplex_v1.DataScanJob.State.SUCCEEDED:
return job
# If no successful job is found after checking all pages.
return None
except NotFound:
print(f"[WARN] No scan history found for '{data_scan_id}'.")
return None
def main():
"""Main execution function."""
# --- Load configuration from environment variables ---
PROJECT_ID = os.environ.get("PROJECT_ID")
LOCATION = os.environ.get("LOCATION")
if not all([PROJECT_ID, LOCATION]):
print("[ERROR] Required environment variables PROJECT_ID or LOCATION are not set.")
sys.exit(1)
print(f"[INFO] Using Project: {PROJECT_ID}, Location: {LOCATION}")
print(f"--- Starting Profile Retrieval for: {TARGET_VIEW} ---")
# Construct the data_scan_id based on the target view name.
data_scan_id = f"profile-scan-{TARGET_VIEW.replace('_', '-')}"
# 1. Initialize client and get the latest successful job.
client = dataplex_v1.DataScanServiceClient()
latest_job = get_latest_successful_job(client, PROJECT_ID, LOCATION, data_scan_id)
if not latest_job:
print(f"\n[ERROR] No successful job record was found for '{data_scan_id}'.")
print("Please ensure the '1_run_scan.py' script has completed successfully.")
return
job_id_short = latest_job.name.split('/')[-1]
print(f"[SUCCESS] Found the latest successful job: '{job_id_short}'.")
# 2. Fetch the full, detailed profile result for the job.
print(f"[INFO] Retrieving detailed profile results for job '{job_id_short}'...")
try:
request = dataplex_v1.GetDataScanJobRequest(
name=latest_job.name,
view=dataplex_v1.GetDataScanJobRequest.DataScanJobView.FULL,
)
job_with_full_results = client.get_data_scan_job(request=request)
except Exception as e:
print(f"[ERROR] Failed to retrieve detailed job results: {e}")
return
# 3. Convert the profile result to a dictionary and save it to a JSON file.
if job_with_full_results.data_profile_result:
profile_dict = MessageToDict(job_with_full_results.data_profile_result._pb)
save_to_json_file(profile_dict, OUTPUT_FILENAME)
else:
print("[WARN] The job completed, but no data profile result was found within it.")
print("\n[INFO] Script finished successfully.")
if __name__ == "__main__":
main()
ターミナルからスクリプトを実行します。
python 2_dq_profile_save.py
正常に完了すると、ディレクトリに dq_profile_results.json という名前の新しいファイルが作成されます。このファイルには、データ品質ルールの生成に使用する詳細な統計メタデータが含まれています。dq_profile_results.json の内容を確認する場合は、次のコマンドを実行します。
cat dq_profile_results.json
4. Antigravity CLI を使用してデータ品質ルールを生成する
これで、Antigravity CLI(agy)を使用してローカル プロファイル スキャン結果を読み取ることができます。複雑なデータセットのデータ品質仕様を手動で作成するのは、時間がかかり、エラーが発生しがちです。ターミナル内で生成 AI エージェントを使用すると、最初の宣言型構成を数秒で作成できるため、このワークフローが加速します。これにより、データチームは手動による構文の作成から、ビジネスに沿った高度な人間参加型(HITL)の監督に移行できます。
Antigravity CLI を起動するには、ターミナルで次のコマンドを実行します。
agy
これで、品質ルールを生成する準備が整いました。CLI は現在のディレクトリ内のファイルを読み取ることができるため、新しいプロファイル スキャンデータを直接使用できます。
エージェントにプランを作成するよう指示する
まず、Antigravity エージェントに統計プロファイルを分析して、アクション プランを提案してもらいます。ここでは、YAML ファイルをまだ書き込まないように明示的に指示しています。これにより、分析と正当化に注意を集中させることができます。
対話型 Antigravity CLI セッションで、次の構造化プロンプトを入力します。
# Context
You are preparing a data quality rule configuration plan for Google Cloud Knowledge Catalog based on data profile statistics.
# Input
- File Path: `./dq_profile_results.json` (contains metrics like null percentage, distinct counts, and distributions)
# Task
Analyze the input statistics and propose a step-by-step plan for establishing automated data quality rules.
*Do not write any YAML code in this step.* Focus only on analytical planning.
# Rule Mapping Strategy
For candidate columns, match the statistical metrics to the most appropriate expectations:
- `nonNullExpectation`: Propose for columns with 0% null values in the profile.
- `setExpectation`: Propose for columns with a highly limited, stable set of categorical values.
- `rangeExpectation`: Propose for numeric columns with consistent and predictable value boundaries.
# Guidelines
- Provide a metric-based justification for each proposed rule (e.g., "Recommend `nonNullExpectation` for column 'user_pseudo_id' because its null percentage is 0%").
- Flag volatile metrics such as hardcoded row counts that could cause false-positive alerts in production.
# Output Format
Provide your analysis and proposed rules as a structured, step-by-step markdown plan with clear headings.
エージェントは JSON ファイルを分析し、次のような構造化されたプランを返します。
Automated Data Quality Rule Configuration Plan
Google Cloud Knowledge Catalog (Dataplex Data Quality)
──────
## Executive Summary
This analytical planning document outlines a step-by-step strategy for configuring automated data quality (DQ) rules in Google Cloud Knowledge Catalog (formerly Dataplex Data Quality) based on profiling statistics.
The dataset contains 26,489 rows representing GA4 event logs. Based on statistical metrics (null ratios, distinct value distributions, and data types), candidate columns are mapped to appropriate expectation rules.
──────
## 1. Data Profile Overview & Statistical Highlights
Column Name │ Data Type │ Null Ratio │ Distinct Count │ Key Value Range / Categories
─────────────────┼───────────┼────────────────┼────────────────┼──────────────────────────────────────────────────
event_date │ STRING │ 0.0% (0) │ 1 (3.78e-05) │ "20210131" (100%)
event_timestamp │ INTEGER │ 0.0% (0) │ ~16,539 (0.62) │ Min: 1612051200657906, Max: 1612137595412363
event_name │ STRING │ 0.0% (0) │ 16 (0.0006) │ page_view (35.8%), user_engagement (18.9%), etc.
user_pseudo_id │ STRING │ 0.0% (0) │ ~2,545 (0.09) │ 18–21 characters string identifiers
user_id │ STRING │ 100.0% (1.0) │ 0 (0.0) │ Entirely NULL
device_category │ STRING │ 0.0% (0) │ 3 (0.0001) │ desktop (57.5%), mobile (40.1%), tablet (2.4%)
... │ ... │ ... │ ... │ ...
──────
## 2. Rule Mapping Strategy & Analytical Justifications
### Step 1: Nullability Rules (nonNullExpectation)
Propose nonNullExpectation for mandatory columns where the data profile demonstrates 0% null values.
• user_pseudo_id, event_timestamp, event_name, event_date, stream_id, platform, device_category (Metric Justification: nullRatio is 0.0%)
│ [!NOTE] Exclusions:
│ • user_id: Has a nullRatio of 100.0% (unauthenticated traffic).
│ • device_language: Has a nullRatio of 37.53%.
──────
### Step 2: Categorical Value Set Validation (setExpectation)
Propose setExpectation for columns with a highly limited, stable set of categorical domain values.
• device_category: Distinct count is exactly 3. Allowed set: ['desktop', 'mobile', 'tablet']
• platform: Distinct count is 1. Allowed set expanded to: ['WEB', 'ANDROID', 'IOS'] to avoid over-fitting.
• geo_continent: Distinct count is 6. Allowed set: ['Americas', 'Asia', 'Europe', 'Africa', 'Oceania', 'Antarctica', '(not set)']
──────
### Step 3: Numeric & Timestamp Boundary Validation (rangeExpectation)
Propose rangeExpectation for numeric columns with consistent and predictable value boundaries.
• event_timestamp: rangeExpectation requiring event_timestamp > 0 (avoid dynamic microsecond range hardcoding)
• stream_id: rangeExpectation requiring positive integer stream IDs (stream_id > 0)
──────
## 3. Risk Warning: Volatile Metrics & Production False Positives
│ [!WARNING] Volatile Metrics Flagged for Risk Mitigation:
1. Hardcoded Total Row Count (rowCount = 26,489) -> Daily event volume fluctuates. Use dynamic volume thresholds.
2. Hardcoded Partition Date (event_date = '20210131') -> Breaks on future runs. Validate against YYYYMMDD regex patterns.
3. Exact Timestamp Range Bounds -> Enforcing these microsecond limits on incoming live pipelines will reject all future data.
4. Single-Value Domain Restrictions -> Single profile sample might lack active streams. Set sets according to enterprise schema.
──────
## Summary Table of Proposed Rules
Target Column │ Rule Type │ Metric-Based Justification │ Operational Considerations
─────────────────┼────────────────────┼────────────────────────────┼──────────────────────────────────────────────────
user_pseudo_id │ nonNullExpectation │ Null Ratio: 0.0% │ Core identifier, strictly required
event_timestamp │ nonNullExpectation │ Null Ratio: 0.0% │ Temporal key, strictly required
event_timestamp │ rangeExpectation │ Min: > 0 (Microseconds) │ Avoid hardcoding epoch min/max
event_name │ nonNullExpectation │ Null Ratio: 0.0% │ Required event taxonomy key
event_name │ setExpectation │ Categorical distribution │ Map to standard GA4 event taxonomy
device_category │ nonNullExpectation │ Null Ratio: 0.0% │ Required form-factor dimension
device_category │ setExpectation │ Distinct Count: 3 values │ ['desktop', 'mobile', 'tablet']
... │ ... │ ... │ ...
データ品質ルールを生成する
これは、ワークフロー全体で最も重要なステップである人間参加型(HITL)のレビューです。エージェントが生成するプランは、データの統計パターンのみに基づいています。ビジネス コンテキスト、将来のデータの変化、データの背後にある特定の意図を理解していません。人間であるエキスパートのロールは、このプランニングをバリデータ、訂正、承認者してから、コードに変換することです。
エージェントから提示されたプランをよくご確認ください。
- ご理解いただけましたでしょうか?
- ビジネスの知識と一致していますか?
- 統計的には正しいが、実際には役に立たないルールはありますか?
エージェントから受け取る出力は異なる場合があります。目標は、それを洗練することです。たとえば、テーブルのサンプルデータに固定数の行があるため、プランで rowCount ルールが提案されるとします。人間であれば、このテーブルのサイズは毎日増加することが予想されるため、厳密な行数ルールは実用的ではなく、誤ったアラートが発生する可能性が高いことを知っているでしょう。これは、AI に欠けているビジネス コンテキストを適用する完璧な例です。
次に、エージェントにフィードバックを提供し、コードを生成する最終的なコマンドを送信します。次のプロンプトは、実際に受け取ったプランと、行いたい修正に基づいて調整する必要があります。
以下のプロンプトはテンプレートです。1 行目に、具体的な修正内容を入力します。エージェントから提示されたプランが完璧で変更の必要がない場合は、その行を削除するだけで済みます。
同じ反重力セッションで、次のプロンプトの適応版を入力します。
# Feedback & Approvals
[YOUR CORRECTIONS AND APPROVAL GO HERE. Examples:
- "The plan looks good. Please proceed."
- "The rowCount rule is not necessary, as the table size changes daily. The rest of the plan is approved. Please proceed."
- "For the setExpectation on the geo_continent column, please also include 'Antarctica'."]
# Objective
Based on the approved analysis plan and the provided feedback, generate the final `dq_rules.yaml` file conforming to the standard `DataQualityRule` schema.
# Instructions
1. **Rule Justifications**: For every generated rule, add a YAML comment (`#`) on the line directly above it, briefly explaining the justification established in the plan.
2. **Schema Alignment**: Ensure the structure strictly adheres to the required Knowledge Catalog data quality scan specification. Refer to the `sample_rule.yaml` file in the current directory and the `DataQualityRule` class definition in the local virtual environment path (`./dq_venv/.../google/cloud/dataplex_v1/types/data_quality.py`) as the schema authority.
3. **Data-Driven Values**: Derive all rule parameters, such as thresholds or expected values, directly from the statistical metrics in `dq_profile_results.json`.
# Constraints
- **Output Purity**: Return ONLY the raw, valid, and properly formatted YAML code block.
- Do not include conversational preambles, introductory sentences, explanations, or markdown blocks around the YAML.
エージェントは、人間が検証した正確な指示に基づいて YAML コンテンツを生成します。完了すると、作業ディレクトリに dq_rules.yaml という名前の新しいファイルが作成されます。
データ品質スキャンを作成して実行する
エージェントが生成し、人が検証した dq_rules.yaml ファイルが作成されたので、安心してデプロイできます。
/quit と入力するか、Ctrl+C キーを 2 回押して Antigravity CLI を終了します。
次の gcloud コマンドは、新しい Knowledge Catalog データスキャン リソースを作成します。スキャンはまだ実行されません。スキャンの定義と構成(YAML ファイル)が Knowledge Catalog に登録されるだけです。
ターミナルで次のコマンドを実行します。
export DQ_SCAN="dq-scan"
gcloud dataplex datascans create data-quality $DQ_SCAN \
--project=$PROJECT_ID \
--location=$LOCATION \
--data-quality-spec-file=dq_rules.yaml \
--data-source-resource="//bigquery.googleapis.com/projects/$PROJECT_ID/datasets/$DATASET_ID/tables/mv_ga4_user_session_flat"
スキャンが定義されたので、ジョブをトリガーして実行できます。
gcloud dataplex datascans run $DQ_SCAN --location=$LOCATION --project=$PROJECT_ID
このコマンドはジョブ ID を出力します。このジョブのステータスは、Google Cloud コンソールの Knowledge Catalog セクションでモニタリングできます。完了すると、結果が分析用に BigQuery テーブルに書き込まれます。
5. 人間参加型(HITL)の重要な役割
Antigravity エージェントを使用してルールの生成を加速することは非常に強力ですが、AI を完全自動運転のパイロットではなく、高度なスキルを持つ副操縦士として扱うことが重要です。人間参加型(HITL)プロセスは、オプションの提案ではなく、データ ガバナンス ワークフローの交渉不可能な基盤となるステップです。AI が生成したアーティファクトを厳格な人間の監督なしでデプロイすることは、失敗につながります。
AI によって生成された dq_rules.yaml は、非常に高速だが経験の浅い AI デベロッパーによって送信された pull リクエストと考えることができます。ガバナンス ポリシーの「メインブランチ」に統合してデプロイするには、上級の人間による専門家(ユーザー)による徹底的なレビューが必要です。このレビューは、大規模言語モデルの固有の弱点を軽減するために不可欠です。
この人間によるレビューが不可欠な理由と、具体的に確認すべき事項を以下に詳しく説明します。
1. コンテキスト検証: AI にビジネス認識がない
- LLM の弱点: LLM はパターンと統計の専門家ですが、ビジネス コンテキストをまったく理解していません。たとえば、列
new_campaign_idの null 比率が 98% の場合、LLM は統計的な理由でこの列を無視することがあります。 - 人間の重要な役割: 人間の専門家であるあなたは、来週の主要な製品リリースに向けて
new_campaign_idフィールドが昨日追加されたことを知っています。現在、null 比率は高いはずですが、大幅に低下することが予想されます。また、入力されたら特定の形式に従う必要があることもわかっています。AI がこの外部ビジネスの知識を推測することはできません。あなたの役割は、このビジネス コンテキストを AI の統計的提案に適用し、必要に応じて提案をオーバーライドまたは補強することです。
2. 正確性と精度: ハルシネーションや微妙なエラーを防ぐ
- LLM の弱点: LLM は「自信を持って間違える」ことがあります。「ハルシネーション」が発生したり、わずかに間違ったコードを生成したりする可能性があります。たとえば、名前が正しく付けられたルールと無効なパラメータを含む YAML ファイルが生成されたり、ルールタイプが誤って入力されたり(正しい
setExpectationではなくsetExpectationsなど)することがあります。このような微妙なエラーはデプロイの失敗につながりますが、見つけにくいことがあります。 - 人間の重要な役割: 最終的なリンターとスキーマ バリデーターとして機能します。生成された YAML を公式の Knowledge Catalog
DataQualityRule仕様と照合して、綿密に確認する必要があります。単に「正しく見える」かどうかを確認するだけでなく、構文と意味の正しさを検証して、ターゲット API に 100% 準拠していることを確認します。そのため、この Codelab では、エラーの可能性を減らすために、エージェントにスキーマ ファイルを参照するよう促していますが、最終的な検証はユーザーが行う必要があります。
3. 安全性とリスク軽減: 下流の影響を防止する
- LLM の弱点: 本番環境にデプロイされたデータ品質ルールに欠陥があると、重大な結果を招く可能性があります。AI が金融取引額の
rangeExpectationを広すぎる範囲で提案すると、不正行為を検出できない可能性があります。逆に、小さなデータサンプルに基づいて厳しすぎるルールが提案された場合、オンコール チームに数千件の偽陽性アラートが殺到し、アラート疲れにつながり、実際の問題が見逃される可能性があります。 - 人間の重要な役割: あなたは安全エンジニアです。AI が提案するすべてのルールの潜在的なダウンストリームの影響を評価する必要があります。「このルールが失敗するとどうなるか?」を自問します。アラートはアクション可能ですか?このルールが誤って合格した場合のリスクは何ですか?」このリスク評価は、チェックのメリットと失敗のコストを比較する、人間ならではの能力です。
4. 継続的なプロセスとしてのガバナンス: 将来を見据えた知識の組み込み
- LLM の弱点: AI の知識は、データの静的なスナップショット(特定の時点のプロファイル結果)に基づいています。今後のイベントについては認識していません。
- 人間の重要な役割: ガバナンス戦略は将来を見据えたものでなければなりません。データソースは来月移行される予定で、stream_id が変更されることがわかっています。
geo_countryリストに新しい国が追加されることを知っています。HITL プロセスでは、この将来の状態に関する知識を注入し、計画されたビジネスまたは技術の進化中に破損を防ぐためにルールを更新または一時的に無効にします。データ品質は 1 回限りの設定ではなく、進化し続けるプロセスであり、その進化を導くことができるのは人間だけです。
要するに、HITL は、AI を活用したガバナンスを斬新ではあるもののリスクの高いアイデアから、責任あるスケーラブルなエンタープライズ グレードのプラクティスへと変革する、不可欠な品質保証と安全性のメカニズムです。これにより、最終的にデプロイされるポリシーは AI によって加速されるだけでなく、人間によって検証され、マシンのスピードと人間の専門家の知恵とコンテキストが組み合わされます。
ただし、人間の監督を重視しても、AI の価値が損なわれるわけではありません。一方、生成 AI は HITL プロセス自体を加速させるうえで重要な役割を果たします。
AI がなければ、データエンジニアは次のことを行う必要があります。
- 複雑な SQL クエリを手動で記述して、データをプロファイリングします(例: 各列の
COUNT DISTINCT、AVG、MIN、MAX)。 - 結果のスプレッドシートを 1 つずつ丁寧に分析します。
- YAML ルールファイルのすべての行をゼロから記述する。これは、面倒でエラーが発生しやすい作業です。
AI は、このような手間のかかる時間のかかる手順を自動化します。統計プロファイルを瞬時に処理し、ポリシーの「ファースト ドラフト」を 80% 完成した状態で提供する、疲れを知らないアナリストとして機能します。
これにより、人間の仕事の性質が根本的に変わります。手動でのデータ処理やボイラープレート コーディングに何時間も費やすのではなく、人間の専門家はすぐに最も価値の高いタスクに集中できます。
- ビジネス上の重要なコンテキストを適用する。
- AI のロジックの正しさを検証する。
- どのルールが本当に重要かについて戦略的な意思決定を行う。
このパートナーシップでは、AI が「何」(統計パターンとは何か)を処理し、人間は「なぜ」(このパターンがビジネスにとって重要なのはなぜか)と「だから何」(ポリシーはどうあるべきか)に集中できます。したがって、AI はループに取って代わるものではなく、ループの各サイクルをより高速かつスマートに、より効果的にします。
6. 環境をクリーンアップする
この Codelab で使用したリソースに対して Google Cloud アカウントで課金されないようにするには、リソースを含むプロジェクトを削除します。ただし、プロジェクトを保持する場合は、作成した個々のリソースを削除できます。
Knowledge Catalog スキャンを削除する
まず、作成したプロファイル スキャンと品質スキャンを削除します。重要なリソースが誤って削除されないように、これらのコマンドでは、この Codelab で作成されたスキャンの特定の名前を使用します。
# Delete the Data Quality Scan
gcloud dataplex datascans delete dq-scan \
--location=us-central1 \
--project=$PROJECT_ID --quiet
# Delete the Data Profile Scans
gcloud dataplex datascans delete profile-scan-mv-ga4-user-session-flat \
--location=us-central1 \
--project=$PROJECT_ID --quiet
gcloud dataplex datascans delete profile-scan-mv-ga4-ecommerce-transactions \
--location=us-central1 \
--project=$PROJECT_ID --quiet
gcloud dataplex datascans delete profile-scan-mv-ga4-ecommerce-items \
--location=us-central1 \
--project=$PROJECT_ID --quiet
BigQuery データセットの削除
次に、BigQuery データセットを削除します。このコマンドは元に戻すことができません。-f(強制)フラグを使用して、確認なしでデータセットとそのすべてのテーブルを削除します。
# Manually type this command to confirm you are deleting the correct dataset
bq rm -r -f --dataset $PROJECT_ID:kc_dq_codelab
7. 完了
この Codelab は終了です。
エンドツーエンドのプログラムによるデータガバナンス ワークフローを構築しました。まず、マテリアライズド ビューを使用して複雑な BigQuery データをフラット化し、分析に適した状態にしました。次に、Knowledge Catalog プロファイル スキャンをプログラムで実行して、統計メタデータを生成しました。最も重要なのは、Antigravity CLI を活用してプロファイル出力を分析し、「ポリシー コード」アーティファクト(dq_rules.yaml)をインテリジェントに生成したことです。次に、CLI を使用してこの構成を自動化されたデータ品質スキャンとしてデプロイし、最新のスケーラブルなガバナンス戦略のループを閉じました。
これで、Google Cloud で信頼性の高い AI アクセラレーションと人間による検証済みのデータ品質システムを構築するための基本的なパターンを理解できました。
次のステップ
- CI/CD と統合する:
dq_rules.yamlファイルを取得して、Git リポジトリに commit します。ルールファイルが更新されるたびに Knowledge Catalog スキャンを自動的にデプロイする CI/CD パイプライン(Cloud Build や GitHub Actions などを使用)を作成します。 - カスタム SQL ルールを調べる: 標準のルールタイプを超えて、Knowledge Catalog は、事前定義されたチェックでは表現できない、より複雑なビジネス固有のロジックを適用するためのカスタム SQL ルールをサポートしています。これは、独自の要件に合わせて検証を調整する機能です。
- 効率と費用を考慮してスキャンを最適化する: 非常に大きなテーブルの場合、データセット全体を常にスキャンしないことで、パフォーマンスを向上させ、費用を削減できます。フィルタを使用してスキャンを特定の期間やデータ セグメントに絞り込んだり、サンプリングされたスキャンを構成してデータの代表的な割合を確認したりできます。
- 結果を可視化する: Knowledge Catalog のデータ品質スキャンの出力はすべて BigQuery テーブルに書き込まれます。このテーブルを Looker Studio に接続して、定義したディメンション(完全性、有効性など)で集計されたデータ品質スコアを時系列で追跡するダッシュボードを作成します。これにより、モニタリングが事前対応型になり、すべての関係者に可視化されます。
- ベスト プラクティスを共有する: 組織内での知識の共有を促進し、集団的な経験を活用してデータ品質戦略を改善します。ガバナンスの取り組みを最大限に活用するには、データの信頼性を高める文化を育むことが重要です。
- ドキュメントを読む: