Resolução de identidade do cliente com o BigQuery Graph

1. Introdução

Neste codelab, você vai criar um mecanismo modular e completo de resolução de identidade do cliente (correspondência de entidades) diretamente no Google Cloud BigQuery. Você vai combinar o Google Cloud Shell para implantação de infraestrutura com o editor de SQL do BigQuery Studio para limpeza de dados, pontuação de candidatos, construção de gráficos de propriedades e travessias de caminhos GQL (Graph Query Language) ISO.

A resolução de identidade é um recurso fundamental para a visão 360 do cliente empresarial, a detecção de fraudes e a consolidação de dados de vários sistemas. Como há muitas abordagens válidas para a resolução de identidade, dependendo da maturidade dos dados e das necessidades comerciais, todas as etapas deste codelab são modulares e opcionais. O pipeline foi projetado para mostrar uma variedade de técnicas comuns do setor de nível de produção, incluindo normalização de endereço de UDF remota, bloqueio fonético Soundex, pesquisa de vetor semântico (AI.EMBED), pontuação de recursos híbridos e clustering de grafo de propriedades GQL. Assim, você pode adotar seletivamente os padrões que se encaixam na sua arquitetura.

Os métodos de correspondência e os limites de pontuação precisam ser ajustados com base na preferência da sua organização por correspondência determinística x probabilística, que é determinada pelo caso de uso de destino. Por exemplo, a conformidade estrita, o faturamento ou as operações financeiras geralmente favorecem regras determinísticas de alta precisão (como correspondências exatas de CPF ou CNPJ) para evitar vinculações falsas. Já a personalização de marketing, a análise e os mecanismos de recomendação costumam usar a correspondência difusa probabilística e a similaridade de vetores semânticos para maximizar o recall e descobrir conexões sutis.

Arquitetura do mecanismo de resolução de identidade do cliente do BigQuery

Atividades deste laboratório

  • Ingerir o conjunto de dados de comparativo de mercado FEBRL3: carregue registros de clientes sintéticos e pares de correspondência de verdade básica no BigQuery.
  • Implantar UDF remota de validação de endereço: implante uma função do Cloud em Python e registre uma função remota do BigQuery para normalizar endereços.
  • Pré-processar dados de perfil e codificações fonéticas: execute a limpeza de dados SQL, invoque a UDF de endereço e calcule as chaves fonéticas SOUNDEX e as distâncias de edição de Levenshtein:
    • Codificação fonética Soundex: um algoritmo fonético para indexar nomes por som, conforme pronunciados em inglês. Ele converte nomes em um código de quatro caracteres (uma letra inicial seguida por três dígitos) que representa grupos de sons consonantais.Por exemplo, "John" e "Jon" são mapeados para J500, enquanto "Smith" e "Smyth" são mapeados para S530. Isso fornece indicadores de correspondência fonética para pontuação de recursos e bloqueio incremental delta em tempo real.
    • Distância de Levenshtein (EDIT_DISTANCE): uma métrica de string que mede o número mínimo de edições de um único caractere (inserções, exclusões ou substituições) necessárias para mudar uma string em outra, permitindo a correspondência precisa de nomes e endereços aproximados.
  • Gerar embeddings de perfil semântico e pesquisa vetorial: gere embeddings de texto diretamente em SQL usando AI.EMBED (text-embedding-005) e encontre os K vizinhos mais próximos usando VECTOR_SEARCH para servir como uma camada sublinear de geração de candidatos.
  • Pontuação de pares candidatos e fusão de recursos de borda híbrida: aproveite os pares candidatos da pesquisa de vetor para eliminar a complexidade de junção cruzada O(N²), calcule pontuações de similaridade ponderada de vários recursos (SSN, distância de edição de Levenshtein, data de nascimento, Jaccard de endereço) e fusione bordas em uma tabela de candidatos unificada.
  • Construção de gráficos de propriedades e travessias de caminhos ISO GQL: crie um PROPERTY GRAPH do BigQuery, execute consultas de caminhos {1, 2} ISO GQL (GRAPH_TABLE) para resolver clusters de clientes conectados, calcule métricas de avaliação individuais e faça um clustering flexível de domicílios usando a ponderação de gráficos de Adamic-Adar.
  • Resolução incremental e estabilidade persistente: processe as entradas de lote diárias com correspondência delta incremental.
  • Consolidação de cluster e estabilidade do cluster (sobreposição de 1 a ε): impõe estabilidade persistente do cluster em execuções de pipeline usando uma garantia de limite de sobreposição (1-ε).

