1. はじめに
この Codelab では、Google Cloud BigQuery 内で直接、モジュール式のエンドツーエンドの顧客 ID 解決(エンティティ マッチング)エンジンを構築します。インフラストラクチャのデプロイには Google Cloud Shell を、データのクリーンアップ、候補のスコアリング、プロパティ グラフの構築、ISO GQL(グラフ クエリ言語)パスのトラバーサルには BigQuery Studio SQL エディタを組み合わせます。
ID 解決は、エンタープライズ Customer 360、不正行為の検出、マルチシステム データ統合の基盤となる機能です。データの成熟度やビジネスニーズに応じて、ID 解決には多くの有効なアプローチがあるため、この Codelab のすべてのステップはモジュール式で省略可能です。このパイプラインは、リモート UDF アドレスの正規化、Soundex 音声ブロッキング、セマンティック ベクトル検索(AI.EMBED)、ハイブリッド特徴スコアリング、GQL プロパティ グラフ クラスタリングなど、さまざまな一般的な本番環境グレードの業界技術を紹介するように設計されています。そのため、アーキテクチャに適合するパターンを選択的に採用できます。
照合方法とスコアリングしきい値は、組織の決定論的照合と確率論的照合のニーズに基づいて調整する必要があります。これは、対象のユースケースによって決まります。たとえば、厳格なコンプライアンス、請求、財務業務では、誤ったリンクを防ぐために、高精度の決定的ルール(SSN や納税者番号の完全一致など)が優先される傾向があります。一方、マーケティングのパーソナライズ、分析、レコメンデーション エンジンでは、再現率を最大化し、微妙なつながりを明らかにするために、確率的ファジー マッチングとセマンティック ベクトル類似性が重視される傾向があります。

