1. Introduzione
In questo codelab, creerai un motore modulare e end-to-end di risoluzione dell'identità del cliente (corrispondenza delle entità) direttamente all'interno di Google Cloud BigQuery. Combinerai Google Cloud Shell per il deployment dell'infrastruttura con l'editor SQL di BigQuery Studio per la pulizia dei dati, l'assegnazione di punteggi ai candidati, la creazione di grafici delle proprietà e gli attraversamenti di percorsi GQL (Graph Query Language) ISO.
La risoluzione dell'identità è una funzionalità fondamentale per Customer 360 aziendale, il rilevamento delle frodi e il consolidamento dei dati di più sistemi. Poiché esistono molti approcci validi alla risoluzione dell'identità a seconda della maturità dei dati e delle esigenze aziendali, tutti i passaggi di questo codelab sono modulari e facoltativi. La pipeline è progettata per mostrare una serie di tecniche industriali comuni di livello di produzione, tra cui la normalizzazione degli indirizzi delle funzioni definite dall'utente remote, il blocco fonetico Soundex, la ricerca vettoriale semantica (AI.EMBED), l'assegnazione di punteggi alle funzionalità ibride e il clustering dei grafici delle proprietà GQL, in modo da poter adottare selettivamente i pattern adatti alla tua architettura.
I metodi di corrispondenza e le soglie di punteggio devono essere ottimizzati in base alla propensione della tua organizzazione alla corrispondenza deterministica e probabilistica, che è dettata dal caso d'uso di destinazione. Ad esempio, le operazioni di conformità, fatturazione o finanziarie rigorose in genere favoriscono regole deterministiche di alta precisione (come corrispondenze esatte di SSN o ID fiscale) per impedire collegamenti errati, mentre i motori di personalizzazione del marketing, analisi e raccomandazione spesso si basano su corrispondenze probabilistiche fuzzy e similarità vettoriale semantica per massimizzare il recupero e scoprire connessioni sottili.