O que é necessário

  • Um navegador da Web, como o Chrome.
  • Ter um projeto do Google Cloud com o faturamento ativado.

Este codelab foi criado para engenheiros de dados, desenvolvedores de banco de dados e profissionais de IA/ML de todos os níveis, inclusive iniciantes.

Duração estimada:45 minutos
Custo estimado:menos de US $2,00 (usa o Cloud Functions no modelo de pagamento por uso e o processamento de consultas do BigQuery).

2. Antes de começar

Criar um projeto do Google Cloud

  1. No console do Google Cloud, na página do seletor de projetos, selecione ou crie um projeto na nuvem do Google Cloud.
  2. Verifique se o faturamento está ativado para seu projeto do Cloud. Saiba como verificar se o faturamento está ativado em um projeto.

Iniciar o Cloud Shell

O Cloud Shell é um ambiente de linha de comando executado no Google Cloud que vem pré-carregado com as ferramentas necessárias.

  1. Clique em Ativar o Cloud Shell na parte de cima do console do Google Cloud.
  2. Verifique sua autenticação:
gcloud auth list
  1. Configure variáveis de ambiente no Cloud Shell:
export GCP_PROJECT=$(gcloud config get-value project)
export REGION="us-central1"
export DATASET_ID="identity_resolution"

Ativar APIs obrigatórias

Execute o comando a seguir no Cloud Shell usando sua conta de usuário para ativar todos os serviços necessários do Google Cloud:

gcloud services enable \
  addressvalidation.googleapis.com \
  cloudbuild.googleapis.com \
  cloudfunctions.googleapis.com \
  cloudresourcemanager.googleapis.com \
  artifactregistry.googleapis.com \
  aiplatform.googleapis.com \
  run.googleapis.com \
  bigqueryconnection.googleapis.com \
  bigqueryreservation.googleapis.com \
  bigquery.googleapis.com

Para garantir a execução perfeita da API e o acesso a Application Default Credentials (ADC), crie uma conta de serviço dedicada ao laboratório e ative a representação 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

Criar conjunto de dados do BigQuery

Crie o conjunto de dados do BigQuery para armazenar nós de clientes, arestas, modelos de gráficos e visualizações de avaliação:

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

Você verá uma saída como:

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

Para garantir capacidade de computação dedicada para pesquisas de índice de vetor, agregações de gráficos e execução de função remota sem restrições de limites de CPU sob demanda ou cotas compartilhadas, crie uma reserva da Enterprise Edition com escalonamento automático no 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. Ingerir o conjunto de dados de nós de clientes do FEBRL3

Antes de implantar a função remota de validação de endereço e realizar a resolução de identidade, carregue o conjunto de dados sintético de comparativo de mercado de resolução de entidade FEBRL3 (link em inglês), que contém 5.000 registros de clientes com clusters de vários duplicados de até 5 duplicados por cliente. Para isso, use a biblioteca recordlinkage do Python e grave os nós de clientes brutos (customer_nodes) e os links de correspondência de verdade básica (ground_truth_links) no BigQuery usando o BigQuery DataFrames (bigframes).

Execute os comandos a seguir no Cloud Shell para instalar dependências e executar o script de ingestão:

# 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

No console do Google Cloud, acesse BigQuery Studio, abra uma nova guia Consulta SQL (+) e execute a consulta abaixo para inspecionar a tabela de nós de clientes ingeridos:

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;

Você verá uma saída como:

rec_id

given_name

sobrenome

street_number

address_1

address_2

subúrbio

código postal

estado

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

Observe como o conjunto de dados de comparativo de mercado introduz dados sujos realistas em clusters duplicados:

  • Variações fonéticas e ortográficas: brent x brnt / bernt, wood x woode / wod e clifton x cliffton.
  • Abreviações e erros de digitação de endereços: girdlestone circuit x girdelstone circut / girdlestone cir / girdlestone crt e número da casa 11 x erro de OCR 15.
  • Transposições de caracteres e valores ausentes: transposições de data de nascimento (19340706 x 19340760), estados ausentes ( ) e IDs de segurança social ausentes ( ).

