Resolución de la identidad del cliente con BigQuery Graph

1. Introducción

En este codelab, compilarás un motor modular de resolución de identidad del cliente (correlación de entidades) de extremo a extremo directamente en Google Cloud BigQuery. Combinarás Google Cloud Shell para la implementación de la infraestructura con el editor de SQL de BigQuery Studio para la limpieza de datos, la puntuación de candidatos, la construcción de gráficos de propiedades y los recorridos de rutas de Graph Query Language (GQL) de ISO.

La resolución de identidades es una capacidad fundamental para la vista del cliente de 360 grados empresarial, la detección de fraudes y la consolidación de datos en varios sistemas. Dado que existen muchos enfoques válidos para la resolución de identidades según la madurez de los datos y las necesidades comerciales, todos los pasos de este codelab son modulares y opcionales. La canalización está diseñada para mostrar una variedad de técnicas industriales comunes y aptas para la producción, como la normalización de direcciones de UDF remotas, el bloqueo fonético de Soundex, la búsqueda de vectores semánticos (AI.EMBED), la puntuación de atributos híbridos y la agrupación en clústeres de grafos de propiedades de GQL, para que puedas adoptar de forma selectiva los patrones que se ajusten a tu arquitectura.

Los métodos de correlación y los umbrales de puntuación se deben ajustar según el interés de tu organización en la correlación determinística frente a la probabilística, que se determina según el caso de uso objetivo. Por ejemplo, las operaciones financieras, de facturación o de cumplimiento estrictas suelen favorecer las reglas determinísticas de alta precisión (como las coincidencias exactas del número de seguridad social o el ID fiscal) para evitar la vinculación falsa, mientras que los motores de personalización de marketing, análisis y recomendación a menudo se inclinan por la coincidencia difusa probabilística y la similitud de vectores semánticos para maximizar la recuperación y descubrir conexiones sutiles.

Arquitectura del motor de resolución de identidades del cliente de BigQuery

Actividades

  • Ingest FEBRL3 Benchmark Dataset: Carga registros de clientes sintéticos y pares de coincidencias de referencia en BigQuery.
  • Implementa la UDF remota de Address Validation: Implementa una Cloud Function de Python y registra una función remota de BigQuery para normalizar las direcciones.
  • Preprocess Profile Data & Phonetic Encodings: Ejecuta la limpieza de datos con SQL, invoca la UDF de dirección y calcula las claves fonéticas SOUNDEX y las distancias de edición de Levenshtein:
    • Codificación fonética de Soundex: Es un algoritmo fonético para indexar nombres por sonido, como se pronuncian en inglés. Convierte los nombres en un código de 4 caracteres (una letra inicial seguida de tres dígitos) que representa grupos de sonidos consonánticos (p.ej., tanto "John" como "Jon" se asignan a J500, mientras que "Smith" y "Smyth" se asignan a S530), lo que proporciona indicadores de coincidencia fonética para la puntuación de funciones y el bloqueo incremental delta en tiempo real.
    • Distancia de Levenshtein (EDIT_DISTANCE): Es una métrica de cadena que mide la cantidad mínima de ediciones de un solo carácter (inserciones, eliminaciones o sustituciones) necesarias para cambiar una cadena por otra, lo que permite una coincidencia difusa precisa de nombres y direcciones.
  • Genera embeddings de perfil semántico y realiza búsquedas de vectores: Genera embeddings de texto directamente en SQL con AI.EMBED (text-embedding-005) y encuentra los K vecinos más cercanos con VECTOR_SEARCH para que sirvan como una capa de generación de candidatos sublineal.
  • Puntuación de pares candidatos y fusión de funciones de borde híbridas: Aprovecha los pares candidatos de la búsqueda de vectores para eliminar la complejidad de la unión cruzada O(N²), calcula las puntuaciones de similitud ponderadas de múltiples funciones (SSN, distancia de edición de Levenshtein, DOB, Jaccard de dirección) y fusiona los bordes en una tabla de candidatos unificada.
  • Construcción de gráficos de propiedades y recorridos de rutas de ISO GQL: Construye un PROPERTY GRAPH de BigQuery, ejecuta consultas de rutas de ISO GQL {1, 2} (GRAPH_TABLE) para resolver clústeres de clientes conectados, calcula métricas de evaluación individuales y realiza un agrupamiento en clústeres flexible de grupos familiares con la ponderación de gráficos de Adamic-Adar.
  • Resolución incremental y estabilidad persistente: Procesa las incorporaciones de lotes diarias con la correlación de delta incremental.
  • Consolidación y estabilidad de clústeres (superposición de 1-ε): Aplica la estabilidad persistente de los clústeres en las ejecuciones de la canalización con una garantía de umbral de superposición de (1-ε).

