1. 简介
在此 Codelab 中,您将直接在 Google Cloud BigQuery 中构建一个模块化、端到端的客户身份解析(实体匹配)引擎。您将结合使用 Google Cloud Shell 进行基础架构部署,并使用 BigQuery Studio SQL 编辑器 进行数据清理、候选评分、属性图构建和 ISO GQL(图查询语言)路径遍历。
身份解析是企业级“全面了解客户”计划、欺诈检测和多系统数据整合的基础功能。由于根据数据成熟度和业务需求,有许多有效的身份解析方法,因此本 Codelab 中的所有步骤都是模块化且可选的。该流水线旨在展示各种常见的生产级行业技术,包括远程 UDF 地址归一化、Soundex 音标分桶、语义向量搜索 (AI.EMBED)、混合特征得分和 GQL 属性图聚类,以便您可以有选择地采用适合您架构的模式。
应根据组织对确定性与概率性匹配的偏好调整匹配方法和得分阈值,而这种偏好取决于目标使用情形。例如,严格的合规性、结算或财务运营通常倾向于采用高精度的确定性规则(例如精确的社会保障号 (SSN) 或税号匹配)来防止错误关联,而营销个性化、分析和商品推荐引擎通常倾向于采用概率性模糊匹配和语义向量相似性来最大限度地提高召回率并发现细微的关联。