Nas próximas etapas, você vai usar codificações fonéticas SOUNDEX, UDFs de normalização de endereço, distância de edição de Levenshtein e pesquisa vetorial AI.EMBED para resolver essas discrepâncias e vincular com precisão os perfis duplicados.

4. Implantar a UDF da função remota de Address Validation

A normalização de endereços padroniza nomes de ruas, limites de bairros e códigos postais antes de realizar a correspondência. A API Address Validation do Google Maps é um serviço que aceita um endereço, identifica os componentes dele e os valida. Nesta etapa, você vai implantar uma função do Cloud em Python no Cloud Shell que expõe uma UDF de validação e normalização de endereços ao BigQuery.

Gravar arquivos de origem da função do Cloud

Execute o seguinte comando no Cloud Shell para criar o diretório de origem da função do Cloud e gravar main.py e requirements.txt:

mkdir -p cloud_function_address_validation && cd cloud_function_address_validation

cat << 'EOF' > main.py
import os
import json
import logging
import requests
from functools import lru_cache
from concurrent.futures import ThreadPoolExecutor

import functions_framework
import google.auth
from google.auth.transport.requests import AuthorizedSession
from requests.adapters import HTTPAdapter
from urllib3.util import Retry

ADDRESS_VALIDATION_URL = "https://addressvalidation.googleapis.com/v1:validateAddress"
ENABLE_ADDRESS_VALIDATION_API = os.environ.get("ENABLE_ADDRESS_VALIDATION_API", "false").lower() == "true"

# ==========================================
# GLOBAL INITIALIZATION (Runs once per Cold Start)
# ==========================================

credentials, _ = google.auth.default(scopes=["https://www.googleapis.com/auth/cloud-platform"])
session = AuthorizedSession(credentials)

retries = Retry(
    total=4, 
    backoff_factor=0.5, 
    status_forcelist=[429, 500, 502, 503, 504],
    allowed_methods=["POST"]
)
adapter = HTTPAdapter(max_retries=retries, pool_connections=100, pool_maxsize=100)
session.mount("https://", adapter)

executor = ThreadPoolExecutor(max_workers=50)


@lru_cache(maxsize=10000)
def call_validation_api(address_text: str) -> str:
    payload = {
        "address": {
            "regionCode": "AU",
            "addressLines": [address_text]
        }
    }
    
    response = session.post(ADDRESS_VALIDATION_URL, json=payload, timeout=10)
    response.raise_for_status()
        
    res_data = response.json()
    result = res_data.get('result', {})
    address_obj = result.get('address', {})
    verdict = result.get('verdict', {})

    formatted = address_obj.get('formattedAddress', address_text).lower()
    has_unconfirmed = verdict.get('hasUnconfirmedComponents', True)
    address_complete = verdict.get('addressComplete', False)
    granularity = verdict.get('validationGranularity', 'UNCONFIRMED')
    actions = verdict.get('possibleNextActions', [])
    next_action = str(actions[0]) if actions else "NONE"

    is_valid = bool(address_complete and not has_unconfirmed)

    return {
        "formatted_address": formatted,
        "address_is_valid": is_valid,
        "validation_granularity": granularity,
        "possible_next_action": next_action
    }


def process_single_call(call):
    call = call or []
    padded = (call + [""] * 6)[:6]
    
    cleaned_parts = [str(p).strip() if p is not None else "" for p in padded]
    street_num, addr_1, addr_2, suburb, state, postcode = cleaned_parts
    
    address_parts = [p for p in cleaned_parts if p]
    address_text = " ".join(address_parts)

    if not address_text:
        return {
            "formatted_address": "",
            "address_is_valid": False,
            "validation_granularity": "EMPTY",
            "possible_next_action": "NONE"
        }

    if not ENABLE_ADDRESS_VALIDATION_API:
        normalized = (
            address_text.lower()
            .replace("street", "st")
            .replace("road", "rd")
            .replace("place", "pl")
            .replace("avenue", "ave")
            .replace("circuit", "cct")
        )
        return {
            "formatted_address": normalized,
            "address_is_valid": bool(len(address_parts) >= 3),
            "validation_granularity": "PREMISE" if postcode and suburb else "SUBURB",
            "possible_next_action": "NONE"
        }

    try:
        return call_validation_api(address_text)
    except Exception as e:
        logging.error(f"Address Validation API Error for '{address_text}': {str(e)}")
        return {
            "formatted_address": address_text.lower(),
            "address_is_valid": False,
            "validation_granularity": "UNCONFIRMED",
            "possible_next_action": "NONE"
        }