Requisitos

  • Un navegador web, como Chrome
  • Un proyecto de Google Cloud con facturación habilitada.

Este codelab está diseñado para ingenieros de datos, desarrolladores de bases de datos y profesionales de la IA/AA de todos los niveles, incluidos los principiantes.

Duración estimada: 45 minutos
Costo estimado: Menos de USD 2.00 (usa Cloud Functions de pago por uso y procesamiento de consultas de BigQuery).

2. Antes de comenzar

Crea un proyecto de Google Cloud

  1. En la consola de Google Cloud, en la página del selector de proyectos, selecciona o crea un proyecto de Google Cloud.
  2. Asegúrate de que la facturación esté habilitada para tu proyecto de Cloud. Obtén información para verificar si la facturación está habilitada en un proyecto.

Inicie Cloud Shell

Cloud Shell es un entorno de línea de comandos que se ejecuta en Google Cloud y que viene precargado con las herramientas necesarias.

  1. Haz clic en Activar Cloud Shell en la parte superior de la consola de Google Cloud.
  2. Verifica tu autenticación:
gcloud auth list
  1. Configura las variables de entorno en Cloud Shell:
export GCP_PROJECT=$(gcloud config get-value project)
export REGION="us-central1"
export DATASET_ID="identity_resolution"

Habilitar las APIs obligatorias

Ejecuta el siguiente comando en Cloud Shell con tu cuenta de usuario para habilitar todos los servicios de Google Cloud necesarios:

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

Para garantizar la ejecución sin problemas de la API y el acceso a las credenciales predeterminadas de la aplicación (ADC), crea una cuenta de servicio de lab dedicada y habilita la suplantación de 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 un conjunto de datos de BigQuery

Crea el conjunto de datos de BigQuery para almacenar tus nodos, aristas, modelos de grafos y vistas de evaluación de clientes:

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

Deberías ver un resultado similar a este:

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

Para garantizar la capacidad de procesamiento dedicada para las búsquedas de índices vectoriales, las agregaciones de gráficos y la ejecución de funciones remotas sin estar limitado por los límites de CPU a pedido ni las cuotas compartidas, crea una reserva de Enterprise Edition con ajuste de escala automático en 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. Transfiere el conjunto de datos de nodos de clientes de FEBRL3

Antes de implementar la función remota de validación de direcciones y realizar la resolución de identidades, cargarás el conjunto de datos de referencia sintético de resolución de entidades FEBRL3 (que contiene 5,000 registros de clientes con clústeres de varios duplicados de hasta 5 duplicados por cliente) con la biblioteca recordlinkage de Python y escribirás los nodos de clientes sin procesar (customer_nodes) y los vínculos de coincidencia de verdad fundamental (ground_truth_links) en BigQuery con BigQuery DataFrames (bigframes).

Ejecuta los siguientes comandos en Cloud Shell para instalar las dependencias y ejecutar la secuencia de comandos de transferencia:

# 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

En la consola de Google Cloud, navega a BigQuery Studio, abre una nueva pestaña de consulta en SQL (+) y ejecuta la siguiente consulta para inspeccionar la tabla de nodos de clientes transferidos:

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;

Deberías ver un resultado similar a este:

rec_id

given_name

apellido

street_number

address_1

address_2

suburbio

código 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

Observa cómo el conjunto de datos de comparativas introduce datos sucios realistas en clústeres duplicados:

  • Variaciones fonéticas y ortográficas: brent en comparación con brnt / bernt, wood en comparación con woode / wod y clifton en comparación con cliffton.
  • Abreviaturas y errores de escritura en las direcciones: girdlestone circuit vs. girdelstone circut / girdlestone cir / girdlestone crt y número de casa 11 vs. error de OCR 15.
  • Transposiciones de caracteres y valores faltantes: Transposiciones de fechas de nacimiento (19340706 vs. 19340760), estados faltantes ( ) y números de identificación personal faltantes ( ).

En los próximos pasos, usarás codificaciones fonéticas SOUNDEX, UDF de normalización de direcciones, distancia de edición de Levenshtein y búsqueda de vectores AI.EMBED para subsanar estas discrepancias y vincular con precisión los perfiles duplicados.

4. Implementa la UDF de la función remota de Address Validation

