Résolution de l'identité client avec BigQuery Graph

1. Introduction

Dans cet atelier de programmation, vous allez créer un moteur modulaire de résolution de l'identité client (mise en correspondance d'entités) de bout en bout directement dans Google Cloud BigQuery. Vous allez combiner Google Cloud Shell pour le déploiement de l'infrastructure avec l'éditeur SQL BigQuery Studio pour le nettoyage des données, la notation des candidats, la construction du graphique de propriétés et les traversées de chemin GQL (Graph Query Language) ISO.

La résolution de l'identité est une fonctionnalité fondamentale pour la vision à 360° du client dans les entreprises, la détection des fraudes et la consolidation des données multisystèmes. Comme il existe de nombreuses approches valides pour la résolution de l'identité en fonction de la maturité des données et des besoins de l'entreprise, toutes les étapes de cet atelier de programmation sont modulaires et facultatives. Le pipeline est conçu pour présenter diverses techniques courantes de niveau production, y compris la normalisation d'adresses UDF à distance, le blocage phonétique Soundex, la recherche vectorielle sémantique (AI.EMBED), le scoring hybride des caractéristiques et le clustering de graphiques de propriétés GQL. Vous pouvez ainsi adopter de manière sélective les modèles qui correspondent à votre architecture.

Les méthodes de mise en correspondance et les seuils de score doivent être ajustés en fonction de l'appétit de votre organisation pour la mise en correspondance déterministe ou probabiliste, qui est dicté par le cas d'utilisation cible. Par exemple, les opérations strictes de conformité, de facturation ou financières privilégient généralement les règles déterministes de haute précision (comme les correspondances exactes de numéros de sécurité sociale ou de numéros d'identification fiscale) pour éviter les faux liens. En revanche, les moteurs de personnalisation marketing, d'analyse et de recommandations s'appuient souvent sur la mise en correspondance floue probabiliste et la similarité vectorielle sémantique pour maximiser le rappel et découvrir des liens subtils.

Architecture du moteur de résolution de l'identité client BigQuery

Objectifs de l'atelier

  • Ingérer l'ensemble de données de référence FEBRL3 : chargez des enregistrements client synthétiques et des paires correspondantes de référence dans BigQuery.
  • Déployer la fonction définie par l'utilisateur distante Address Validation : déployez une fonction Cloud Python et enregistrez une fonction distante BigQuery pour normaliser les adresses postales.
  • Prétraiter les données de profil et les encodages phonétiques : exécutez le nettoyage des données SQL, appelez l'UDF d'adresse et calculez les clés phonétiques SOUNDEX et les distances d'édition de Levenshtein :
    • Encodage phonétique Soundex : algorithme phonétique permettant d'indexer les noms par leur prononciation en anglais. Il convertit les noms en un code de quatre caractères (une lettre initiale suivie de trois chiffres) représentant des groupes de sons consonantiques (par exemple, "John" et "Jon" sont mappés sur J500, tandis que "Smith" et "Smyth" sont mappés sur S530). Il fournit ainsi des signaux de correspondance phonétique pour le scoring des caractéristiques et le blocage incrémental en temps réel des deltas.
    • Distance de Levenshtein (EDIT_DISTANCE) : métrique de chaîne mesurant le nombre minimal de modifications d'un seul caractère (insertions, suppressions ou substitutions) nécessaires pour transformer une chaîne en une autre. Elle permet une mise en correspondance floue précise des noms et des adresses.
  • Générer des embeddings de profil sémantique et effectuer une recherche vectorielle : générez des embeddings de texte directement en SQL à l'aide de AI.EMBED (text-embedding-005) et trouvez les K voisins les plus proches à l'aide de VECTOR_SEARCH pour servir de couche de génération de candidats sous-linéaire.
  • Scoring des paires de candidats et fusion hybride des caractéristiques des arêtes : exploitez les paires de candidats de la recherche vectorielle pour éliminer la complexité de la jointure croisée O(N²), calculer les scores de similarité pondérés multifonctionnels (SSN, distance d'édition de Levenshtein, date de naissance, adresse Jaccard) et fusionner les arêtes dans un tableau de candidats unifié.
  • Construction de graphiques de propriétés et traversées de chemins ISO GQL : construisez un PROPERTY GRAPH BigQuery, exécutez des requêtes de chemin {1, 2} ISO GQL (GRAPH_TABLE) pour résoudre les clusters de clients connectés, calculez les métriques d'évaluation individuelles et effectuez un clustering souple des foyers à l'aide de la pondération de graphe Adamic-Adar.
  • Résolution incrémentielle et stabilité persistante : traitez les entrées par lot quotidiennes avec une mise en correspondance incrémentielle des deltas.
  • Consolidation du clustering et stabilité des clusters (chevauchement 1-ε) : assurez une stabilité persistante des clusters lors des exécutions de pipeline à l'aide d'une garantie de seuil de chevauchement (1-ε).

Prérequis

  • Un navigateur Web tel que Chrome.
  • Un projet Google Cloud avec facturation activée.

Cet atelier de programmation s'adresse aux ingénieurs de données, aux développeurs de bases de données et aux spécialistes de l'IA/ML de tous niveaux, y compris aux débutants.

Durée estimée : 45 minutes
Coût estimé : moins de 2 $ (utilise les fonctions Cloud et le traitement des requêtes BigQuery avec paiement à l'utilisation).

2. Avant de commencer

Créer un projet Google Cloud

  1. Dans la console Google Cloud, sur la page de sélection du projet, sélectionnez ou créez un projet Google Cloud.
  2. Assurez-vous que la facturation est activée pour votre projet Cloud. Découvrez comment vérifier si la facturation est activée sur un projet.

Démarrer Cloud Shell

Cloud Shell est un environnement de ligne de commande exécuté dans Google Cloud et fourni avec les outils nécessaires.

  1. Cliquez sur Activer Cloud Shell en haut de la console Google Cloud.
  2. Vérifiez votre authentification :
gcloud auth list
  1. Configurez les variables d'environnement dans Cloud Shell :
export GCP_PROJECT=$(gcloud config get-value project)
export REGION="us-central1"
export DATASET_ID="identity_resolution"

Activer les API requises

Exécutez la commande suivante dans Cloud Shell à l'aide de votre compte utilisateur pour activer tous les services Google Cloud requis :

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

Pour garantir une exécution fluide des API et un accès aux identifiants par défaut de l'application (ADC), créez un compte de service d'atelier dédié et activez l'emprunt d'identité 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

Créer un ensemble de données BigQuery

Créez l'ensemble de données BigQuery pour stocker vos nœuds clients, vos arêtes, vos modèles de graphiques et vos vues d'évaluation :

bq mk --location=US --dataset ${GCP_PROJECT}:${DATASET_ID}

Un résultat semblable à celui-ci s'affiche :

Dataset 'your-project-id:identity_resolution' successfully created.

Pour garantir une capacité de calcul dédiée aux recherches d'index vectoriels, aux agrégations de graphiques et à l'exécution de fonctions à distance sans être limité par les limites de processeur à la demande ni les quotas partagés, créez une réservation Enterprise Edition avec autoscaling dans 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. Ingérer l'ensemble de données de nœuds client FEBRL3

Avant de déployer la fonction distante de validation d'adresse et d'effectuer la résolution d'identité, vous allez charger l'ensemble de données de référence pour la résolution d'entités FEBRL3 synthétique (qui contient 5 000 fiches client avec des clusters de doublons multiples pouvant comporter jusqu'à cinq doublons par client) à l'aide de la bibliothèque recordlinkage de Python,puis écrire les nœuds client bruts (customer_nodes) et les liens de correspondance de référence (ground_truth_links) dans BigQuery à l'aide de BigQuery DataFrames (bigframes).

Exécutez les commandes suivantes dans Cloud Shell pour installer les dépendances et exécuter le script d'ingestion :

# 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

Dans la console Google Cloud, accédez à BigQuery Studio, ouvrez un nouvel onglet Requête SQL (+), puis exécutez la requête ci-dessous pour inspecter la table des nœuds clients ingérés :

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;

Un résultat semblable à celui-ci s'affiche :

rec_id

given_name

nom de famille

street_number

address_1

address_2

banlieue

code postal

state

date_of_birth

soc_sec_id

dataset_source

rec-10-org

brent

wood

11

girdlestone circuit

kingston tower

clifton springs

4152

nsw

19340706

1075870

febrl3

rec-10-dup-0

brnt

woode

11

girdelstone circut

cliffton springs

4152

nsw

19340706

1075870

febrl3

rec-10-dup-1

bernt

wood

11

girdlestone cir

kingston twr

clifton spngs

4152

19340760

1075870

febrl3

rec-10-dup-2

brent

wod

15

girdlestone crt

clifton springs

4152

nsw

19340706

febrl3

rec-25-org

mccarthy

henry

13

beasley street

crystal brook farm

ingleburn

6164

vic

19770913

1347524

febrl3

Notez comment l'ensemble de données de référence introduit des données sales réalistes dans les clusters en double :

  • Variantes phonétiques et orthographiques : brent vs brnt / bernt, wood vs woode / wod et clifton vs cliffton.
  • Abréviations et fautes de frappe dans les adresses : girdlestone circuit vs girdelstone circut / girdlestone cir / girdlestone crt et numéro de maison 11 vs erreur OCR 15.
  • Transpositions de caractères et valeurs manquantes : transpositions de la date de naissance (19340706 au lieu de 19340760), états manquants ( ) et numéros de sécurité sociale manquants ( ).

Dans les étapes à venir, vous utiliserez des encodages phonétiques SOUNDEX, des UDF de normalisation d'adresses, la distance d'édition de Levenshtein et la recherche vectorielle AI.EMBED pour combler ces écarts et associer précisément les profils en double.

4. Déployer la fonction distante de validation d'adresse

La normalisation des adresses permet de standardiser les noms de rues, les limites des banlieues et les codes postaux avant d'effectuer la mise en correspondance. L'API Address Validation de Google Maps est un service qui accepte une adresse, identifie ses composants et les valide. Dans cette étape, vous allez déployer une fonction Cloud Python dans Cloud Shell qui expose une UDF de validation et de normalisation d'adresse à BigQuery.

Écrire des fichiers sources de fonctions Cloud

Exécutez la commande suivante dans Cloud Shell pour créer le répertoire source Cloud Functions et écrire main.py et 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

Déployer une fonction Cloud Functions et configurer les autorisations IAM

Exécutez ces commandes dans Cloud Shell pour déployer la fonction Cloud de 2e génération et configurer une connexion aux ressources 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

Le résultat qui s'affiche doit indiquer que le déploiement de la fonction Cloud est terminé et que les liaisons IAM ont bien été appliquées.

Enregistrer la fonction de normalisation des adresses distantes

Vous allez maintenant enregistrer le DDL (Data Definition Language) de la fonction distante BigQuery (validate_address_udf) qui connecte les lignes de la table BigQuery au point de terminaison de votre fonction Cloud déployée (${FUNCTION_URL}).

Exécutez la commande suivante dans Cloud Shell pour récupérer l'URL de votre fonction Cloud déployée et enregistrer automatiquement la fonction à distance :

# 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. Prétraiter les données de profil et les encodages phonétiques

Au cours de cette étape, vous allez exécuter une requête SQL de prétraitement BigQuery sur votre table customer_nodes ingérée.

Exécuter une requête de nettoyage des données et de caractéristiques phonétiques

Dans l'éditeur SQL BigQuery Studio, exécutez la requête ci-dessous pour créer customer_nodes_cleaned. Cette requête :

  1. Appelle validate_address_udf pour obtenir des adresses normalisées et des verdicts de validation.
  2. Génère des encodages phonétiques SOUNDEX pour given_name et surname afin de gérer les variantes orthographiques.
  3. Construit un champ profile_text structuré.
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;

Interrogez la table des nœuds nettoyés :

SELECT rec_id, given_name, given_name_soundex, surname, surname_soundex, formatted_address 
FROM `identity_resolution.customer_nodes_cleaned`
LIMIT 5;

Un résultat semblable à celui-ci s'affiche :

rec_id

given_name

given_name_soundex

nom de famille

surname_soundex

formatted_address

rec-001-A

John

J500

Smith

S530

12 high st richmond vic 3121

rec-001-B

Jon

J500

Smith

S530

12 high st richmond vic 3121

rec-002-A

Elizabeth

E421

Taylor

T460

45 park rd suite 4 south yarra vic 3141

6. Générer des embeddings de profil sémantique et utiliser la recherche vectorielle

En plus de la normalisation des adresses, des clés phonétiques Soundex et de la distance d'édition Levenshtein, BigQuery est compatible avec les fonctions d'embedding d'IA générative intégrées via AI.EMBED.

À l'aide de AI.EMBED, BigQuery génère des embeddings de texte directement en SQL à l'aide de modèles de base (comme text-embedding-005) sans nécessiter de DDL d'index vectoriel manuel :

Exécutez les requêtes ci-dessous dans l'éditeur SQL de 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. Scoring des paires de candidats et fusion hybride des caractéristiques Edge

L'évaluation de toutes les paires d'enregistrements client possibles (croissance quadratique O(N²)) devient prohibitive en termes de calcul à mesure que l'ensemble de données augmente. Dans les moteurs de base de données relationnelles comme BigQuery, la tentative d'implémentation du blocage basé sur des règles à l'aide de conditions de jointure OR complexes sur plusieurs colonnes (par exemple, une jointure sur a.soc_sec_id = b.soc_sec_id OR a.given_name_soundex = b.given_name_soundex OR ...) empêche l'optimiseur de requête d'utiliser des jointures hachées ou des jointures par fusion-tri sur une seule clé d'équi-jointure. Au lieu de cela, le moteur revient à une jointure croisée O(N²) et filtre chaque paire, ce qui échoue à grande échelle.

Dans cette étape, vous allez prendre les paires candidates générées par notre table de recherche vectorielle (vector_candidate_edges) et les joindre à customer_nodes_cleaned à l'aide de jointures équi-jointes indexées rapides (ON c.source_id = a.rec_id et ON c.target_id = b.rec_id). Vous allez ensuite calculer un score de correspondance pondéré combinant :

  • Score de correspondance du numéro de sécurité sociale (pondération : 0.30)
  • Similarité de modification du nom de famille à l'aide de la distance de Levenshtein EDIT_DISTANCE (pondération : 0.20)
  • Similarité de modification du prénom (pondération : 0.20)
  • Score de correspondance de la date de naissance (pondération : 0.15)
  • Similarité Jaccard des jetons d'adresse (pondération : 0.15) sur SPLIT(LOWER(formatted_address), ' ')

Calculer les arêtes candidates et les scores de similarité pondérés

Exécutez la requête suivante dans l'éditeur SQL de BigQuery Studio pour remplir 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;

Inspectez les correspondances d'arêtes candidates :

SELECT source_id, target_id, match_score 
FROM `identity_resolution.matched_edges`
ORDER BY match_score DESC;

Un résultat semblable à celui-ci s'affiche :

source_id

target_id

match_score

rec-001-A

rec-001-B

0.9400

rec-002-A

rec-002-B

0.9100

rec-003-A

rec-003-B

0.7300

Fusionner les nœuds de recherche vectorielle et basée sur des règles dans un tableau unifié

Combinez les arêtes candidates issues de la correspondance approximative basée sur des règles et de la recherche vectorielle sémantique dans une seule table final_matched_edges dédupliquée :

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. Construction de graphiques de propriétés et traversées de chemins ISO GQL

BigQuery est compatible avec le langage GQL (Graph Query Language) ISO de façon native via les graphes de propriétés. Un graphe de propriétés crée une vue logique de graphe sur les tables relationnelles BigQuery sans duplication des données.

Dans cette étape, vous allez créer un graphique de propriétés customer_identity_graph à l'aide de votre table d'arêtes de candidats unifiée (final_matched_edges) et interroger les connexions du graphique sur les profils client à l'aide du parcours k-hop {1, 2}.

Connectivité transitive et traversées à K sauts

Créer un LDD de graphe de propriété BigQuery

Exécutez l'instruction DDL suivante dans l'éditeur SQL de 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
);

Visualiser les clusters de graphiques K-Hop

Exécutez la requête ci-dessous pour visualiser les clusters de clients correspondants sur un ou deux niveaux de relations :

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;

Visualisation des clusters de graphiques K-Hop

Résoudre les clusters de clients canoniques

Exécutez la requête suivante pour résoudre les clusters d'entités en 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;

Interrogez la table des clusters résolus :

SELECT canonical_customer_id, record_count, customer_records 
FROM `identity_resolution.resolved_customers`
ORDER BY record_count DESC;

Un résultat semblable à celui-ci s'affiche :

canonical_customer_id

record_count

customer_records

rec-001-A

2

['rec-001-A', 'rec-001-B']

rec-002-A

2

['rec-002-A', 'rec-002-B']

rec-003-A

2

['rec-003-A', 'rec-003-B']

Créer une vue des métriques d'évaluation

Pour calculer la précision, le rappel et le score F1 par rapport à la table ground_truth_links, exécutez la commande suivante :

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;

Interrogez la vue des métriques d'évaluation :

SELECT * FROM `identity_resolution.evaluation_metrics`;

Un résultat semblable à celui-ci s'affiche :

total_ground_truth

total_predictions

true_positives (vrais positifs)

false_positives (faux positifs)

false_negatives (faux négatifs)

precision

recall

f1_score

6538

5620

5608

13

930

0.9977

0.8579

0.9225

Résoudre les clusters de foyers à l'aide de la pondération du graphique Adamic-Adar

Alors que la résolution d'identité individuelle permet de résoudre les enregistrements appartenant à la même personne, les architectures Customer 360 d'entreprise nécessitent souvent un regroupement d'entités de foyer de niveau supérieur pour les personnes résidant à la même adresse.

Étant donné que les benchmarks synthétiques (comme FEBRL3) évaluent la vérité terrain au niveau individuel, la résolution des foyers est effectuée en tant qu'étape en aval. En l'absence d'historique des déménagements horodaté, les personnes associées à plusieurs adresses peuvent entraîner une fusion excessive ou une fragmentation des clusters. Pour résoudre ce problème, nous utilisons la pondération de graphe Adamic-Adar afin de créer des appartenances souples aux foyers.

Exécutez la requête ci-dessous dans l'éditeur SQL BigQuery Studio pour remplir household_clusters à l'aide de la pondération de graphe 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;

Interrogez la table des clusters de foyers résolus :

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;

Visualiser la hiérarchie des identités de bout en bout avec GQL

Pour suivre visuellement l'intégralité de la hiérarchie d'identité à trois niveaux (en reliant les clients bruts non regroupés aux entités client résolues, puis aux entités foyer résolues), exécutez la requête DDL et ISO GQL suivante dans 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'exécution de cette requête GQL dans BigQuery Studio affiche un canevas de visualisation graphique interactif à trois niveaux. Il présente les enregistrements bruts des profils client (RawCustomer) résolus en entités canoniques individuelles (ResolvedCustomer), qui sont liées à des entités de foyers partagés à plusieurs résidents (ResolvedHousehold).

Visualisation des clusters de graphiques des foyers à plusieurs résidents

9. Résolution incrémentielle et stabilité persistante

Dans les applications d'entreprise réelles, de nouveaux enregistrements de clients arrivent en continu via des lots quotidiens ou en temps réel. Plutôt que de réexécuter la résolution complète du graphique sur l'ensemble de l'ensemble de données historiques, un moteur de correspondance delta incrémentiel compare les nouveaux enregistrements entrants aux clusters de référence résolus existants (resolved_customers).

Pour ce faire efficacement, le moteur utilise Vector Search

VECTOR_SEARCH

) comme une forme de clustering dynamique. En traitant chaque enregistrement entrant comme un point de requête, VECTOR_SEARCH récupère l'ensemble des K voisins les plus proches à partir de l'index d'embedding de référence historique. Si un enregistrement entrant correspond à un profil client existant au-dessus du seuil de similarité, il est fusionné de manière dynamique dans ce cluster et hérite de la référence canonical_customer_id (MATCHED_TO_EXISTING_CLUSTER). Si aucun voisin le plus proche de référence n'est trouvé au-dessus du seuil, un nouvel UUID d'entité est créé (NEW_CUSTOMER_ENTITY).