In questo lab proverai a:
- Importa il set di dati di benchmark FEBRL3: carica record cliente sintetici e coppie di corrispondenze di dati empirici reali in BigQuery.
- Deploy Address Validation Remote UDF: esegui il deployment di una funzione Cloud Python e registra una funzione remota BigQuery per normalizzare gli indirizzi stradali.
- Pre-elabora i dati del profilo e le codifiche fonetiche: esegui la pulizia dei dati SQL, richiama la UDF dell'indirizzo e calcola le chiavi fonetiche
SOUNDEXe le distanze di modifica di Levenshtein:- Codifica fonetica Soundex: un algoritmo fonetico per l'indicizzazione dei nomi in base al suono pronunciato in inglese. Converte i nomi in un codice di quattro caratteri (una lettera iniziale seguita da tre cifre) che rappresentano gruppi di suoni consonantici (ad es. sia
"John"che"Jon"corrispondono aJ500, mentre"Smith"e"Smyth"corrispondono aS530), fornendo indicatori di corrispondenza fonetica per l'assegnazione del punteggio delle funzionalità e il blocco incrementale in tempo reale. - Distanza di Levenshtein (
EDIT_DISTANCE): una metrica delle stringhe che misura il numero minimo di modifiche di un singolo carattere (inserimenti, eliminazioni o sostituzioni) necessarie per trasformare una stringa in un'altra, consentendo una corrispondenza approssimativa precisa di nomi e indirizzi.
- Codifica fonetica Soundex: un algoritmo fonetico per l'indicizzazione dei nomi in base al suono pronunciato in inglese. Converte i nomi in un codice di quattro caratteri (una lettera iniziale seguita da tre cifre) che rappresentano gruppi di suoni consonantici (ad es. sia
- Genera embedding del profilo semantico e Vector Search: genera text embedding direttamente in SQL utilizzando
AI.EMBED(text-embedding-005) e trova i Top-K vicini più prossimi utilizzandoVECTOR_SEARCHper fungere da livello di generazione di candidati sublineare. - Punteggio delle coppie di candidati e fusione ibrida delle funzionalità perimetrali: sfrutta le coppie di candidati della ricerca vettoriale per eliminare la complessità del cross-join O(N²), calcola i punteggi di similarità ponderata multifunzione (SSN, distanza di modifica di Levenshtein, data di nascita, Jaccard dell'indirizzo) e unisci i bordi in una tabella di candidati unificata.
- Costruzione di grafici delle proprietà e attraversamenti di percorsi ISO GQL: costruisci un
PROPERTY GRAPHBigQuery, esegui query sui percorsi ISO GQL{1, 2}(GRAPH_TABLE) per risolvere i cluster di clienti connessi, calcola le metriche di valutazione individuali ed esegui il clustering soft delle famiglie utilizzando la ponderazione del grafico Adamic-Adar. - Risoluzione incrementale e stabilità persistente: elabora le importazioni batch giornaliere con la corrispondenza delta incrementale.
- Consolidamento del clustering e stabilità del cluster (sovrapposizione 1-ε): garantisci la stabilità persistente del cluster durante le esecuzioni della pipeline utilizzando una garanzia di soglia di sovrapposizione (1-ε).
Che cosa ti serve
- Un browser web come Chrome.
- Un progetto Google Cloud con la fatturazione abilitata.
Questo codelab è progettato per data engineer, sviluppatori di database e professionisti dell'AI/ML di tutti i livelli, inclusi i principianti.
Durata stimata: 45 minuti
Costo stimato: meno di 2 € (utilizza l'elaborazione di query Cloud Functions e BigQuery con pagamento a consumo).
2. Prima di iniziare
Crea un progetto Google Cloud
- Nella console Google Cloud, nella pagina di selezione del progetto, seleziona o crea un progetto Google Cloud.
- Verifica che la fatturazione sia attivata per il tuo progetto Cloud. Scopri come verificare se la fatturazione è abilitata per un progetto.
Avvia Cloud Shell
Cloud Shell è un ambiente a riga di comando in esecuzione in Google Cloud che viene precaricato con gli strumenti necessari.
- Fai clic su Attiva Cloud Shell nella parte superiore della console Google Cloud.
- Verifica l'autenticazione:
gcloud auth list
- Configura le variabili di ambiente in Cloud Shell:
export GCP_PROJECT=$(gcloud config get-value project)
export REGION="us-central1"
export DATASET_ID="identity_resolution"
Abilita le API richieste
Esegui questo comando in Cloud Shell utilizzando il tuo account utente per abilitare tutti i servizi Google Cloud richiesti:
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
Configura il service account e la simulazione dell'identità (consigliato)
Per garantire l'esecuzione perfetta delle API e l'accesso alle credenziali predefinite dell'applicazione (ADC), crea un service account lab dedicato e abilita l'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
Crea set di dati BigQuery
Crea il set di dati BigQuery per archiviare i nodi, gli archi, i modelli di grafici e le viste di valutazione dei clienti:
bq mk --location=US --dataset ${GCP_PROJECT}:${DATASET_ID}
Dovresti visualizzare un output simile al seguente:
Dataset 'your-project-id:identity_resolution' successfully created.
Crea prenotazione e assegnazione BigQuery (facoltativo / consigliato)
Per garantire una capacità di calcolo dedicata per le ricerche di indici vettoriali, le aggregazioni di grafi e l'esecuzione di funzioni remote senza essere vincolato dai limiti della CPU on demand o dalle quote condivise, crea una prenotazione Enterprise Edition con scalabilità automatica in Cloud Shell:
# 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. Importa il set di dati del nodo cliente FEBRL3
Prima di eseguire il deployment della funzione remota di convalida degli indirizzi ed eseguire la risoluzione delle identità, caricherai il set di dati di benchmark per la risoluzione delle entità FEBRL3 sintetico (che contiene 5000 record di clienti con cluster multiduplicati fino a 5 duplicati per cliente) utilizzando la libreria recordlinkage di Python e scriverai i nodi cliente non elaborati (customer_nodes) e i link di corrispondenza di riferimento (ground_truth_links) in BigQuery utilizzando BigQuery DataFrames (bigframes).
Esegui questi comandi in Cloud Shell per installare le dipendenze ed eseguire lo script di importazione:
# 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
Nella console Google Cloud, vai a BigQuery Studio, apri una nuova scheda query SQL (+) ed esegui la query riportata di seguito per esaminare la tabella dei nodi dei clienti inseriti:
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;
Dovresti visualizzare un output simile al seguente:
rec_id | given_name | cognome | street_number | address_1 | address_2 | periferia | codice postale | stato | date_of_birth | soc_sec_id | dataset_source |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| |
|
|
|
|
|
|
|
|
|
|
| |
|
|
|
|
|
|
|
|
|
| ||
|
|
|
|
|
|
|
|
|
|
|
|
Nota come il set di dati di benchmark introduce dati sporchi realistici nei cluster duplicati:
- Varianti fonetiche e ortografiche:
brentrispetto abrnt/bernt,woodrispetto awoode/wodecliftonrispetto acliffton. - Abbreviazioni e errori di battitura negli indirizzi:
girdlestone circuitanzichégirdelstone circut/girdlestone cir/girdlestone crte numero civico11anziché errore OCR15. - Trasposizioni di caratteri e valori mancanti: trasposizioni della data di nascita (
19340706vs.19340760), stati mancanti ( ) e ID Social Security mancanti ( ).
Nei passaggi successivi, utilizzerai le SOUNDEX codifiche fonetiche, le funzioni definite dall'utente per la normalizzazione degli indirizzi, la distanza di modifica di Levenshtein e la ricerca vettoriale AI.EMBED per colmare queste discrepanze e collegare con precisione i profili duplicati.
4. Esegui il deployment della funzione remota UDF di Address Validation
La normalizzazione degli indirizzi standardizza i nomi delle vie, i confini dei sobborghi e i codici postali prima di eseguire la corrispondenza. L'API Google Maps Address Validation è un servizio che accetta un indirizzo, ne identifica i componenti e li convalida. In questo passaggio, eseguirai il deployment di una funzione Cloud Python in Cloud Shell che espone una UDF di convalida e normalizzazione degli indirizzi a BigQuery.
Scrivi i file sorgente della funzione Cloud
Esegui il comando seguente in Cloud Shell per creare la directory di origine di Cloud Functions e scrivere main.py e 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
Esegui il deployment della funzione Cloud e configura le autorizzazioni IAM
Esegui questi comandi in Cloud Shell per eseguire il deployment della funzione Cloud Functions di seconda generazione e configurare una connessione alle risorse Cloud BigQuery:
# 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
Dovresti visualizzare un output che indica che il deployment della funzione Cloud è stato completato e che i binding IAM sono stati applicati correttamente.
Registra la funzione di normalizzazione dell'indirizzo remoto
Ora registrerai il DDL della funzione remota BigQuery (validate_address_udf) che collega le righe della tabella BigQuery all'endpoint Cloud Functions di cui è stato eseguito il deployment (${FUNCTION_URL}).
Esegui questo comando in Cloud Shell per recuperare l'URL della funzione Cloud Functions di cui è stato eseguito il deployment e registrare automaticamente la funzione remota:
# 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. Pre-elaborare i dati del profilo e le codifiche fonetiche
In questo passaggio, eseguirai una query di pre-elaborazione SQL di BigQuery sulla tabella customer_nodes importata.
Esegui la pulizia dei dati e la query delle funzionalità fonetiche
Nell'editor SQL di BigQuery Studio, esegui la query riportata di seguito per creare customer_nodes_cleaned. Questa query:
- Chiamate
validate_address_udfper ottenere indirizzi normalizzati e verdetti di convalida. - Genera codifiche fonetiche
SOUNDEXpergiven_nameesurnameper gestire le varianti ortografiche. - Crea un campo
profile_textstrutturato.
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;
Esegui una query sulla tabella dei nodi pulita:
SELECT rec_id, given_name, given_name_soundex, surname, surname_soundex, formatted_address
FROM `identity_resolution.customer_nodes_cleaned`
LIMIT 5;
Dovresti visualizzare un output simile al seguente:
rec_id | given_name | given_name_soundex | cognome | surname_soundex | formatted_address |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
6. Generare embedding del profilo semantico e Vector Search
Oltre alla normalizzazione degli indirizzi, alle chiavi fonetiche Soundex e alla distanza di modifica Levenshtein, BigQuery supporta le funzioni di incorporamento dell'AI generativa integrate tramite AI.EMBED.
Utilizzando AI.EMBED, BigQuery genera incorporamenti di testo direttamente in SQL utilizzando modelli di base (come text-embedding-005) senza richiedere DDL di indice vettoriale manuale:
Generare incorporamenti dei profili ed eseguire la ricerca vettoriale Top-K
Esegui le query riportate di seguito nell'editor SQL di BigQuery Studio:
-- 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. Punteggio delle coppie di candidati e fusione delle funzionalità ibride Edge
La valutazione di tutte le possibili coppie di record cliente (crescita quadratica O(N²)) diventa proibitiva dal punto di vista computazionale man mano che la scala del set di dati aumenta. Nei motori di database relazionali come BigQuery, il tentativo di implementare il blocco basato su regole utilizzando condizioni di join OR complesse su più colonne (ad esempio il join su a.soc_sec_id = b.soc_sec_id OR a.given_name_soundex = b.given_name_soundex OR ...) impedisce all'ottimizzatore di query di utilizzare join hash scalabili o join sort-merge su una singola chiave di equi-join. Al contrario, il motore esegue un cross join O(N²) e filtra ogni coppia, il che non funziona su larga scala.
In questo passaggio, prenderai le coppie di candidati generate dalla nostra tabella di ricerca vettoriale (vector_candidate_edges) e le unirai a customer_nodes_cleaned tramite equi-join veloci e indicizzati (ON c.source_id = a.rec_id e ON c.target_id = b.rec_id). Dopodiché, calcolerai un punteggio di corrispondenza ponderato combinando:
- Punteggio di corrispondenza del numero di previdenza sociale (peso:
0.30) - Somiglianza di modifica del cognome utilizzando la distanza di Levenshtein
EDIT_DISTANCE(peso:0.20) - Given Name Edit Similarity (Similarità di modifica del nome) (peso:
0.20) - Punteggio corrispondenza data di nascita (peso:
0.15) - Similarità di Jaccard dei token indirizzo (peso:
0.15) suSPLIT(LOWER(formatted_address), ' ')
Calcola i bordi candidati e i punteggi di somiglianza ponderati
Esegui la seguente query nell'editor SQL di BigQuery Studio per compilare 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;
Ispeziona le corrispondenze esatte candidate:
SELECT source_id, target_id, match_score
FROM `identity_resolution.matched_edges`
ORDER BY match_score DESC;
Dovresti visualizzare un output simile al seguente:
source_id | target_id | match_score |
|
|
|
|
|
|
|
|
|
Unire gli archi della ricerca basata su regole e Vector Search in una tabella unificata
Combina i bordi candidati della corrispondenza fuzzy basata su regole e della ricerca vettoriale semantica in una singola tabella final_matched_edges deduplicata:
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. Costruzione di grafici delle proprietà e attraversamenti di percorsi ISO GQL
BigQuery supporta ISO GQL (Graph Query Language) in modo nativo tramite grafi delle proprietà. Un grafo a proprietà crea una visualizzazione logica del grafo sulle tabelle BigQuery relazionali senza duplicazione dei dati.
In questo passaggio, creerai un grafico delle proprietà customer_identity_graph utilizzando la tabella degli archi dei candidati unificata (final_matched_edges) ed eseguirai query sulle connessioni del grafico tra i profili dei clienti utilizzando l'attraversamento del percorso k-hop {1, 2}.

Crea DDL del grafico delle proprietà BigQuery
Esegui la seguente istruzione DDL nell'editor SQL di BigQuery Studio:
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
);
Visualizzare i cluster del grafico K-hop
Esegui la query riportata di seguito per visualizzare i cluster di clienti corrispondenti in 1-2 hop di relazione:
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;

Risolvi i cluster di clienti canonici
Esegui la seguente query per risolvere i cluster di entità in 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;
Esegui una query sulla tabella dei cluster risolti:
SELECT canonical_customer_id, record_count, customer_records
FROM `identity_resolution.resolved_customers`
ORDER BY record_count DESC;
Dovresti visualizzare un output simile al seguente:
canonical_customer_id | record_count | customer_records |
|
|
|
|
|
|
|
|
|
Crea la visualizzazione delle metriche di valutazione
Per calcolare precisione, richiamo e punteggio F1 rispetto alla tabella ground_truth_links, esegui:
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;
Esegui una query sulla visualizzazione delle metriche di valutazione:
SELECT * FROM `identity_resolution.evaluation_metrics`;
Dovresti visualizzare un output simile al seguente:
total_ground_truth | total_predictions | true_positives | false_positives | false_negatives | precisione | richiamo | f1_score |
|
|
|
|
|
|
|
|
Risolvere i cluster di nuclei familiari tramite la ponderazione del grafico di Adamic-Adar
Mentre la risoluzione dell'identità individuale risolve i record appartenenti alla stessa persona, le architetture Customer 360 aziendali spesso richiedono un raggruppamento di entità nucleo familiare di livello superiore per le persone che condividono un indirizzo.
Poiché i benchmark sintetici (come FEBRL3) valutano i dati di fatto a livello individuale, la risoluzione a livello di nucleo familiare viene eseguita come passaggio successivo. In assenza di una cronologia dei trasferimenti con timestamp, le persone collegate a più indirizzi potrebbero causare un'unione eccessiva o una frammentazione dei cluster. Per risolvere questo problema, utilizziamo la ponderazione del grafico Adamic-Adar per creare appartenenze soft al nucleo familiare.
Esegui la query riportata di seguito nell'editor SQL di BigQuery Studio per popolare household_clusters utilizzando la ponderazione del grafico Adamic-Adar:
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;
Esegui una query sulla tabella dei cluster di nuclei familiari risolti:
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;
Visualizzare la gerarchia delle identità end-to-end tramite GQL
Per tracciare visivamente la gerarchia completa delle identità a tre livelli, collegando i Clienti non raggruppati alle Entità cliente risolte e a loro volta alle Entità nucleo familiare risolte, esegui la seguente query DDL e ISO GQL in BigQuery Studio:
-- 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;
L'esecuzione di questa query GQL in BigQuery Studio esegue il rendering di un canvas di visualizzazione del grafico interattivo a tre livelli che mostra i record del profilo cliente non elaborati (RawCustomer) risolti in singole entità canoniche (ResolvedCustomer), che sono collegate a entità di nuclei familiari con più residenti condivisi (ResolvedHousehold).

9. Risoluzione incrementale e stabilità persistente
Nelle applicazioni aziendali reali, i nuovi record dei clienti arrivano continuamente tramite importazioni batch giornaliere o in tempo reale. Anziché eseguire nuovamente la risoluzione completa del grafico sull'intero set di dati storici, un motore di corrispondenza delta incrementale confronta i nuovi record in arrivo con i cluster di base risolti esistenti (resolved_customers).
Per ottenere questo risultato in modo efficiente, il motore utilizza Vector Search (
VECTOR_SEARCH
) come forma di clustering dinamico. Trattando ogni record in entrata come un punto di query, VECTOR_SEARCH recupera l'insieme dei K vicini più prossimi dall'indice di incorporamento di base storico. Se un record in entrata corrisponde a un profilo cliente esistente al di sopra della soglia di somiglianza, viene unito dinamicamente a quel cluster ed eredita la baseline canonical_customer_id (MATCHED_TO_EXISTING_CLUSTER). Se non viene trovato alcun vicino più prossimo della baseline al di sopra della soglia, viene creato un nuovo UUID entità (NEW_CUSTOMER_ENTITY).
Importa record di importazione batch incrementale di esempio
Incolla ed esegui il seguente DDL nell'editor SQL di BigQuery Studio per creare 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
)
]);
Esegui query di corrispondenza delta incrementale
Esegui la seguente query nell'editor SQL di BigQuery per eseguire la corrispondenza delta rispetto al set di dati di base risolto:
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;
Esegui una query sui risultati della risoluzione incrementale:
SELECT record_id, persistent_canonical_customer_id, assigned_household_id, assignment_type, match_score, match_strategy
FROM `identity_resolution.incremental_resolved_customers`;
Dovresti visualizzare un output simile al seguente:
record_id | persistent_canonical_customer_id | assigned_household_id | assignment_type | match_score | match_strategy |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
10. Consolidamento del clustering e stabilità del cluster (sovrapposizione 1-ε)
Nei sistemi aziendali di produzione, i team in genere raggruppano nuovamente l'intero grafico su base ricorrente (ad es. settimanale o mensile) per incorporare nuovi archi e origini dati. Man mano che si formano nuove relazioni, il raggruppamento completo del grafico può causare lo spostamento o l'inversione arbitraria degli identificatori dei cluster nelle esecuzioni della pipeline.
Per mantenere gli ID cliente persistenti per i sistemi CRM, CDP e di fatturazione downstream, Stabilità cluster valuta la sovrapposizione dei nodi tra i cluster dell'esecuzione corrente (t) e quelli dell'esecuzione precedente (t-1) utilizzando una soglia di sovrapposizione (1 - ε) (dove ε = 0, 30, che richiede una sovrapposizione minima del 70% dei nodi).
Se un cluster appena calcolato in Esecuzione (t) condivide almeno il 70% dei record dei membri con un cluster di Esecuzione (t-1), eredita l'ID cliente persistente storico (STABLE_EVOLUTION). I cluster nuovi di zecca ricevono UUID appena generati (NEW_CLUSTER_CREATED).
Esegui la query su sovrapposizione e stabilità dei cluster
Esegui la seguente query nell'editor SQL di BigQuery Studio per compilare 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;
Visualizzare l'anteprima filtrata per i record incrementali
Esegui questa query per verificare lo stato di stabilità del cluster per i record di assunzione giornaliera:
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;
Dovresti visualizzare un output simile al seguente:
persistent_canonical_customer_id | record_id | assigned_household_id | cluster_status |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
11. Esegui la pulizia
Per evitare addebiti continui al tuo account Google Cloud, libera spazio dalle risorse di cui è stato eseguito il deployment e dal set di dati BigQuery.
In Cloud Shell, esegui:
# 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
Se hai creato un progetto Google Cloud dedicato per questo lab, puoi eliminarlo:
gcloud projects delete ${GCP_PROJECT}
12. Complimenti
Complimenti! Hai creato correttamente un motore di risoluzione dell'identità del cliente end-to-end all'interno di Google Cloud BigQuery utilizzando BigQuery Property Graph, query ISO GQL, corrispondenza di similarità ibrida, corrispondenza delta incrementale e garanzie di stabilità del cluster persistente.
Cosa hai imparato
- Come eseguire il deployment di una funzione Cloud Functions di seconda generazione ed esporla come funzione remota BigQuery.
- Come pre-elaborare i dati demografici dei clienti utilizzando le codifiche fonetiche
SOUNDEXe la convalida dell'indirizzo. - Come eseguire il blocco dei candidati e calcolare i punteggi di similarità ibrida utilizzando la distanza di Levenshtein (
EDIT_DISTANCE) e la similarità di Jaccard dei token. - Come costruire un BigQuery Property Graph (
CREATE PROPERTY GRAPH) su tabelle di nodi e archi. - Come eseguire query sui percorsi del grafo utilizzando GQL (
GRAPH_TABLE) ISO con quantificatori di{1, 2}k-hop. - Come risolvere i cluster di clienti canonici e valutare le prestazioni del modello rispetto alle metriche basate su dati empirici reali.
- Come eseguire la corrispondenza delta incrementale per le importazioni batch giornaliere senza rielaborare l'intero set di dati.
- Come applicare una garanzia di soglia di sovrapposizione (1 - ε) per mantenere la stabilità persistente del cluster durante le esecuzioni della pipeline.
Passaggi successivi
- Esplora la documentazione di BigQuery Property Graph.
- Prova a utilizzare la ricerca vettoriale di BigQuery e gli incorporamenti di testo di Vertex AI per la generazione di candidati semantici.
- Scopri di più sulle funzioni remote di BigQuery per integrazioni scalabili di API esterne.