La normalización de direcciones estandariza los nombres de las calles, los límites de los suburbios y los códigos postales antes de realizar la correlación. La API de Address Validation de Google Maps es un servicio que acepta una dirección, identifica sus componentes y los valida. En este paso, implementarás una Cloud Function de Python en Cloud Shell que exponga una UDF de validación y normalización de direcciones a BigQuery.

Escribe archivos fuente de Cloud Functions

Ejecuta el siguiente comando en Cloud Shell para crear el directorio del código fuente de Cloud Functions y escribir main.py y 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

Implementa Cloud Function y configura los permisos de IAM

Ejecuta estos comandos en Cloud Shell para implementar la Cloud Function de 2ª gen. y configurar una conexión a recursos de Cloud de 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

Deberías ver un resultado que indica que se completó la implementación de Cloud Function y que las vinculaciones de IAM se aplicaron correctamente.

Registra la función de normalización de direcciones remotas

Ahora registrarás el DDL de la función remota de BigQuery (validate_address_udf) que conecta las filas de la tabla de BigQuery con el extremo de Cloud Functions implementado (${FUNCTION_URL}).

Ejecuta el siguiente comando en Cloud Shell para recuperar la URL de la Cloud Function implementada y registrar la función remota automáticamente:

# 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. Procesamiento previo de datos de perfil y codificaciones fonéticas

En este paso, ejecutarás una consulta de preprocesamiento en SQL de BigQuery en tu tabla customer_nodes transferida.

Ejecuta la limpieza de datos y la consulta de funciones fonéticas

En el editor de SQL de BigQuery Studio, ejecuta la siguiente consulta para crear customer_nodes_cleaned. La consulta realiza las siguientes acciones:

  1. Llama a validate_address_udf para obtener direcciones normalizadas y veredictos de validación.
  2. Genera codificaciones fonéticas SOUNDEX para given_name y surname para controlar las variaciones ortográficas.
  3. Construye un campo profile_text estructurado.
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;

Consulta la tabla de nodos limpios:

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

Deberías ver un resultado similar a este:

rec_id

given_name

given_name_soundex

apellido

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. Genera embeddings de perfiles semánticos y búsqueda de vectores

Además de la normalización de direcciones, las claves fonéticas de Soundex y la distancia de edición de Levenshtein, BigQuery admite funciones de embedding de IA generativa integradas a través de AI.EMBED.

Con AI.EMBED, BigQuery genera incorporaciones de texto directamente en SQL con modelos básicos (como text-embedding-005) sin necesidad de DDL de índice de vectores manual:

Ejecuta las siguientes consultas en el editor de 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. Puntuación de pares de candidatos y fusión de funciones híbridas de Edge

Evaluar todos los pares posibles de registros de clientes (crecimiento cuadrático O(N²)) se vuelve prohibitivo desde el punto de vista computacional a medida que aumenta la escala del conjunto de datos. En los motores de bases de datos relacionales, como BigQuery, intentar implementar el bloqueo basado en reglas con condiciones de unión OR complejas en varias columnas (como la unión en a.soc_sec_id = b.soc_sec_id OR a.given_name_soundex = b.given_name_soundex OR ...) impide que el optimizador de consultas use uniones hash escalables o uniones de ordenamiento y combinación en una sola clave de unión por igualdad. En cambio, el motor recurre a una unión cruzada O(N²) y filtra cada par, lo que falla a gran escala.

En este paso, tomarás los pares candidatos generados por nuestra tabla de búsqueda de vectores (vector_candidate_edges) y los unirás con customer_nodes_cleaned a través de uniones de igualdad indexadas rápidas (ON c.source_id = a.rec_id y ON c.target_id = b.rec_id). Luego, calcularás una puntuación de coincidencia ponderada que combine lo siguiente:

  • Puntuación de coincidencia del NSS (peso: 0.30)
  • Surname Edit Similarity con la distancia de Levenshtein EDIT_DISTANCE (peso: 0.20)
  • Given Name Edit Similarity (peso: 0.20)
  • Puntuación de coincidencia de la fecha de nacimiento (peso: 0.15)
  • Similitud de Jaccard de tokens de dirección (peso: 0.15) sobre SPLIT(LOWER(formatted_address), ' ')

Calcula los bordes candidatos y las puntuaciones de similitud ponderadas

Ejecuta la siguiente consulta en el editor de SQL de BigQuery Studio para completar 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;

Inspecciona las coincidencias de borde candidatas:

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