您将执行的操作
- 注入 FEBRL3 基准数据集:将合成客户记录和实际匹配对加载到 BigQuery 中。
- 部署地址验证远程 UDF:部署 Python Cloud Functions 函数并注册 BigQuery 远程函数以规范化街道地址。
- 预处理商家资料数据和音标编码:执行 SQL 数据清理、调用地址 UDF,并计算
SOUNDEX音标键和 Levenshtein 编辑距离:- Soundex 语音编码:一种语音算法,用于按名称在英语中的发音对其进行索引。它会将名称转换为 4 个字符的代码(一个首字母后跟三个数字),表示辅音音组(例如,
"John"和"Jon"都映射到J500,而"Smith"和"Smyth"都映射到S530),从而为特征评分和实时增量增量屏蔽提供音标匹配信号。 - Levenshtein 距离 (
EDIT_DISTANCE):一种字符串指标,用于衡量将一个字符串更改为另一个字符串所需的最少单字符编辑次数(插入、删除或替换),从而实现精确的模糊名称和地址匹配。
- Soundex 语音编码:一种语音算法,用于按名称在英语中的发音对其进行索引。它会将名称转换为 4 个字符的代码(一个首字母后跟三个数字),表示辅音音组(例如,
- 生成语义个人资料嵌入和 Vector Search:使用
AI.EMBED(text-embedding-005) 直接在 SQL 中生成文本嵌入,并使用VECTOR_SEARCH查找前 K 个最近邻,以用作次线性候选集生成层。 - 候选对评分和混合边特征融合:利用向量搜索候选对消除 O(N²) 交叉联接复杂性,计算多特征加权相似度得分(SSN、Levenshtein 编辑距离、出生日期、地址 Jaccard),并将边融合到统一的候选表中。
- 属性图构建和 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 版预留:
# 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 客户节点数据集
在部署地址验证远程函数并执行身份解析之前,您将使用 Python 的 recordlinkage 库加载合成的 FEBRL3 实体解析基准数据集(其中包含 5,000 条客户记录,每个客户最多有 5 个重复项的多重复集群),并使用 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 | - | street_number | address_1 | address_2 | 郊区 | 邮政编码 | 州 | 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、Levenshtein 编辑距离和 AI.EMBED 向量搜索来弥合这些差异,并准确关联重复的个人资料。
4. 部署地址验证远程函数 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 远程函数 DDL (validate_address_udf),该 DDL 可将 BigQuery 表行连接到已部署的 Cloud Functions 端点 (${FUNCTION_URL})。
在 Cloud Shell 中运行以下命令,以检索已部署的 Cloud Functions 函数网址并自动注册远程函数:
# 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_soundex | formatted_address |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
6. 生成语义个人资料嵌入和向量搜索
除了地址归一化、Soundex 音标键和 Levenshtein 编辑距离之外,BigQuery 还通过 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 等关系型数据库引擎中,尝试使用多个列(例如基于 a.soc_sec_id = b.soc_sec_id OR a.given_name_soundex = b.given_name_soundex OR ... 进行联接)的复杂 OR 联接条件来实现基于规则的屏蔽,会阻止查询优化器在单个等值联接键上使用可扩容的哈希联接或排序合并联接。相反,引擎会回退到 O(N²) 交叉联接并过滤每个对,这在规模上会失败。
在此步骤中,您将获取由向量搜索表 (vector_candidate_edges) 生成的候选配对,并通过快速的索引等值联接(ON c.source_id = a.rec_id 和 ON c.target_id = b.rec_id)将它们与 customer_nodes_cleaned 联接起来。然后,您将计算一个加权匹配得分,该得分结合了以下因素:
- 社会保障号匹配得分(权重:
0.30) - 使用 Levenshtein 距离
EDIT_DISTANCE的姓氏编辑相似度(权重:0.20) - Given Name Edit Similarity(权重:
0.20) - 出生日期匹配得分(权重:
0.15) - 地址令牌 Jaccard 相似度(权重:
0.15),超过SPLIT(LOWER(formatted_address), ' ')
计算候选边和加权相似度得分
在 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 |
|
|
|
|
|
|
|
|
|
将基于规则的边和向量搜索边融合到统一的表中
将基于规则的模糊匹配和语义向量搜索的候选边合并到单个去重后的 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(图查询语言)。属性图可在关系型 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 | 精确度 | 召回 | f1_score |
|
|
|
|
|
|
|
|
通过 Adamic-Adar 图加权解决家庭聚类问题
虽然个人身份解析可以解析属于同一人的记录,但企业客户 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 可视化端到端身份层次结构
如需直观地跟踪完整的三层身份层次结构(将未聚类的原始客户连接到已解析的客户实体,再连接到已解析的家庭实体),请在 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 查询会呈现一个三层交互式图表可视化画布,其中显示了已解析为各个规范实体 (ResolvedCustomer) 的原始客户个人资料记录 (RawCustomer),这些规范实体与共享的多住户家庭实体 (ResolvedHousehold) 相关联。

9. 增量分辨率和持久稳定性
在实际的企业应用中,每天或实时批量接收新客户记录。增量增量匹配引擎不会针对整个历史数据集重新运行完整图表解析,而是将新传入的记录与现有的已解析基准聚类 (resolved_customers) 进行比较。
为了高效实现这一目标,该引擎使用Vector Search (
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-ε 重叠)
在生产企业系统中,团队通常会定期(例如每周或每月)重新聚类整个图,以纳入新的边和数据源。随着新关系的形成,完整的图重新聚类可能会导致聚类标识符在流水线执行过程中任意移动或翻转。
为了保持下游 CRM、CDP 和结算系统的持久性客户 ID,集群稳定性会使用 (1 - ε) 的重叠阈值(其中 ε = 0.30,要求节点重叠率至少为 70%)评估当前运行 (t) 集群与上一次运行 (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;
查看增量记录的过滤预览
执行此查询,以验证每日摄入量记录的集群稳定性状态:
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. 恭喜
恭喜!您已成功在 Google Cloud BigQuery 中构建了一个端到端客户身份解析引擎,该引擎使用 BigQuery 属性图、ISO GQL 查询、混合相似度匹配、增量增量匹配和持久性集群稳定性保证。
您学到的内容
- 如何部署第 2 代 Cloud Functions 函数并将其公开为 BigQuery 远程函数。
- 如何使用
SOUNDEX语音编码和地址验证来预处理客户人口统计信息。 - 如何使用 Levenshtein 距离 (
EDIT_DISTANCE) 和令牌 Jaccard 相似度执行候选屏蔽并计算混合相似度得分。 - 如何基于节点表和边表构建 BigQuery 属性图 (
CREATE PROPERTY GRAPH)。 - 如何使用 ISO GQL (
GRAPH_TABLE) 和{1, 2}k-hop 量词查询图路径。 - 如何解决规范的客户聚类问题,并根据标准答案指标评估模型性能。
- 如何在不重新处理完整数据集的情况下,针对每日批量接收的数据执行增量增量匹配。
- 如何应用 (1 - ε) 重叠阈值保证,以在流水线运行期间保持持久的集群稳定性。
后续步骤
- 探索 BigQuery 属性图文档。
- 尝试使用 BigQuery Vector Search 和 Vertex AI 文本嵌入来生成语义候选内容。
- 如需了解可扩缩的外部 API 集成,请参阅 BigQuery 远程函数。