Ingérer des exemples d'enregistrements d'ingestion incrémentielle par lot

Collez et exécutez la DDL suivante dans l'éditeur SQL de BigQuery Studio pour créer 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
  )
]);

Exécuter une requête de correspondance incrémentielle Delta

Exécutez la requête suivante dans l'éditeur SQL BigQuery pour effectuer une mise en correspondance delta par rapport à votre ensemble de données de référence résolu :

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;

Interrogez les résultats de résolution incrémentaux :

SELECT record_id, persistent_canonical_customer_id, assigned_household_id, assignment_type, match_score, match_strategy 
FROM `identity_resolution.incremental_resolved_customers`;

Un résultat semblable à celui-ci s'affiche :

record_id

persistent_canonical_customer_id

assigned_household_id

assignment_type

match_score

match_strategy

rec-9999-new-1

rec-529-dup-0

hh-8745613133408211212

MATCHED_TO_EXISTING_CLUSTER

1.0000

BOTH

rec-9999-new-2

rec-875-dup-0

hh-4802217335228845918

MATCHED_TO_EXISTING_CLUSTER

1.0000

BOTH

rec-9999-new-3

rec-359-dup-0

hh-8891673998566910207

MATCHED_TO_EXISTING_CLUSTER

0.8739

VECTOR_SEARCH