Deberías ver un resultado similar a este:

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

Cómo fusionar bordes de búsqueda basados en reglas y vectores en una tabla unificada

Combina las aristas candidatas de la coincidencia aproximada basada en reglas y la búsqueda de vectores semánticos en una sola tabla final_matched_edges sin duplicados:

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. Construcción de gráficos de propiedades y recorridos de rutas de GQL según la norma ISO

BigQuery admite ISO GQL (Graph Query Language) de forma nativa a través de Property Graphs. Un grafo de propiedades crea una vista lógica del grafo sobre las tablas relacionales de BigQuery sin duplicación de datos.

En este paso, compilarás un gráfico de propiedades customer_identity_graph con tu tabla de aristas candidatas unificada (final_matched_edges) y consultarás las conexiones del gráfico de consulta en los perfiles de clientes con el recorrido de ruta de k-saltos {1, 2}.

Conectividad transitiva y recorridos de K-saltos

Crea un DDL de BigQuery Property Graph

Ejecuta la siguiente instrucción DDL en el editor de 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
);

Visualiza los clústeres del gráfico de K-Hop

Ejecuta la siguiente consulta para visualizar los clústeres de clientes coincidentes en 1 o 2 saltos de relación:

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;

Visualización de clústeres de gráficos de K-Hop

Resuelve los clústeres de clientes canónicos

Ejecuta la siguiente consulta para resolver clústeres de entidades 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;

Consulta la tabla de clústeres resueltos:

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

Deberías ver un resultado similar a este:

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']

Crea la vista de métricas de evaluación

Para calcular la precisión, la recuperación y la puntuación F1 en comparación con la tabla ground_truth_links, ejecuta lo siguiente:

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;

Consulta la vista de las métricas de evaluación:

SELECT * FROM `identity_resolution.evaluation_metrics`;

Deberías ver un resultado similar a este:

total_ground_truth

total_predictions

true_positives

false_positives

false_negatives

precision

recall

f1_score

6538

5620

5608

13

930

0.9977

0.8579

0.9225

Resolución de clústeres familiares a través de la ponderación del gráfico de Adamic-Adar

Si bien la resolución de identidad individual resuelve los registros que pertenecen a la misma persona, las arquitecturas de Customer 360 empresariales a menudo requieren una agrupación de entidades familiares de nivel superior que agrupe a las personas que residen en la misma dirección.

Dado que las comparativas sintéticas (como FEBRL3) evalúan la verdad fundamental a nivel individual, la resolución de la familia se realiza como un paso posterior. En ausencia de un historial de reubicación con marcas de tiempo, las personas vinculadas a varias direcciones podrían provocar una fusión excesiva o una fragmentación del clúster. Para resolver este problema, usamos la ponderación del gráfico de Adamic-Adar para construir membresías de grupos familiares flexibles.

Ejecuta la siguiente consulta en el editor de SQL de BigQuery Studio para completar household_clusters con la ponderación del gráfico de 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;

Consulta la tabla de clústeres de familias resueltos:

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;

Visualiza la jerarquía de identidad de extremo a extremo a través de GQL

Para hacer un seguimiento visual de la jerarquía de identidad completa de 3 niveles, que conecta los Clientes sin agrupar con las Entidades de clientes resueltas y, luego, con las Entidades de familias resueltas, ejecuta la siguiente consulta en DDL y en ISO GQL en 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;

Si ejecutas esta consulta de GQL en BigQuery Studio, se renderizará un lienzo de visualización de gráficos interactivos de 3 niveles que muestra registros sin procesar de perfiles de clientes (RawCustomer) resueltos en entidades canónicas individuales (ResolvedCustomer), que están vinculadas a entidades de grupo familiar compartido de varios residentes (ResolvedHousehold).

Visualización de clústeres del gráfico de grupos familiares con varios residentes

9. Resolución incremental y estabilidad persistente

En las aplicaciones empresariales del mundo real, los registros de clientes nuevos llegan de forma continua a través de incorporaciones por lotes diarias o en tiempo real. En lugar de volver a ejecutar la resolución completa del gráfico en todo el conjunto de datos históricos, un Incremental Delta Matching Engine compara los nuevos registros entrantes con los clústeres de referencia resueltos existentes (resolved_customers).