@functions_framework.http
def validate_address_udf(request):
    request_json = request.get_json(silent=True) or {}
    calls = request_json.get('calls', [])

    if not calls:
        return {'replies': []}

    try:
        # executor.map inherently preserves array input order (Strictly required by BigQuery)
        replies = list(executor.map(process_single_call, calls))
        return {'replies': replies}
    except Exception as e:
        logging.error(f"Batch execution failed: {e}")
        return {'errorMessage': str(e)}, 400
EOF

cat << 'EOF' > requirements.txt
functions-framework==3.*
requests==2.*
google-auth==2.*
urllib3==2.*
EOF

Implantar uma função do Cloud e configurar permissões do IAM

Execute estes comandos no Cloud Shell para implantar a função do Cloud de 2ª geração e configurar uma conexão a recursos do Cloud do 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

Você vai ver uma saída indicando que a implantação do Cloud Functions foi concluída e que as vinculações do IAM foram aplicadas com sucesso.

Registrar função remota de normalização de endereços

Agora, registre a DDL da função remota do BigQuery (validate_address_udf) que conecta as linhas da tabela do BigQuery ao endpoint da função do Cloud implantada (${FUNCTION_URL}).

Execute o comando a seguir no Cloud Shell para recuperar o URL da função do Cloud implantada e registrar a função remota automaticamente:

# 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é-processar dados de perfil e codificações fonéticas

Nesta etapa, você vai executar uma consulta de pré-processamento SQL do BigQuery na tabela customer_nodes ingerida.

Executar limpeza de dados e consulta de recursos fonéticos

No editor de SQL do BigQuery Studio, execute a consulta abaixo para criar customer_nodes_cleaned. Esta consulta:

  1. Chama validate_address_udf para receber endereços normalizados e veredictos de validação.
  2. Gera codificações fonéticas SOUNDEX para given_name e surname para lidar com variações de ortografia.
  3. Cria um campo profile_text estruturado.
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;

Consulte a tabela de nós limpos:

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

Você verá uma saída como:

rec_id

given_name

given_name_soundex

sobrenome

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. Gerar embeddings de perfil semântico e pesquisa vetorial

Além da normalização de endereços, das chaves fonéticas Soundex e da distância de edição de Levenshtein, o BigQuery oferece suporte a funções de incorporação de IA generativa integradas usando AI.EMBED.

Usando AI.EMBED, o BigQuery gera incorporações de texto diretamente em SQL usando modelos de base (como text-embedding-005) sem exigir DDL de índice de vetor manual:

Execute as consultas abaixo no editor de SQL do 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. Pontuação de pares de candidatos e fusão de recursos híbridos de borda

A avaliação de todos os pares possíveis de registros de clientes (crescimento quadrático O(N²)) se torna computacionalmente proibitiva à medida que a escala do conjunto de dados aumenta. Em mecanismos de banco de dados relacionais, como o BigQuery, tentar implementar o bloqueio baseado em regras usando condições de junção OR complexas em várias colunas (como junção em a.soc_sec_id = b.soc_sec_id OR a.given_name_soundex = b.given_name_soundex OR ...) impede que o otimizador de consultas use junções hash escalonáveis ou junções de classificação e mesclagem em uma única chave de equi-junção. Em vez disso, o mecanismo volta para uma correlação O(N²) e filtra cada par, o que falha em grande escala.

Nesta etapa, você vai usar os pares candidatos gerados pela nossa tabela de pesquisa vetorial (vector_candidate_edges) e uni-los a customer_nodes_cleaned usando equi-joins rápidos e indexados (ON c.source_id = a.rec_id e ON c.target_id = b.rec_id). Em seguida, você vai calcular uma pontuação de correspondência ponderada combinando:

  • Pontuação de correspondência do CPF (peso: 0.30)
  • Similaridade de edição de sobrenome usando a distância de Levenshtein EDIT_DISTANCE (peso: 0.20)
  • Similaridade de edição do nome próprio (peso: 0.20)
  • Pontuação de correspondência de data de nascimento (peso: 0.15)
  • Similaridade de Jaccard de token de endereço (peso: 0.15) em SPLIT(LOWER(formatted_address), ' ')

Calcular arestas candidatas e pontuações de similaridade ponderadas