rec-9999-new-4

197e02f9-7175-4484-9c32-b7edbc731c9e

NULL

NEW_CUSTOMER_ENTITY

NULL

NONE

10. Consolidation du clustering et stabilité des clusters (chevauchement de 1-ε)

Dans les systèmes d'entreprise de production, les équipes regroupent généralement l'ensemble du graphique de manière récurrente (par exemple, toutes les semaines ou tous les mois) pour intégrer de nouvelles arêtes et sources de données. À mesure que de nouvelles relations se forment, le re-clustering complet du graphique peut entraîner un déplacement ou une inversion arbitraires des identifiants de cluster lors des exécutions du pipeline.

Pour conserver des ID client persistants pour les systèmes CRM, CDP et de facturation en aval, la stabilité des clusters évalue le chevauchement des nœuds entre les clusters de l'exécution actuelle (t) et les clusters de l'exécution précédente (t-1) à l'aide d'un seuil de chevauchement (1 - ε), où ε = 0, 30 (ce qui nécessite un chevauchement des nœuds d'au moins 70 %).

Si un cluster nouvellement calculé dans Run (t) partage au moins 70% de ses enregistrements membres avec un cluster de Run (t-1), il hérite de l'ID client persistant historique (STABLE_EVOLUTION). Les tout nouveaux clusters reçoivent des UUID nouvellement générés (NEW_CLUSTER_CREATED).

Exécuter une requête sur le chevauchement et la stabilité des clusters

Exécutez la requête suivante dans l'éditeur SQL de BigQuery Studio pour remplir 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;

Afficher l'aperçu filtré des enregistrements incrémentaux

Exécutez cette requête pour vérifier l'état de stabilité du cluster pour vos enregistrements d'ingestion quotidienne :

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;

Un résultat semblable à celui-ci s'affiche :

persistent_canonical_customer_id

record_id

assigned_household_id

cluster_status

rec-529-dup-0

rec-9999-new-1

hh-8745613133408211212

STABLE_EVOLUTION

rec-875-dup-0

rec-9999-new-2

hh-4802217335228845918

STABLE_EVOLUTION

rec-359-dup-0

rec-9999-new-3

hh-8891673998566910207

STABLE_EVOLUTION

589f8102-1204-4530-8910-bc10294810a4

rec-9999-new-4

NULL

NEW_CLUSTER_CREATED

11. Effectuer un nettoyage

Pour éviter que des frais ne soient facturés en permanence sur votre compte Google Cloud, nettoyez les ressources déployées et l'ensemble de données BigQuery.

Dans Cloud Shell, exécutez la commande suivante :

# 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

Si vous avez créé un projet Google Cloud dédié pour cet atelier, vous pouvez le supprimer :

gcloud projects delete ${GCP_PROJECT}

12. Félicitations

Félicitations ! Vous avez réussi à créer un moteur de résolution d'identité client de bout en bout dans Google Cloud BigQuery à l'aide de BigQuery Property Graph, de requêtes ISO GQL, de la mise en correspondance hybride de similarité, de la mise en correspondance incrémentielle des deltas et des garanties de stabilité des clusters persistants.

Connaissances acquises

  • Découvrez comment déployer une fonction Cloud Functions de 2e génération et l'exposer en tant que fonction distante BigQuery.
  • Comment prétraiter les données démographiques des clients à l'aide des encodages phonétiques SOUNDEX et de la validation des adresses.
  • Comment exécuter le blocage des candidats et calculer les scores de similarité hybrides à l'aide de la distance de Levenshtein (EDIT_DISTANCE) et de la similarité Jaccard des jetons.
  • Comment construire un graphique de propriétés BigQuery (CREATE PROPERTY GRAPH) sur des tables de nœuds et d'arêtes.
  • Comment interroger des chemins de graphe à l'aide de GQL (GRAPH_TABLE) ISO avec des quantificateurs k-hop {1, 2}.
  • Découvrez comment résoudre les problèmes liés aux clusters de clients canoniques et évaluer les performances du modèle par rapport aux métriques de vérité terrain.
  • Découvrez comment exécuter la mise en correspondance incrémentielle des deltas pour les ingérences de lots quotidiens sans retraitement complet de l'ensemble de données.
  • Découvrez comment appliquer une garantie de seuil de chevauchement (1 – ε) pour maintenir la stabilité persistante des clusters lors des exécutions de pipeline.

Étapes suivantes

Documents de référence