演習内容
- FEBRL3 ベンチマーク データセットを取り込む: 合成顧客レコードとグラウンド トゥルース一致ペアを BigQuery に読み込みます。
- Address Validation リモート UDF をデプロイする: Python Cloud Functions をデプロイし、BigQuery リモート関数を登録して、住所を正規化します。
- プロフィール データと音声エンコードを前処理する: SQL データ クリーニングを実行し、アドレス UDF を呼び出し、
SOUNDEX音声キーとレーベンシュタイン編集距離を計算します。- Soundex 音声エンコード: 英語の発音どおりに音声で名前をインデックスに登録するための音声アルゴリズム。名前を子音の音のグループを表す 4 文字のコード(最初の文字と 3 桁の数字)に変換し(たとえば、
"John"と"Jon"はJ500に、"Smith"と"Smyth"はS530にマッピングされます)、特徴スコアリングとリアルタイムの増分デルタ ブロックのための音声一致シグナルを提供します。 - レーベンシュタイン距離(
EDIT_DISTANCE): ある文字列を別の文字列に変更するために必要な最小の 1 文字の編集回数(挿入、削除、置換)を測定する文字列指標。あいまいな名前と住所を正確に照合できます。
- Soundex 音声エンコード: 英語の発音どおりに音声で名前をインデックスに登録するための音声アルゴリズム。名前を子音の音のグループを表す 4 文字のコード(最初の文字と 3 桁の数字)に変換し(たとえば、
- セマンティック プロファイル エンベディングとベクトル検索を生成する:
AI.EMBED(text-embedding-005)を使用して SQL でテキスト エンベディングを直接生成し、VECTOR_SEARCHを使用して上位 K 個の最近傍を見つけて、準線形候補生成レイヤとして機能させます。 - 候補ペアのスコアリングとハイブリッド エッジ特徴融合: ベクトル検索候補ペアを活用して O(N²) クロス結合の複雑さを解消し、マルチ特徴の重み付き類似度スコア(SSN、レーベンシュタイン編集距離、生年月日、住所のジャカード)を計算し、エッジを統合候補テーブルに融合します。
- プロパティ グラフの構築と ISO GQL パス トラバーサル: BigQuery
PROPERTY GRAPHを構築し、ISO GQL{1, 2}パス クエリ(GRAPH_TABLE)を実行して、接続された顧客クラスタを解決し、個々の評価指標を計算し、Adamic-Adar グラフの重み付けを使用してソフト世帯クラスタリングを実行します。 - 増分解決と永続的な安定性: 増分デルタ マッチングを使用して、毎日のバッチ取り込みを処理します。
- クラスタリングの統合とクラスタの安定性(1-ε の重複): (1-ε)の重複しきい値保証を使用して、パイプライン実行全体でクラスタの安定性を維持します。
必要なもの
- ウェブブラウザ(Chrome など)。
- 課金を有効にした Google Cloud プロジェクト
この Codelab は、初心者を含むあらゆるレベルのデータ エンジニア、データベース デベロッパー、AI/ML 実務担当者を対象としています。
推定所要時間: 45 分
推定費用: 2.00 米ドル未満(従量課金制の Cloud Functions と BigQuery クエリ処理を使用)。
2. 始める前に
Google Cloud プロジェクトの作成
- Google Cloud コンソールのプロジェクト セレクタ ページで、Google Cloud プロジェクトを選択または作成します。
- Cloud プロジェクトに対して課金が有効になっていることを確認します。プロジェクトで課金が有効になっているかどうかを確認する方法をご覧ください。
Cloud Shell の起動
Cloud Shell は、必要なツールがプリロードされた Google Cloud で動作するコマンドライン環境です。
- Google Cloud コンソールの上部にある [Cloud Shell をアクティブにする] をクリックします。
- 認証を確認します。
gcloud auth list
- Cloud Shell で環境変数を構成します。
export GCP_PROJECT=$(gcloud config get-value project)
export REGION="us-central1"
export DATASET_ID="identity_resolution"
必要な API の有効化
ユーザー アカウントを使用して Cloud Shell で次のコマンドを実行し、必要なすべての Google Cloud サービスを有効にします。
gcloud services enable \
addressvalidation.googleapis.com \
cloudbuild.googleapis.com \
cloudfunctions.googleapis.com \
cloudresourcemanager.googleapis.com \
artifactregistry.googleapis.com \
aiplatform.googleapis.com \
run.googleapis.com \
bigqueryconnection.googleapis.com \
bigqueryreservation.googleapis.com \
bigquery.googleapis.com
サービス アカウントと権限借用を設定する(推奨)
API のシームレスな実行とアプリケーションのデフォルト認証情報(ADC)へのアクセスを確保するには、専用のラボ サービス アカウントを作成し、gcloud の権限借用を有効にします。
# 1. Create a Service Account for the lab (if it does not already exist)
gcloud iam service-accounts create identity-res-sa \
--display-name="Identity Resolution Service Account" 2>/dev/null || true
# Wait 5 seconds for IAM propagation
sleep 5
# Extract Project Number for default build and compute service accounts
export PROJECT_NUMBER=$(gcloud projects describe ${GCP_PROJECT} --format="value(projectNumber)")
# 2. Grant specific required least-privilege roles to the lab Service Account
for role in roles/bigquery.admin \
roles/bigquery.resourceAdmin \
roles/run.admin \
roles/cloudfunctions.admin \
roles/resourcemanager.projectIamAdmin \
roles/cloudbuild.builds.editor \
roles/cloudbuild.builds.builder \
roles/artifactregistry.repoAdmin \
roles/artifactregistry.writer \
roles/storage.admin \
roles/logging.logWriter \
roles/iam.serviceAccountUser \
roles/aiplatform.user; do
gcloud projects add-iam-policy-binding ${GCP_PROJECT} \
--member="serviceAccount:identity-res-sa@${GCP_PROJECT}.iam.gserviceaccount.com" \
--role="${role}" --quiet
done
# 3. Grant required build & storage permissions to default Compute Engine & Cloud Build service accounts (required for 2nd-gen Cloud Functions container builds)
for role in roles/cloudbuild.builds.builder \
roles/logging.logWriter \
roles/artifactregistry.writer \
roles/storage.objectAdmin; do
gcloud projects add-iam-policy-binding ${GCP_PROJECT} \
--member="serviceAccount:${PROJECT_NUMBER}-compute@developer.gserviceaccount.com" \
--role="${role}" --quiet || true
gcloud projects add-iam-policy-binding ${GCP_PROJECT} \
--member="serviceAccount:${PROJECT_NUMBER}@cloudbuild.gserviceaccount.com" \
--role="${role}" --quiet || true
done
# 4. Grant Service Account Token Creator role to your user account
export SA_EMAIL="identity-res-sa@${GCP_PROJECT}.iam.gserviceaccount.com"
gcloud iam service-accounts add-iam-policy-binding \
"${SA_EMAIL}" \
--member="user:$(gcloud config get-value account)" \
--role="roles/iam.serviceAccountTokenCreator" --quiet
# 5. Enable Service Account impersonation for gcloud
gcloud config set auth/impersonate_service_account "${SA_EMAIL}"
# 6. Wait for IAM role assignments and impersonation caches to propagate
echo "Waiting 90 seconds for IAM policies and impersonation caches to propagate..."
sleep 90
BigQuery データセットを作成する
顧客ノード、エッジ、グラフモデル、評価ビューを保存する BigQuery データセットを作成します。
bq mk --location=US --dataset ${GCP_PROJECT}:${DATASET_ID}
次のような出力が表示されます。
Dataset 'your-project-id:identity_resolution' successfully created.
BigQuery の予約と割り当てを作成する(省略可 / 推奨)
オンデマンド CPU の上限や共有割り当てに制約されることなく、ベクトル インデックス検索、グラフ集計、リモート関数の実行専用のコンピューティング容量を確保するには、Cloud Shell で自動スケーリングを使用して Enterprise Edition の予約を作成します。
# 1. Create a BigQuery Enterprise reservation with 0 baseline slots and 100 max autoscaling slots
bq mk --reservation \
--project_id=${GCP_PROJECT} \
--location=US \
--edition=ENTERPRISE \
--slots=0 \
--autoscale_max_slots=100 \
--ignore_idle_slots=true \
identity-res-reservation
# 2. Assign your Cloud project to the newly created reservation for query execution
bq mk --reservation_assignment \
--project_id=${GCP_PROJECT} \
--location=US \
--reservation_id=identity-res-reservation \
--job_type=QUERY \
--assignee_type=PROJECT \
--assignee_id=${GCP_PROJECT}
3. FEBRL3 の顧客ノード データセットを取り込む
アドレス検証リモート関数をデプロイして ID 解決を実行する前に、Python の recordlinkage ライブラリを使用して合成の FEBRL3 エンティティ解決ベンチマーク データセット(顧客ごとに最大 5 つの重複を含むマルチ重複クラスタを含む 5,000 件の顧客レコードを含む)を読み込み、BigQuery DataFrames(bigframes)を使用して未加工の顧客ノード(customer_nodes)とグラウンド トゥルース一致リンク(ground_truth_links)を BigQuery に書き込みます。
Cloud Shell で次のコマンドを実行して、依存関係をインストールし、取り込みスクリプトを実行します。
# 1. Install recordlinkage dataset library & bigframes (if outside Cloud Shell, activate your virtual environment first)
pip install recordlinkage bigframes --quiet
# 2. Write and execute the FEBRL3 dataset ingestion script
cat << 'EOF' > ingest_febrl.py
import os
import pandas as pd
import bigframes.pandas as bpd
from recordlinkage.datasets import load_febrl3
GCP_PROJECT = os.environ.get("GCP_PROJECT", "your-project-id")
DATASET_ID = "identity_resolution"
table_raw_id = f"{GCP_PROJECT}.{DATASET_ID}.customer_nodes"
table_gt_id = f"{GCP_PROJECT}.{DATASET_ID}.ground_truth_links"
print("Loading FEBRL3 benchmark dataset...")
df_nodes, true_links = load_febrl3(return_links=True)
df_nodes = df_nodes.reset_index()
df_nodes['dataset_source'] = 'febrl3'
for col in df_nodes.columns:
if df_nodes[col].dtype == 'object':
df_nodes[col] = df_nodes[col].fillna('')
print("Ingesting raw customer nodes into BigQuery via BigQuery DataFrames...")
bf_nodes = bpd.read_pandas(df_nodes)
bf_nodes.to_gbq(table_raw_id, if_exists="replace")
print("Ingesting ground truth links into BigQuery via BigQuery DataFrames...")
df_gt = pd.DataFrame(list(true_links), columns=["source_id", "target_id"])
bf_gt = bpd.read_pandas(df_gt)
bf_gt.to_gbq(table_gt_id, if_exists="replace")
print(f"Raw customer nodes ingested into `{table_raw_id}` ({len(df_nodes):,} rows).")
print(f"Ground truth links ingested into `{table_gt_id}` ({len(df_gt):,} pairs).")
EOF
python3 ingest_febrl.py
Google Cloud コンソールで BigQuery Studio に移動し、新しい SQL クエリタブ(+)を開いて、次のクエリを実行して取り込まれた顧客ノード テーブルを検査します。
SELECT rec_id, given_name, surname, street_number, address_1, address_2, suburb, postcode, state, date_of_birth, soc_sec_id, dataset_source
FROM `identity_resolution.customer_nodes`
LIMIT 5;
次のような出力が表示されます。
rec_id | given_name | surname | street_number(番地) | address_1 | address_2 | 郊外 | 郵便番号 | state | date_of_birth | soc_sec_id | dataset_source |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| |
|
|
|
|
|
|
|
|
|
|
| |
|
|
|
|
|
|
|
|
|
| ||
|
|
|
|
|
|
|
|
|
|
|
|
ベンチマーク データセットでは、重複するクラスタ全体に現実的なダーティデータが導入されています。
- 発音とスペルのバリエーション:
brentとbrnt/bernt、woodとwoode/wod、cliftonとcliffton。 - 住所の略語と誤字脱字:
girdlestone circuitとgirdelstone circut/girdlestone cir/girdlestone crt、番地11と OCR エラー15。 - 文字の転置と欠損値: 生年月日の転置(
19340706と19340760)、州の欠損()、社会保障 ID の欠損()。
以降の手順では、SOUNDEX 音声エンコード、住所正規化 UDF、レーベンシュタイン編集距離、AI.EMBED ベクトル検索を使用して、これらの不一致を解消し、重複するプロファイルを正確にリンクします。
4. Address Validation リモート関数 UDF をデプロイする
住所の正規化では、照合を行う前に、番地、郊外の境界、郵便番号を標準化します。Google Maps Address Validation API は、住所を受け取り、住所の構成要素を識別して検証するサービスです。このステップでは、Cloud Shell に Python Cloud Functions 関数をデプロイします。この関数は、BigQuery にアドレスの検証と正規化の UDF を公開します。
Cloud Functions のソースファイルを記述する
Cloud Shell で次のコマンドを実行して、Cloud Functions のソース ディレクトリを作成し、main.py と requirements.txt を書き込みます。
mkdir -p cloud_function_address_validation && cd cloud_function_address_validation
cat << 'EOF' > main.py
import os
import json
import logging
import requests
from functools import lru_cache
from concurrent.futures import ThreadPoolExecutor
import functions_framework
import google.auth
from google.auth.transport.requests import AuthorizedSession
from requests.adapters import HTTPAdapter
from urllib3.util import Retry
ADDRESS_VALIDATION_URL = "https://addressvalidation.googleapis.com/v1:validateAddress"
ENABLE_ADDRESS_VALIDATION_API = os.environ.get("ENABLE_ADDRESS_VALIDATION_API", "false").lower() == "true"
# ==========================================
# GLOBAL INITIALIZATION (Runs once per Cold Start)
# ==========================================
credentials, _ = google.auth.default(scopes=["https://www.googleapis.com/auth/cloud-platform"])
session = AuthorizedSession(credentials)
retries = Retry(
total=4,
backoff_factor=0.5,
status_forcelist=[429, 500, 502, 503, 504],
allowed_methods=["POST"]
)
adapter = HTTPAdapter(max_retries=retries, pool_connections=100, pool_maxsize=100)
session.mount("https://", adapter)
executor = ThreadPoolExecutor(max_workers=50)
@lru_cache(maxsize=10000)
def call_validation_api(address_text: str) -> str:
payload = {
"address": {
"regionCode": "AU",
"addressLines": [address_text]
}
}
response = session.post(ADDRESS_VALIDATION_URL, json=payload, timeout=10)
response.raise_for_status()
res_data = response.json()
result = res_data.get('result', {})
address_obj = result.get('address', {})
verdict = result.get('verdict', {})
formatted = address_obj.get('formattedAddress', address_text).lower()
has_unconfirmed = verdict.get('hasUnconfirmedComponents', True)
address_complete = verdict.get('addressComplete', False)
granularity = verdict.get('validationGranularity', 'UNCONFIRMED')
actions = verdict.get('possibleNextActions', [])
next_action = str(actions[0]) if actions else "NONE"
is_valid = bool(address_complete and not has_unconfirmed)
return {
"formatted_address": formatted,
"address_is_valid": is_valid,
"validation_granularity": granularity,
"possible_next_action": next_action
}
def process_single_call(call):
call = call or []
padded = (call + [""] * 6)[:6]
cleaned_parts = [str(p).strip() if p is not None else "" for p in padded]
street_num, addr_1, addr_2, suburb, state, postcode = cleaned_parts
address_parts = [p for p in cleaned_parts if p]
address_text = " ".join(address_parts)
if not address_text:
return {
"formatted_address": "",
"address_is_valid": False,
"validation_granularity": "EMPTY",
"possible_next_action": "NONE"
}
if not ENABLE_ADDRESS_VALIDATION_API:
normalized = (
address_text.lower()
.replace("street", "st")
.replace("road", "rd")
.replace("place", "pl")
.replace("avenue", "ave")
.replace("circuit", "cct")
)
return {
"formatted_address": normalized,
"address_is_valid": bool(len(address_parts) >= 3),
"validation_granularity": "PREMISE" if postcode and suburb else "SUBURB",
"possible_next_action": "NONE"
}
try:
return call_validation_api(address_text)
except Exception as e:
logging.error(f"Address Validation API Error for '{address_text}': {str(e)}")
return {
"formatted_address": address_text.lower(),
"address_is_valid": False,
"validation_granularity": "UNCONFIRMED",
"possible_next_action": "NONE"
}
@functions_framework.http
def validate_address_udf(request):
request_json = request.get_json(silent=True) or {}
calls = request_json.get('calls', [])
if not calls:
return {'replies': []}
try:
# executor.map inherently preserves array input order (Strictly required by BigQuery)
replies = list(executor.map(process_single_call, calls))
return {'replies': replies}
except Exception as e:
logging.error(f"Batch execution failed: {e}")
return {'errorMessage': str(e)}, 400
EOF
cat << 'EOF' > requirements.txt
functions-framework==3.*
requests==2.*
google-auth==2.*
urllib3==2.*
EOF
Cloud Functions をデプロイして IAM 権限を構成する
Cloud Shell で次のコマンドを実行して、第 2 世代の Cloud Functions をデプロイし、BigQuery Cloud リソース接続を構成します。
# 1. Deploy 2nd-Gen Cloud Function
gcloud functions deploy validate_address_udf \
--gen2 \
--runtime=python311 \
--region=${REGION} \
--source=. \
--entry-point=validate_address_udf \
--trigger-http \
--no-allow-unauthenticated \
--memory=512Mi \
--cpu=1 \
--concurrency=80 \
--quiet
# 2. Extract Function Endpoint URI
export FUNCTION_URL=$(gcloud functions describe validate_address_udf --region=${REGION} --gen2 --format="value(serviceConfig.uri)")
# 3. Create BigQuery Cloud Resource Connection
bq mk --connection --location=US --project_id=${GCP_PROJECT} --connection_type=CLOUD_RESOURCE address_val_conn || true
# 4. Extract Connection Service Account Email
export BQ_SA_EMAIL=$(bq show --format=prettyjson --connection US.address_val_conn | grep -o '"serviceAccountId": "[^"]*"' | cut -d'"' -f4)
# 5. Bind Cloud Run Invoker and Vertex AI User IAM Roles to BigQuery Connection Service Account
gcloud run services add-iam-policy-binding validate-address-udf \
--region=${REGION} \
--member="serviceAccount:${BQ_SA_EMAIL}" \
--role="roles/run.invoker" --quiet
gcloud projects add-iam-policy-binding ${GCP_PROJECT} \
--member="serviceAccount:${BQ_SA_EMAIL}" \
--role="roles/aiplatform.user" --quiet
# 6. Wait for connection IAM policy propagation
echo "Waiting 60 seconds for BigQuery connection IAM policy to propagate..."
sleep 60
Cloud Functions のデプロイが完了し、IAM バインディングが正常に適用されたことを示す出力が表示されます。
リモート アドレス正規化関数を登録する
次に、BigQuery テーブルの行をデプロイされた Cloud Functions エンドポイント(${FUNCTION_URL})に接続する BigQuery リモート関数 DDL(validate_address_udf)を登録します。
Cloud Shell で次のコマンドを実行して、デプロイされた Cloud Functions の URL を取得し、リモート関数を自動的に登録します。
# 1. Retrieve deployed Cloud Function URL
export FUNCTION_URL=$(gcloud functions describe validate_address_udf --region=${REGION:-us-central1} --gen2 --format="value(serviceConfig.uri)")
# 2. Register Remote Function DDL in BigQuery
bq query --use_legacy_sql=false \
"CREATE OR REPLACE FUNCTION \`${GCP_PROJECT}.${DATASET_ID}.validate_address_udf\`(
street_number STRING,
address_1 STRING,
address_2 STRING,
suburb STRING,
state STRING,
postcode STRING
) RETURNS JSON
REMOTE WITH CONNECTION \`us.address_val_conn\`
OPTIONS (
endpoint = '${FUNCTION_URL}',
max_batching_rows = 100
);"
5. プロフィール データと音声エンコードを前処理する
このステップでは、取り込んだ customer_nodes テーブルに対して BigQuery SQL 前処理クエリを実行します。
データ クリーニングと音声特徴クエリを実行する
BigQuery Studio の SQL エディタで、次のクエリを実行して customer_nodes_cleaned を作成します。このクエリは次の処理を行います。
validate_address_udfを呼び出して、正規化された住所と検証結果を取得します。given_nameとsurnameのSOUNDEX表音エンコードを生成して、スペル バリエーションを処理します。- 構造化された
profile_textフィールドを構築します。
CREATE OR REPLACE TABLE `identity_resolution.customer_nodes_cleaned` AS
WITH raw_data AS (
SELECT
rec_id, dataset_source,
TRIM(LOWER(given_name)) AS given_name_clean,
TRIM(LOWER(surname)) AS surname_clean,
`identity_resolution.validate_address_udf`(street_number, address_1, address_2, suburb, state, postcode) AS addr_json,
TRIM(suburb) AS suburb, TRIM(state) AS state, TRIM(postcode) AS postcode,
TRIM(date_of_birth) AS date_of_birth, TRIM(soc_sec_id) AS soc_sec_id
FROM `identity_resolution.customer_nodes`
)
SELECT
rec_id, dataset_source,
given_name_clean AS given_name,
SOUNDEX(given_name_clean) AS given_name_soundex,
surname_clean AS surname,
SOUNDEX(surname_clean) AS surname_soundex,
CONCAT(given_name_clean, ' ', surname_clean) AS full_name,
STRING(addr_json.formatted_address) AS formatted_address,
BOOL(addr_json.address_is_valid) AS address_is_valid,
STRING(addr_json.validation_granularity) AS validation_granularity,
STRING(addr_json.possible_next_action) AS possible_next_action,
suburb, state, postcode, date_of_birth, soc_sec_id,
CONCAT('Name: ', CONCAT(given_name_clean, ' ', surname_clean), '; Address: ', STRING(addr_json.formatted_address), '; DOB: ', date_of_birth, '; SSN: ', soc_sec_id) AS profile_text
FROM raw_data;
クリーンアップされたノード テーブルに対してクエリを実行します。
SELECT rec_id, given_name, given_name_soundex, surname, surname_soundex, formatted_address
FROM `identity_resolution.customer_nodes_cleaned`
LIMIT 5;
次のような出力が表示されます。
rec_id | given_name | given_name_soundex | surname | surname_soundex | formatted_address |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
6. セマンティック プロファイル エンベディングとベクトル検索を生成する
BigQuery は、住所の正規化、Soundex 音標文字キー、レーベンシュタイン編集距離に加えて、AI.EMBED を介して組み込みの生成 AI エンベディング関数をサポートしています。
AI.EMBED を使用すると、BigQuery は基盤モデル(text-embedding-005 など)を使用して SQL でテキスト エンベディングを直接生成します。ベクトル インデックス DDL を手動で作成する必要はありません。
プロファイル エンベディングを生成して Top-K ベクトル検索を実行する
BigQuery Studio SQL エディタで次のクエリを実行します。
-- 1. Generate Customer Profile Embeddings using AI.EMBED (offloaded to Vertex AI)
CREATE OR REPLACE TABLE `identity_resolution.customer_embeddings` AS
SELECT
rec_id,
dataset_source,
profile_text,
AI.EMBED(profile_text, connection_id => 'us.address_val_conn', endpoint => 'text-embedding-005').result AS text_embedding
FROM `identity_resolution.customer_nodes_cleaned`;
-- 2. Execute VECTOR_SEARCH for Top-K Nearest Neighbors Candidate Generation
CREATE OR REPLACE TABLE `identity_resolution.vector_candidate_edges` AS
SELECT
query.rec_id AS source_id,
base.rec_id AS target_id,
distance AS vector_distance
FROM VECTOR_SEARCH(
TABLE `identity_resolution.customer_embeddings`,
'text_embedding',
TABLE `identity_resolution.customer_embeddings`,
top_k => 5,
distance_type => 'COSINE'
)
WHERE query.rec_id < base.rec_id AND distance <= 0.20;
7. 候補ペアのスコアリングとハイブリッド エッジ特徴の融合
データセットの規模が大きくなると、考えられるすべての顧客レコードのペアを評価する(O(N²) の二次関数的な増加)ことは計算上不可能になります。BigQuery などのリレーショナル データベース エンジンでは、複数の列にわたる複雑な OR 結合条件(a.soc_sec_id = b.soc_sec_id OR a.given_name_soundex = b.given_name_soundex OR ... での結合など)を使用してルールベースのブロックを実装しようとすると、クエリ オプティマイザーが単一の等結合キーでスケーラブルなハッシュ結合またはソートマージ結合を使用できなくなります。代わりに、エンジンは O(N²) クロス結合にフォールバックして各ペアをフィルタしますが、これは大規模な場合には失敗します。
このステップでは、ベクトル検索テーブル(vector_candidate_edges)で生成された候補ペアを取得し、高速なインデックス付き等結合(ON c.source_id = a.rec_id と ON c.target_id = b.rec_id)を使用して customer_nodes_cleaned と結合します。次に、次の要素を組み合わせて重み付けされた一致スコアを計算します。
- SSN 一致スコア(重み:
0.30) - レーベンシュタイン距離
EDIT_DISTANCEを使用した姓の編集類似度(重み:0.20) - Given Name Edit Similarity(重み:
0.20) - 生年月日の一致スコア(重み:
0.15) SPLIT(LOWER(formatted_address), ' ')の Address Token Jaccard Similarity(重み:0.15)
候補エッジと重み付けされた類似性スコアを計算する
BigQuery Studio SQL エディタで次のクエリを実行して、matched_edges にデータを入力します。
CREATE OR REPLACE TABLE `identity_resolution.matched_edges` AS
WITH candidate_pairs AS (
SELECT
c.source_id, c.target_id,
a.given_name AS a_given_name, b.given_name AS b_given_name,
a.surname AS a_surname, b.surname AS b_surname,
a.given_name_soundex AS a_gn_snd, b.given_name_soundex AS b_gn_snd,
a.surname_soundex AS a_sn_snd, b.surname_soundex AS b_sn_snd,
a.date_of_birth AS a_dob, b.date_of_birth AS b_dob,
a.soc_sec_id AS a_ssn, b.soc_sec_id AS b_ssn,
SPLIT(LOWER(a.formatted_address), ' ') AS a_tokens,
SPLIT(LOWER(b.formatted_address), ' ') AS b_tokens
FROM `identity_resolution.vector_candidate_edges` c
JOIN `identity_resolution.customer_nodes_cleaned` a ON c.source_id = a.rec_id
JOIN `identity_resolution.customer_nodes_cleaned` b ON c.target_id = b.rec_id
),
scored_pairs AS (
SELECT
source_id, target_id,
CASE WHEN a_ssn = b_ssn AND a_ssn != '' THEN 1.0 ELSE 0.0 END AS ssn_match,
CASE WHEN a_dob = b_dob THEN 1.0 ELSE 0.0 END AS dob_match,
CASE WHEN a_gn_snd = b_gn_snd THEN 1.0 ELSE 0.0 END AS given_name_soundex_match,
CASE WHEN a_sn_snd = b_sn_snd THEN 1.0 ELSE 0.0 END AS surname_soundex_match,
GREATEST(
(1.0 - (EDIT_DISTANCE(a_given_name, b_given_name) / GREATEST(LENGTH(a_given_name), LENGTH(b_given_name), 1))),
(1.0 - (EDIT_DISTANCE(a_given_name, b_surname) / GREATEST(LENGTH(a_given_name), LENGTH(b_surname), 1)))
) AS given_name_edit_sim,
GREATEST(
(1.0 - (EDIT_DISTANCE(a_surname, b_surname) / GREATEST(LENGTH(a_surname), LENGTH(b_surname), 1))),
(1.0 - (EDIT_DISTANCE(a_surname, b_given_name) / GREATEST(LENGTH(a_surname), LENGTH(b_given_name), 1)))
) AS surname_edit_sim,
(
(SELECT COUNT(DISTINCT t) FROM UNNEST(a_tokens) t JOIN UNNEST(b_tokens) t2 ON t = t2)
/
GREATEST(1.0, (SELECT COUNT(DISTINCT t) FROM UNNEST(ARRAY_CONCAT(a_tokens, b_tokens)) t))
) AS address_jaccard_sim
FROM candidate_pairs
)
SELECT
source_id, target_id,
ROUND((0.30 * ssn_match) + (0.20 * surname_edit_sim) + (0.20 * given_name_edit_sim) + (0.15 * dob_match) + (0.15 * address_jaccard_sim), 4) AS match_score
FROM scored_pairs
WHERE ((0.30 * ssn_match) + (0.20 * surname_edit_sim) + (0.20 * given_name_edit_sim) + (0.15 * dob_match) + (0.15 * address_jaccard_sim)) >= 0.55;
候補のエッジ一致を調べます。
SELECT source_id, target_id, match_score
FROM `identity_resolution.matched_edges`
ORDER BY match_score DESC;
次のような出力が表示されます。
source_id | target_id | match_score |
|
|
|
|
|
|
|
|
|
ルールベースとベクトル検索のエッジを統合テーブルに統合する
ルールベースのファジー マッチングとセマンティック ベクトル検索の候補エッジを 1 つの重複除去された final_matched_edges テーブルに結合します。
CREATE OR REPLACE TABLE `identity_resolution.final_matched_edges` AS
SELECT
source_id,
target_id,
MAX(edge_weight) AS edge_weight,
IF(COUNT(DISTINCT edge_type) > 1, 'HYBRID', MAX(edge_type)) AS edge_type
FROM (
SELECT source_id, target_id, match_score AS edge_weight, 'RULE_BASED' AS edge_type
FROM `identity_resolution.matched_edges`
UNION ALL
SELECT source_id, target_id, ROUND(1.0 - vector_distance, 4) AS edge_weight, 'VECTOR_SEARCH' AS edge_type
FROM `identity_resolution.vector_candidate_edges`
WHERE vector_distance <= 0.05
)
GROUP BY source_id, target_id;
8. プロパティ グラフの構築と ISO GQL パス トラバーサル
BigQuery は、プロパティ グラフを介して ISO GQL(Graph Query Language)をネイティブにサポートしています。プロパティ グラフは、データ重複なしで、リレーショナル BigQuery テーブル上に論理グラフビューを作成します。
このステップでは、統合された候補エッジテーブル(final_matched_edges)を使用してプロパティ グラフ customer_identity_graph を構築し、{1, 2} k-hop パス トラバーサルを使用して顧客プロファイル間のグラフ接続をクエリします。

BigQuery プロパティ グラフの DDL を作成する
BigQuery Studio の SQL エディタで次の DDL ステートメントを実行します。
CREATE OR REPLACE PROPERTY GRAPH `identity_resolution.customer_identity_graph`
NODE TABLES (
`identity_resolution.customer_nodes_cleaned` AS `Customer`
KEY (rec_id)
)
EDGE TABLES (
`identity_resolution.final_matched_edges`
KEY (source_id, target_id)
SOURCE KEY (source_id) REFERENCES `Customer`(rec_id)
DESTINATION KEY (target_id) REFERENCES `Customer`(rec_id)
LABEL MATCHED_TO
);
K-Hop グラフ クラスタを可視化する
次のクエリを実行して、1 対 2 の関係ホップで一致する顧客クラスタを可視化します。
GRAPH `identity_resolution.customer_identity_graph`
MATCH p = (c1:Customer)-[e:MATCHED_TO]->{1, 2}(c2:Customer)
RETURN TO_JSON(p) AS graph_cluster_path
LIMIT 10;

正規のお客様クラスタを解決する
次のクエリを実行して、エンティティ クラスタを resolved_customers に解決します。
CREATE OR REPLACE TABLE `identity_resolution.resolved_customers` AS
WITH graph_paths AS (
SELECT
source_node_id, target_node_id
FROM GRAPH_TABLE(
`identity_resolution.customer_identity_graph`
MATCH (c1:Customer)-[e:MATCHED_TO]->{1, 2}(c2:Customer)
COLUMNS (c1.rec_id AS source_node_id, c2.rec_id AS target_node_id)
)
),
all_connections AS (
SELECT source_node_id AS node_id, target_node_id AS connected_id FROM graph_paths
UNION DISTINCT
SELECT target_node_id AS node_id, source_node_id AS connected_id FROM graph_paths
UNION DISTINCT
SELECT rec_id AS node_id, rec_id AS connected_id FROM `identity_resolution.customer_nodes_cleaned`
),
clusters AS (
SELECT
node_id,
MIN(connected_id) AS canonical_customer_id
FROM all_connections
GROUP BY node_id
)
SELECT
canonical_customer_id,
ARRAY_AGG(node_id) AS customer_records,
COUNT(node_id) AS record_count
FROM clusters
GROUP BY canonical_customer_id;
解決済みのクラスタ テーブルに対してクエリを実行します。
SELECT canonical_customer_id, record_count, customer_records
FROM `identity_resolution.resolved_customers`
ORDER BY record_count DESC;
次のような出力が表示されます。
canonical_customer_id | record_count | customer_records |
|
|
|
|
|
|
|
|
|
評価指標ビューを作成する
ground_truth_links テーブルに対して適合率、再現率、F1 スコアを計算するには、次のコマンドを実行します。
CREATE OR REPLACE VIEW `identity_resolution.evaluation_metrics` AS
WITH predictions AS (
SELECT
LEAST(source_id, target_id) AS source_id,
GREATEST(source_id, target_id) AS target_id
FROM `identity_resolution.final_matched_edges`
WHERE edge_weight >= 0.55
),
ground_truth AS (
SELECT
LEAST(source_id, target_id) AS source_id,
GREATEST(source_id, target_id) AS target_id
FROM `identity_resolution.ground_truth_links`
),
stats AS (
SELECT
COUNT(g.source_id) AS total_ground_truth,
COUNT(p.source_id) AS total_predictions,
COUNTIF(p.source_id IS NOT NULL AND g.source_id IS NOT NULL) AS true_positives,
COUNTIF(p.source_id IS NOT NULL AND g.source_id IS NULL) AS false_positives,
COUNTIF(p.source_id IS NULL AND g.source_id IS NOT NULL) AS false_negatives
FROM ground_truth g
FULL OUTER JOIN predictions p ON g.source_id = p.source_id AND g.target_id = p.target_id
)
SELECT
total_ground_truth, total_predictions, true_positives, false_positives, false_negatives,
ROUND(true_positives / NULLIF(true_positives + false_positives, 0), 4) AS precision,
ROUND(true_positives / NULLIF(true_positives + false_negatives, 0), 4) AS recall,
ROUND(2 * true_positives / NULLIF((2 * true_positives) + false_positives + false_negatives, 0), 4) AS f1_score
FROM stats;
評価指標ビューに対してクエリを実行します。
SELECT * FROM `identity_resolution.evaluation_metrics`;
次のような出力が表示されます。
total_ground_truth | total_predictions | true_positives | false_positives | false_negatives | precision | recall | f1_score |
|
|
|
|
|
|
|
|
Adamic-Adar グラフの重み付けによる世帯クラスタの解決
個々の ID の解決では同じ人物に属するレコードが解決されますが、エンタープライズの Customer 360 アーキテクチャでは、住所を共有する同居者をグループ化する上位レベルの世帯エンティティが必要になることがよくあります。
合成ベンチマーク(FEBRL3 など)は個人レベルのグラウンド トゥルースを評価するため、世帯の解決はダウンストリーム ステップとして実行されます。タイムスタンプ付きの転居履歴がない場合、複数の住所にリンクされているユーザーが原因で、過剰な統合やクラスタの断片化が発生する可能性があります。この問題を解決するために、Adamic-Adar グラフ重み付けを使用して、ソフト世帯メンバーシップを構築します。
BigQuery Studio SQL エディタで次のクエリを実行して、Adamic-Adar グラフの重み付けを使用して household_clusters にデータを入力します。
CREATE OR REPLACE TABLE `identity_resolution.household_clusters` AS
WITH customer_addresses AS (
SELECT DISTINCT
r.canonical_customer_id,
c.formatted_address
FROM `identity_resolution.resolved_customers` r,
UNNEST(r.customer_records) AS rec_id
JOIN `identity_resolution.customer_nodes_cleaned` c ON rec_id = c.rec_id
WHERE c.formatted_address IS NOT NULL AND c.formatted_address != ''
),
-- Adamic-Adar Exclusivity Weighting: 1.0 / LN(GREATEST(degree, 2))
address_degrees AS (
SELECT
formatted_address,
COUNT(DISTINCT canonical_customer_id) AS address_degree,
1.0 / LN(GREATEST(COUNT(DISTINCT canonical_customer_id), 2)) AS address_exclusivity_weight
FROM customer_addresses
GROUP BY formatted_address
),
customer_household_affinity AS (
SELECT
ca.canonical_customer_id,
ca.formatted_address AS household_address,
ad.address_degree,
ad.address_exclusivity_weight,
ad.address_exclusivity_weight * COUNT(DISTINCT ca2.canonical_customer_id) AS raw_household_affinity
FROM customer_addresses ca
JOIN address_degrees ad ON ca.formatted_address = ad.formatted_address
LEFT JOIN customer_addresses ca2
ON ca.formatted_address = ca2.formatted_address
AND ca.canonical_customer_id != ca2.canonical_customer_id
GROUP BY ca.canonical_customer_id, ca.formatted_address, ad.address_degree, ad.address_exclusivity_weight
),
ranked_households AS (
SELECT
canonical_customer_id,
household_address,
address_degree AS total_residents,
ROUND(
COALESCE(SAFE_DIVIDE(raw_household_affinity, SUM(raw_household_affinity) OVER(PARTITION BY canonical_customer_id)), 1.0),
4
) AS household_membership_weight,
ROW_NUMBER() OVER(PARTITION BY canonical_customer_id ORDER BY raw_household_affinity DESC) AS household_rank
FROM customer_household_affinity
)
SELECT
CONCAT('hh-', ABS(FARM_FINGERPRINT(household_address))) AS canonical_household_id,
canonical_customer_id,
household_address,
total_residents,
household_membership_weight,
household_rank
FROM ranked_households;
解決済みの世帯クラスタ テーブルに対してクエリを実行します。
SELECT canonical_household_id, canonical_customer_id, household_address, total_residents, household_membership_weight, household_rank
FROM `identity_resolution.household_clusters`
ORDER BY total_residents DESC;
GQL を使用してエンドツーエンドの ID 階層を可視化する
クラスタリングされていない Raw Customers を解決済みの Customer Entities に接続し、さらに解決済みの Household Entities に接続して、3 階層の ID 階層全体を視覚的にトレースするには、BigQuery Studio で次の DDL と ISO GQL クエリを実行します。
-- 1. Create Household Node Table
CREATE OR REPLACE TABLE `identity_resolution.household_nodes` AS
SELECT DISTINCT
canonical_household_id,
household_address,
total_residents
FROM `identity_resolution.household_clusters`;
-- 2. Create Unresolved Record to Resolved Entity Edge Table
CREATE OR REPLACE TABLE `identity_resolution.customer_entity_edges` AS
SELECT DISTINCT
rec_id,
canonical_customer_id
FROM `identity_resolution.resolved_customers`,
UNNEST(customer_records) AS rec_id;
-- 3. Create Primary Household Edge Table (Highest Weighted Household Rank = 1)
CREATE OR REPLACE TABLE `identity_resolution.primary_household_edges` AS
SELECT
canonical_customer_id,
canonical_household_id,
household_membership_weight,
household_rank
FROM `identity_resolution.household_clusters`
WHERE household_rank = 1;
-- 4. Update Unified Property Graph DDL
CREATE OR REPLACE PROPERTY GRAPH `identity_resolution.customer_identity_graph`
NODE TABLES (
`identity_resolution.customer_nodes_cleaned` AS `RawCustomer`
KEY (rec_id),
`identity_resolution.resolved_customers` AS `ResolvedCustomer`
KEY (canonical_customer_id),
`identity_resolution.household_nodes` AS `ResolvedHousehold`
KEY (canonical_household_id)
)
EDGE TABLES (
`identity_resolution.final_matched_edges`
KEY (source_id, target_id)
SOURCE KEY (source_id) REFERENCES `RawCustomer`(rec_id)
DESTINATION KEY (target_id) REFERENCES `RawCustomer`(rec_id)
LABEL MATCHED_TO,
`identity_resolution.customer_entity_edges`
KEY (rec_id, canonical_customer_id)
SOURCE KEY (rec_id) REFERENCES `RawCustomer`(rec_id)
DESTINATION KEY (canonical_customer_id) REFERENCES `ResolvedCustomer`(canonical_customer_id)
LABEL RESOLVED_TO,
`identity_resolution.primary_household_edges`
KEY (canonical_customer_id, canonical_household_id)
SOURCE KEY (canonical_customer_id) REFERENCES `ResolvedCustomer`(canonical_customer_id)
DESTINATION KEY (canonical_household_id) REFERENCES `ResolvedHousehold`(canonical_household_id)
LABEL BELONGS_TO_HOUSEHOLD
);
-- 5. Execute 3-Tier GQL Query for Multi-Resident Household Visualization
GRAPH `identity_resolution.customer_identity_graph`
MATCH p = (raw:RawCustomer)-[e1:RESOLVED_TO]->(c:ResolvedCustomer)-[e2:BELONGS_TO_HOUSEHOLD]->(h:ResolvedHousehold)
WHERE h.total_residents > 1
RETURN TO_JSON(p) AS multi_resident_household_hierarchy_path
LIMIT 15;
BigQuery Studio でこの GQL クエリを実行すると、3 階層のインタラクティブなグラフ可視化キャンバスがレンダリングされます。このキャンバスには、個々の正規エンティティ(ResolvedCustomer)に解決された未加工の顧客プロファイル レコード(RawCustomer)が表示されます。これらのレコードは、共有の複数居住世帯エンティティ(ResolvedHousehold)にリンクされています。

9. 増分解決と持続的な安定性
実際のエンタープライズ アプリケーションでは、新規顧客レコードが毎日またはリアルタイムのバッチ取り込みによって継続的に到着します。増分デルタ マッチング エンジンは、履歴データセット全体でグラフの完全な解決を再実行するのではなく、新しい受信レコードを既存の解決済みベースライン クラスタ(resolved_customers)と比較します。
これを効率的に実現するために、エンジンはベクトル検索(
VECTOR_SEARCH
)を動的クラスタリングの形式として使用します。各受信レコードをクエリポイントとして扱うことで、VECTOR_SEARCH は過去のベースライン エンベディング インデックスから上位 K 個の最近傍のセットを取得します。受信レコードが類似性しきい値を超える既存の顧客プロファイルと一致する場合、そのレコードはクラスタに動的に統合され、ベースライン canonical_customer_id(MATCHED_TO_EXISTING_CLUSTER)を継承します。しきい値を超えるベースラインの最近傍が見つからない場合は、新しいエンティティ UUID が作成されます(NEW_CUSTOMER_ENTITY)。
サンプル増分バッチ取り込みレコードを取り込む
BigQuery Studio SQL エディタに次の DDL を貼り付けて実行し、incremental_daily_intake を作成します。
CREATE OR REPLACE TABLE `identity_resolution.incremental_daily_intake` AS
SELECT * FROM UNNEST([
STRUCT(
'rec-9999-new-1' AS rec_id, 'erin' AS given_name, 'donaldson' AS surname,
'E650' AS given_name_soundex, 'D543' AS surname_soundex,
'19810427' AS date_of_birth, '2955815' AS soc_sec_id,
'13 hawkesbury crescent aralee lewiston 7018' AS formatted_address, '7018' AS postcode
),
STRUCT(
'rec-9999-new-2' AS rec_id, 'hollie' AS given_name, 'lillie-hinrichs' AS surname,
'H400' AS given_name_soundex, 'L446' AS surname_soundex,
'19251130' AS date_of_birth, '4920253' AS soc_sec_id,
'27 hemmings crescent kilvinton village banyo 4030' AS formatted_address, '4030' AS postcode
),
STRUCT(
'rec-9999-new-3' AS rec_id, 'sarah' AS given_name, 'ryan' AS surname,
'S600' AS given_name_soundex, 'R500' AS surname_soundex,
'20010101' AS date_of_birth, '999999999' AS soc_sec_id,
'500 market st melbourne vic 3000' AS formatted_address, '3000' AS postcode
),
STRUCT(
'rec-9999-new-4' AS rec_id, 'zzyzx' AS given_name, 'qx-vonderland' AS surname,
'Z220' AS given_name_soundex, 'Q215' AS surname_soundex,
'19991231' AS date_of_birth, '999887766' AS soc_sec_id,
'9999 zulu orbit station moon-base alpha 9999' AS formatted_address, '9999' AS postcode
)
]);
増分 Delta 一致クエリを実行する
BigQuery SQL エディタで次のクエリを実行して、解決済みのベースライン データセットに対して差分照合を実行します。
CREATE OR REPLACE TABLE `identity_resolution.incremental_resolved_customers` AS
WITH historical_resolved_base AS (
SELECT c.rec_id, c.given_name, c.surname, c.given_name_soundex, c.surname_soundex, c.date_of_birth, c.soc_sec_id, c.formatted_address, c.postcode, r.canonical_customer_id
FROM `identity_resolution.customer_nodes_cleaned` c
JOIN (
SELECT canonical_customer_id, node_id
FROM `identity_resolution.resolved_customers`, UNNEST(customer_records) AS node_id
) r ON c.rec_id = r.node_id
),
incremental_intake AS (
SELECT
rec_id AS new_rec_id, given_name, surname, given_name_soundex, surname_soundex, date_of_birth, soc_sec_id, formatted_address, postcode,
CONCAT('Name: ', CONCAT(given_name, ' ', surname), '; Address: ', formatted_address, '; DOB: ', date_of_birth, '; SSN: ', soc_sec_id) AS profile_text
FROM `identity_resolution.incremental_daily_intake`
),
rule_delta_matches AS (
SELECT
i.new_rec_id,
h.canonical_customer_id AS matched_canonical_id,
h.rec_id AS matched_baseline_rec_id,
(
0.30 * (CASE WHEN i.soc_sec_id = h.soc_sec_id AND i.soc_sec_id != '' THEN 1.0 ELSE 0.0 END) +
0.20 * (1.0 - (EDIT_DISTANCE(i.surname, h.surname) / GREATEST(LENGTH(i.surname), LENGTH(h.surname), 1))) +
0.20 * (1.0 - (EDIT_DISTANCE(i.given_name, h.given_name) / GREATEST(LENGTH(i.given_name), LENGTH(h.given_name), 1))) +
0.15 * (CASE WHEN i.date_of_birth = h.date_of_birth THEN 1.0 ELSE 0.0 END) +
0.15 * (1.0 - (EDIT_DISTANCE(i.formatted_address, h.formatted_address) / GREATEST(LENGTH(i.formatted_address), LENGTH(h.formatted_address), 1)))
) AS match_score
FROM incremental_intake i
JOIN historical_resolved_base h
ON (i.soc_sec_id = h.soc_sec_id AND i.soc_sec_id != '')
OR (i.date_of_birth = h.date_of_birth AND i.given_name_soundex = h.given_name_soundex)
),
vector_intake_embeddings AS (
SELECT new_rec_id, AI.EMBED(profile_text, connection_id => 'us.address_val_conn', endpoint => 'text-embedding-005').result AS text_embedding
FROM incremental_intake
),
vector_delta_matches AS (
SELECT
v.query.new_rec_id,
h.canonical_customer_id AS matched_canonical_id,
h.rec_id AS matched_baseline_rec_id,
ROUND(1.0 - v.distance, 4) AS match_score
FROM VECTOR_SEARCH(
TABLE `identity_resolution.customer_embeddings`,
'text_embedding',
TABLE vector_intake_embeddings,
top_k => 3,
distance_type => 'COSINE'
) v
JOIN historical_resolved_base h ON v.base.rec_id = h.rec_id
WHERE v.distance <= 0.20
),
combined_delta AS (
SELECT
new_rec_id, matched_canonical_id, matched_baseline_rec_id, match_score,
'RULE_BASED' AS match_strategy
FROM rule_delta_matches WHERE match_score >= 0.55
UNION ALL
SELECT
new_rec_id, matched_canonical_id, matched_baseline_rec_id, match_score,
'VECTOR_SEARCH' AS match_strategy
FROM vector_delta_matches
),
aggregated_delta AS (
SELECT
new_rec_id, matched_canonical_id, matched_baseline_rec_id,
MAX(match_score) AS match_score,
CASE
WHEN COUNT(DISTINCT match_strategy) > 1 THEN 'BOTH'
ELSE MAX(match_strategy)
END AS match_strategy
FROM combined_delta
GROUP BY new_rec_id, matched_canonical_id, matched_baseline_rec_id
),
best_matches AS (
SELECT
new_rec_id, matched_canonical_id, matched_baseline_rec_id, match_score, match_strategy,
ROW_NUMBER() OVER(PARTITION BY new_rec_id ORDER BY match_score DESC) AS rank
FROM aggregated_delta
)
SELECT
i.new_rec_id AS record_id,
i.given_name, i.surname,
COALESCE(b.matched_canonical_id, GENERATE_UUID()) AS persistent_canonical_customer_id,
h.canonical_household_id AS assigned_household_id,
CASE WHEN b.matched_canonical_id IS NOT NULL THEN 'MATCHED_TO_EXISTING_CLUSTER' ELSE 'NEW_CUSTOMER_ENTITY' END AS assignment_type,
b.matched_baseline_rec_id,
b.match_score,
COALESCE(b.match_strategy, 'NONE') AS match_strategy
FROM incremental_intake i
LEFT JOIN best_matches b ON i.new_rec_id = b.new_rec_id AND b.rank = 1
LEFT JOIN `identity_resolution.primary_household_edges` h ON b.matched_canonical_id = h.canonical_customer_id;
増分解決の結果をクエリします。
SELECT record_id, persistent_canonical_customer_id, assigned_household_id, assignment_type, match_score, match_strategy
FROM `identity_resolution.incremental_resolved_customers`;
次のような出力が表示されます。
record_id | persistent_canonical_customer_id | assigned_household_id | assignment_type | match_score | match_strategy |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
10. クラスタリングの統合とクラスタの安定性(1-ε の重複)
本番環境のエンタープライズ システムでは、通常、チームは新しいエッジとデータソースを組み込むために、定期的に(毎週または毎月など)グラフ全体を再クラスタリングします。新しい関係が形成されると、グラフの完全な再クラスタリングにより、クラスタ ID がパイプライン実行間で任意にシフトまたは反転する可能性があります。
ダウンストリームの CRM、CDP、課金システムで永続的な顧客 ID を維持するため、クラスタの安定性は、(1 - ε)の重複しきい値(ε = 0.30、最小 70% のノード重複が必要)を使用して、現在の Run(t)クラスタと以前の Run(t-1)クラスタ間のノード重複を評価します。
Run(t)で新たに計算されたクラスタが、Run(t-1)のクラスタとメンバー レコードの 70% 以上を共有している場合、そのクラスタは過去の永続的な顧客 ID(STABLE_EVOLUTION)を継承します。新しいクラスタには、新たに生成された UUID(NEW_CLUSTER_CREATED)が割り当てられます。
クラスタの重複と安定性のクエリを実行する
BigQuery Studio SQL エディタで次のクエリを実行して、stable_resolved_customers にデータを入力します。
CREATE OR REPLACE TABLE `identity_resolution.stable_resolved_customers` AS
WITH previous_run_clusters AS (
SELECT
canonical_customer_id AS previous_persistent_id,
node_id
FROM `identity_resolution.resolved_customers`, UNNEST(customer_records) AS node_id
),
current_run_clusters AS (
SELECT
canonical_customer_id AS new_cluster_id,
node_id,
COUNT(*) OVER (PARTITION BY canonical_customer_id) AS current_cluster_size
FROM `identity_resolution.resolved_customers`, UNNEST(customer_records) AS node_id
UNION ALL
SELECT
persistent_canonical_customer_id AS new_cluster_id,
record_id AS node_id,
COUNT(*) OVER (PARTITION BY persistent_canonical_customer_id) AS current_cluster_size
FROM `identity_resolution.incremental_resolved_customers`
),
cluster_intersections AS (
SELECT
c.new_cluster_id,
p.previous_persistent_id,
c.current_cluster_size,
COUNT(c.node_id) AS shared_node_count,
COUNT(c.node_id) / c.current_cluster_size AS overlap_fraction
FROM current_run_clusters c
JOIN previous_run_clusters p ON c.node_id = p.node_id
GROUP BY c.new_cluster_id, p.previous_persistent_id, c.current_cluster_size
),
best_matching_previous_cluster AS (
SELECT
new_cluster_id,
previous_persistent_id,
shared_node_count,
overlap_fraction,
ROW_NUMBER() OVER (PARTITION BY new_cluster_id ORDER BY overlap_fraction DESC) AS rank
FROM cluster_intersections
WHERE overlap_fraction >= 0.70
)
SELECT
c.new_cluster_id AS raw_cluster_id,
COALESCE(b.previous_persistent_id, GENERATE_UUID()) AS persistent_canonical_customer_id,
ARRAY_AGG(c.node_id) AS customer_records,
COUNT(c.node_id) AS record_count,
COALESCE(MAX(b.shared_node_count), 0) AS shared_node_count,
COALESCE(MAX(b.overlap_fraction), 0.0) AS overlap_fraction,
CASE
WHEN b.previous_persistent_id IS NOT NULL THEN 'STABLE_EVOLUTION'
ELSE 'NEW_CLUSTER_CREATED'
END AS cluster_status
FROM current_run_clusters c
LEFT JOIN best_matching_previous_cluster b
ON c.new_cluster_id = b.new_cluster_id AND b.rank = 1
GROUP BY c.new_cluster_id, b.previous_persistent_id;
増分レコードのフィルタされたプレビューを表示する
このクエリを実行して、1 日の摂取量の記録のクラスタ安定性ステータスを確認します。
SELECT
s.persistent_canonical_customer_id,
node_id AS record_id,
h.canonical_household_id AS assigned_household_id,
s.cluster_status
FROM `identity_resolution.stable_resolved_customers` s, UNNEST(s.customer_records) AS node_id
LEFT JOIN `identity_resolution.primary_household_edges` h
ON s.persistent_canonical_customer_id = h.canonical_customer_id
WHERE node_id IN ('rec-9999-new-1', 'rec-9999-new-2', 'rec-9999-new-3', 'rec-9999-new-4')
ORDER BY record_id;
次のような出力が表示されます。
persistent_canonical_customer_id | record_id | assigned_household_id | cluster_status |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
11. クリーンアップ
Google Cloud アカウントへの継続的な課金を回避するには、デプロイされたリソースと BigQuery データセットをクリーンアップします。
Cloud Shell で、次のコマンドを実行します。
# 1. Delete BigQuery Dataset
bq rm -r -f -d ${GCP_PROJECT}:${DATASET_ID}
# 2. Delete Cloud Function (2nd-Gen)
gcloud functions delete validate_address_udf --region=${REGION} --gen2 --quiet
# 3. Delete BigQuery Cloud Connection
bq rm -f --connection US.address_val_conn
このラボ専用の Google Cloud プロジェクトを作成した場合は、プロジェクトを削除できます。
gcloud projects delete ${GCP_PROJECT}
12. 完了
おめでとうございます!BigQuery Property Graph、ISO GQL クエリ、ハイブリッド類似性照合、増分デルタ照合、永続クラスタの安定性保証を使用して、Google Cloud BigQuery 内にエンドツーエンドの顧客 ID 解決エンジンを正常に構築しました。
学習した内容
- 第 2 世代の Cloud Functions をデプロイし、BigQuery リモート関数として公開する方法。
SOUNDEX音声エンコードと住所検証を使用して顧客の属性を前処理する方法。- レーベンシュタイン距離(
EDIT_DISTANCE)とトークン ジャカード類似度を使用して、候補のブロックとハイブリッド類似度スコアの計算を実行する方法。 - ノードテーブルとエッジテーブルに BigQuery プロパティ グラフ(
CREATE PROPERTY GRAPH)を構築する方法。 {1, 2}k-hop 量化子を使用して ISO GQL(GRAPH_TABLE)でグラフパスをクエリする方法。- 標準の顧客クラスタを解決し、グラウンド トゥルース指標に対してモデルのパフォーマンスを評価する方法。
- 完全なデータセットの再処理を行わずに、毎日のバッチ取り込みの増分デルタ マッチングを実行する方法。
- パイプライン実行全体で永続的なクラスタの安定性を維持するために、(1 - ε)の重複しきい値保証を適用する方法。
次のステップ
- BigQuery プロパティ グラフのドキュメントを確認する。
- セマンティック候補の生成に BigQuery ベクトル検索と Vertex AI テキスト エンベディングを使用してみてください。
- スケーラブルな外部 API 統合については、BigQuery リモート関数をご覧ください。