Para lograr esto de manera eficiente, el motor usa la búsqueda de vectores (

VECTOR_SEARCH

) como una forma de agrupamiento dinámico. Al tratar cada registro entrante como un punto de consulta, VECTOR_SEARCH recupera el conjunto de los K vecinos más cercanos del índice de incorporación del modelo de referencia histórico. Si un registro entrante coincide con un perfil de cliente existente por encima del umbral de similitud, se combina de forma dinámica en ese clúster y hereda el canonical_customer_id (MATCHED_TO_EXISTING_CLUSTER) de referencia. Si no se encuentra ningún vecino más cercano de referencia por encima del umbral, se crea un nuevo UUID de entidad (NEW_CUSTOMER_ENTITY).

Transfiere registros de muestra de la transferencia incremental por lotes

Pega y ejecuta el siguiente DDL en el editor de SQL de BigQuery Studio para crear 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
  )
]);

Ejecuta la consulta de coincidencia delta incremental

Ejecuta la siguiente consulta en el editor de SQL de BigQuery para realizar la correlación de delta con tu conjunto de datos de referencia resuelto:

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;

Consulta los resultados de la resolución incremental:

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

Deberías ver un resultado similar a este:

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. Consolidación del agrupamiento y estabilidad del clúster (superposición de 1-ε)

En los sistemas empresariales de producción, los equipos suelen volver a agrupar todo el gráfico de forma recurrente (p.ej., semanal o mensual) para incorporar nuevos bordes y fuentes de datos. A medida que se forman nuevas relaciones, el re-agrupamiento completo del grafo puede hacer que los identificadores de clústeres cambien o se inviertan de forma arbitraria en las ejecuciones de la canalización.

Para mantener los IDs de cliente persistentes en los sistemas de CRM, CDP y facturación posteriores, Cluster Stability evalúa la superposición de nodos entre los clústeres de la ejecución actual (t) y los clústeres de la ejecución anterior (t-1) con un umbral de superposición de (1 - ε) (en el que ε = 0.30, lo que requiere una superposición de nodos mínima del 70%).

Si un clúster calculado recientemente en Ejecutar (t) comparte al menos el 70% de sus registros de miembros con un clúster de Ejecutar (t-1), hereda el ID de cliente persistente histórico (STABLE_EVOLUTION). Los clústeres nuevos reciben UUIDs generados recientemente (NEW_CLUSTER_CREATED).

Ejecuta la consulta de estabilidad y superposición de clústeres

Ejecuta la siguiente consulta en el editor de SQL de BigQuery Studio para completar 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;

Cómo ver la vista previa filtrada de los registros incrementales

Ejecuta esta consulta para verificar el estado de estabilidad del clúster en tus registros de admisión diarios:

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;

Deberías ver un resultado similar a este:

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. Limpia

Para evitar que se apliquen cargos continuos a tu cuenta de Google Cloud, libera espacio en los recursos implementados y el conjunto de datos de BigQuery.

En Cloud Shell, ejecuta el siguiente comando:

# 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 creaste un proyecto de Google Cloud dedicado para este lab, puedes borrarlo:

gcloud projects delete ${GCP_PROJECT}

12. ¡Felicitaciones!

¡Felicitaciones! Compilaste con éxito un motor de resolución de identidad del cliente integral dentro de Google Cloud BigQuery con el gráfico de propiedades de BigQuery, las consultas de ISO GQL, la correlación de similitud híbrida, la correlación delta incremental y las garantías de estabilidad de clúster persistente.

Qué aprendiste

  • Cómo implementar una función de Cloud Functions de 2ª gen. y exponerla como una función remota de BigQuery
  • Cómo preprocesar los datos demográficos de los clientes con codificaciones fonéticas SOUNDEX y validación de direcciones
  • Cómo ejecutar el bloqueo de candidatos y calcular las puntuaciones de similitud híbrida con la distancia de Levenshtein (EDIT_DISTANCE) y la similitud de Jaccard de tokens
  • Cómo construir un grafo de propiedad de BigQuery (CREATE PROPERTY GRAPH) sobre tablas de nodos y aristas
  • Cómo consultar rutas de gráficos con GQL (GRAPH_TABLE) de ISO con cuantificadores de k-saltos {1, 2}.
  • Cómo resolver clústeres de clientes canónicos y evaluar el rendimiento del modelo en función de las métricas de verdad fundamental
  • Cómo ejecutar la correlación delta incremental para las incorporaciones de lotes diarias sin volver a procesar el conjunto de datos completo
  • Cómo aplicar una garantía de umbral de superposición (1 - ε) para mantener la estabilidad persistente del clúster en las ejecuciones de la canalización

Próximos pasos

Documentos de referencia