Execute a consulta a seguir no editor de SQL do BigQuery Studio para preencher 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;

Inspecione as correspondências de borda candidatas:

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

Você verá uma saída como:

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

Combinar arestas de pesquisa vetorial e baseada em regras em uma tabela unificada

Combine as arestas candidatas da correspondência aproximada baseada em regras e da pesquisa semântica de vetores em uma única tabela final_matched_edges sem duplicações:

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. Construção de gráficos de propriedades e travessias de caminhos GQL ISO

O BigQuery oferece suporte nativo à ISO GQL (Graph Query Language) usando gráficos de propriedades. Um grafo de propriedades cria uma visualização lógica de grafo em tabelas relacionais do BigQuery sem duplicação de dados.

Nesta etapa, você vai criar um gráfico de propriedades customer_identity_graph usando sua tabela unificada de arestas candidatas (final_matched_edges) e consultar conexões de gráficos em perfis de clientes usando a travessia de caminhos de {1, 2} k-hop.

Conectividade transitiva e travessias de K-Hop

Criar DDL de gráfico de propriedades do BigQuery

Execute a seguinte instrução DDL no editor de SQL do 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
);

Visualizar clusters de gráficos K-Hop

Execute a consulta abaixo para visualizar os clusters de clientes correspondentes em 1 a 2 hops de relacionamento:

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;

Visualização de clusters de gráficos K-Hop

Resolver clusters de clientes canônicos

Execute a consulta a seguir para resolver clusters de entidades em 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;

Consulte a tabela de clusters resolvidos:

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

Você verá uma saída como:

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

Criar visualização de métricas de avaliação

Para calcular a precisão, o recall e a pontuação F1 em relação à tabela ground_truth_links, execute:

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;

Consulte a visualização de métricas de avaliação:

SELECT * FROM `identity_resolution.evaluation_metrics`;

Você verá uma saída como:

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

Resolver clusters domésticos usando ponderação de gráfico de Adamic-Adar

Enquanto a resolução de identidade individual resolve registros pertencentes à mesma pessoa, as arquiteturas empresariais do Customer 360 geralmente exigem um agrupamento de entidade familiar de nível superior, que reúne pessoas que moram juntas e compartilham um endereço.

Como os comparativos sintéticos (como o FEBRL3) avaliam informações empíricas no nível individual, a resolução de domicílios é realizada como uma etapa downstream. Sem um histórico de mudança de endereço com carimbo de data/hora, pessoas vinculadas a vários endereços podem causar fusão excessiva ou fragmentação de clusters. Para resolver isso, usamos a ponderação de gráficos de Adamic-Adar para construir associações flexíveis de domicílios.

Execute a consulta abaixo no editor de SQL do BigQuery Studio para preencher household_clusters usando a ponderação de 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;

Consulte a tabela de clusters domésticos resolvidos:

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;

Visualizar a hierarquia de identidade completa usando o GQL

Para rastrear visualmente a hierarquia completa de identidade de três níveis, conectando Clientes brutos não agrupados a Entidades de cliente resolvidas e, em seguida, a Entidades de domicílio resolvidas, execute a seguinte consulta DDL e ISO GQL no 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;

Executar essa consulta GQL no BigQuery Studio renderiza uma tela de visualização de gráfico interativo de três níveis mostrando registros brutos de perfil do cliente (RawCustomer) resolvidos para entidades canônicas individuais (ResolvedCustomer), que estão vinculadas a entidades compartilhadas de domicílio com vários residentes (ResolvedHousehold).

Visualização de clusters de gráficos de domicílios com vários residentes

9. Resolução incremental e estabilidade persistente

Em aplicativos empresariais do mundo real, novos registros de clientes chegam continuamente por lotes diários ou em tempo real. Em vez de executar novamente a resolução completa do gráfico em todo o conjunto de dados históricos, um mecanismo de correspondência delta incremental compara os novos registros recebidos com os clusters de referência resolvidos (resolved_customers).

Para fazer isso de maneira eficiente, o mecanismo usa a Pesquisa Vetorial (

VECTOR_SEARCH

) como uma forma de clusterização dinâmica. Ao tratar cada registro recebido como um ponto de consulta, o VECTOR_SEARCH recupera o conjunto dos K vizinhos mais próximos do índice de incorporação de base histórica. Se um registro recebido corresponder a um perfil de cliente acima do limite de similaridade, ele será mesclado dinamicamente nesse cluster e herdará o canonical_customer_id (MATCHED_TO_EXISTING_CLUSTER) de base. Se nenhum vizinho mais próximo de base for encontrado acima do limite, um novo UUID de entidade será criado (NEW_CUSTOMER_ENTITY).

Ingerir registros de ingestão de lote incremental de amostra

Cole e execute a seguinte DDL no editor de SQL do BigQuery Studio para criar 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
  )
]);

Executar consulta incremental de Delta Match

Execute a consulta a seguir no editor de SQL do BigQuery para fazer a correspondência delta com seu conjunto de dados de comparativo resolvido:

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;

Consulte os resultados da resolução incremental:

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

Você verá uma saída como:

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. Consolidação de clustering e estabilidade de cluster (sobreposição de 1 a ε)

Em sistemas corporativos de produção, as equipes geralmente reagrupam todo o gráfico de maneira recorrente (por exemplo, semanal ou mensalmente) para incorporar novas arestas e fontes de dados. À medida que novas relações são formadas, o reagrupamento completo do gráfico pode fazer com que os identificadores de cluster mudem ou sejam invertidos arbitrariamente nas execuções do pipeline.

Para manter IDs de clientes persistentes para sistemas de CRM, CDP e faturamento downstream, a estabilidade do cluster avalia a sobreposição de nós entre os clusters da execução atual (t) e da execução anterior (t-1) usando um limite de sobreposição (1 - ε), em que ε = 0,30, exigindo uma sobreposição mínima de 70% dos nós.

Se um cluster recém-calculado na execução (t) compartilhar pelo menos 70% dos registros de membros com um cluster da execução (t-1), ele vai herdar o ID do cliente persistente histórico (STABLE_EVOLUTION). Clusters totalmente novos recebem UUIDs recém-gerados (NEW_CLUSTER_CREATED).

Executar consulta de sobreposição e estabilidade de cluster

Execute a consulta a seguir no editor de SQL do BigQuery Studio para preencher 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;

Ver visualização filtrada para registros incrementais

Execute esta consulta para verificar o status de estabilidade do cluster nos seus registros diários de ingestão:

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;

Você verá uma saída como:

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

Para evitar cobranças contínuas na sua conta do Google Cloud, limpe os recursos implantados e o conjunto de dados do BigQuery.

No Cloud Shell, execute:

# 1. Delete BigQuery Dataset
bq rm -r -f -d ${GCP_PROJECT}:${DATASET_ID}

# 2. Delete Cloud Function (2nd-Gen)
gcloud functions delete validate_address_udf --region=${REGION} --gen2 --quiet

# 3. Delete BigQuery Cloud Connection
bq rm -f --connection US.address_val_conn

Se você criou um projeto dedicado do Google Cloud para este laboratório, é possível excluir o projeto:

gcloud projects delete ${GCP_PROJECT}

12. Parabéns

Parabéns! Você criou um mecanismo de resolução de identidade do cliente de ponta a ponta no Google Cloud BigQuery usando o gráfico de propriedades do BigQuery, consultas ISO GQL, correspondência de similaridade híbrida, correspondência delta incremental e garantias de estabilidade de cluster persistente.

O que você aprendeu

  • Como implantar uma função do Cloud de segunda geração e expô-la como uma função remota do BigQuery.
  • Como pré-processar dados demográficos dos clientes usando codificações fonéticas SOUNDEX e validação de endereço.
  • Como executar o bloqueio de candidatos e calcular pontuações de similaridade híbrida usando a distância de Levenshtein (EDIT_DISTANCE) e a similaridade de Jaccard de tokens.
  • Como construir um gráfico de propriedades do BigQuery (CREATE PROPERTY GRAPH) em tabelas de nós e arestas.
  • Como consultar caminhos de gráficos usando a GQL (GRAPH_TABLE) ISO com quantificadores de {1, 2} k-hop.
  • Como resolver clusters de clientes canônicos e avaliar a performance do modelo em relação às métricas de informações empíricas.
  • Como executar a correspondência delta incremental para ingestões diárias em lote sem o reprocessamento completo do conjunto de dados.
  • Como aplicar uma garantia de limite de sobreposição (1 - ε) para manter a estabilidade persistente do cluster em execuções de pipeline.

Próximas etapas

Documentos de referência