Rozwiązywanie problemów z tożsamością klienta za pomocą grafu BigQuery

1. Wprowadzenie

W tym ćwiczeniu utworzysz modułowy silnik rozpoznawania tożsamości klientów (dopasowywania podmiotów) typu end-to-end bezpośrednio w Google Cloud BigQuery. Do wdrażania infrastruktury użyjesz Google Cloud Shell, a do czyszczenia danych, oceniania kandydatów, tworzenia wykresu właściwości i przechodzenia ścieżek w języku GQL (Graph Query Language) zgodnym z ISO – edytora SQL w BigQuery Studio.

Rozpoznawanie tożsamości to podstawowa funkcja w przypadku całościowego profilu klienta w przedsiębiorstwie, wykrywania oszustw i konsolidacji danych z wielu systemów. Istnieje wiele prawidłowych podejść do rozwiązywania problemu tożsamości, w zależności od dojrzałości danych i potrzeb biznesowych, dlatego wszystkie kroki w tym samouczku są modułowe i opcjonalne. Ten potok ma na celu zaprezentowanie różnych typowych technik stosowanych w branży, w tym normalizacji adresów za pomocą zdalnej funkcji zdefiniowanej przez użytkownika, blokowania fonetycznego Soundex, wyszukiwania wektorów semantycznych (AI.EMBED), hybrydowego oceniania cech i klastrowania grafów właściwości GQL. Dzięki temu możesz selektywnie wdrażać wzorce pasujące do Twojej architektury.

Metody dopasowywania i progi punktacji należy dostosować do preferencji organizacji w zakresie dopasowywania deterministycznego i probabilistycznego, które są uzależnione od docelowego przypadku użycia. Na przykład w przypadku ścisłej zgodności z przepisami, rozliczeń lub operacji finansowych zwykle preferowane są deterministyczne reguły o wysokiej precyzji (np. dokładne dopasowania numeru SSN lub identyfikatora podatkowego), aby zapobiec fałszywym połączeniom. Z kolei w przypadku personalizacji marketingowej, analiz i systemów rekomendacji często stosuje się probabilistyczne dopasowywanie przybliżone i podobieństwo wektorów semantycznych, aby zmaksymalizować liczbę wyników i odkryć subtelne powiązania.

Architektura silnika rozwiązywania tożsamości klientów BigQuery

Jakie zadania wykonasz

  • Pozyskiwanie zbioru danych testowych FEBRL3: wczytywanie do BigQuery syntetycznych rekordów klientów i par dopasowań danych podstawowych.
  • Wdróż zdalną funkcję UDF w usłudze Weryfikacja adresów: wdróż funkcję w Cloud Functions w Pythonie i zarejestruj zdalną funkcję BigQuery, aby normalizować adresy ulic.
  • Wstępne przetwarzanie danych profilu i kodowanie fonetyczne: wykonaj czyszczenie danych SQL, wywołaj funkcję UDF adresu i oblicz klucze fonetyczne SOUNDEX oraz odległości edycji Levenshteina:
    • Kodowanie fonetyczne Soundex: algorytm fonetyczny do indeksowania nazw na podstawie ich brzmienia w języku angielskim. Konwertuje nazwy na 4-znakowy kod (pierwsza litera i 3 cyfry) reprezentujący grupy dźwięków spółgłoskowych (np. "John""Jon" są mapowane na J500, a "Smith""Smyth" na S530), co zapewnia sygnały dopasowania fonetycznego na potrzeby oceniania funkcji i blokowania przyrostowego w czasie rzeczywistym.
    • Odległość Levenshteina (EDIT_DISTANCE): miara podobieństwa ciągów znaków określająca minimalną liczbę zmian pojedynczych znaków (wstawień, usunięć lub zamian) wymaganych do przekształcenia jednego ciągu znaków w inny, co umożliwia precyzyjne dopasowywanie nazw i adresów.
  • Generowanie wektorów dystrybucyjnych profilu semantycznego i wyszukiwanie wektorowe: generuj wektory dystrybucyjne tekstu bezpośrednio w SQL za pomocą funkcji AI.EMBED (text-embedding-005) i znajdź K najbliższych sąsiadów za pomocą funkcji VECTOR_SEARCH, aby służyły jako podliniowa warstwa generowania kandydatów.
  • Ocena par kandydatów i hybrydowe łączenie cech krawędzi: wykorzystaj pary kandydatów wyszukiwania wektorowego, aby wyeliminować złożoność połączenia krzyżowego O(N²), obliczyć ważone wyniki podobieństwa wielu cech (SSN, odległość edycji Levenshteina, data urodzenia, adres Jaccarda) i połączyć krawędzie w ujednoliconą tabelę kandydatów.
  • Tworzenie wykresu właściwości i przechodzenie ścieżek ISO GQL: utwórz PROPERTY GRAPH BigQuery, uruchom zapytania o ścieżki ISO GQL {1, 2} (GRAPH_TABLE), aby rozwiązać połączone klastry klientów, obliczyć indywidualne wskaźniki oceny i wykonać miękkie klastrowanie gospodarstw domowych za pomocą ważenia wykresu Adamic-Adar.
  • Przyrostowa rozdzielczość i trwała stabilność: przetwarzaj codzienne partie danych za pomocą przyrostowego dopasowywania różnic.
  • Konsolidacja klastrów i stabilność klastrów (nakładanie się 1-ε): wymusza trwałą stabilność klastrów w różnych przebiegach potoku za pomocą gwarancji progu nakładania się (1-ε).

Czego potrzebujesz

  • przeglądarka, np. Chrome;
  • projekt Google Cloud z włączonymi płatnościami;

To ćwiczenie jest przeznaczone dla inżynierów danych, programistów baz danych i praktyków AI/ML na wszystkich poziomach zaawansowania, w tym dla początkujących.

Szacowany czas trwania: 45 minut.
Szacowany koszt: poniżej 2 USD (korzysta z płatności według wykorzystania w przypadku Cloud Functions i przetwarzania zapytań w BigQuery).

2. Zanim zaczniesz

Tworzenie projektu Google Cloud

  1. W konsoli Google Cloud na stronie wyboru projektu wybierz lub utwórz projekt w chmurze Google Cloud.
  2. Sprawdź, czy w projekcie Cloud włączone są płatności. Dowiedz się, jak sprawdzić, czy w projekcie są włączone płatności.

Uruchamianie Cloud Shell

Cloud Shell to środowisko wiersza poleceń działające w Google Cloud, które zawiera niezbędne narzędzia.

  1. U góry konsoli Google Cloud kliknij Aktywuj Cloud Shell.
  2. Potwierdź uwierzytelnianie:
gcloud auth list
  1. Skonfiguruj zmienne środowiskowe w Cloud Shell:
export GCP_PROJECT=$(gcloud config get-value project)
export REGION="us-central1"
export DATASET_ID="identity_resolution"

Włącz wymagane interfejsy API

Aby włączyć wszystkie wymagane usługi Google Cloud, uruchom w Cloud Shell to polecenie, używając swojego konta użytkownika:

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

Aby zapewnić płynne wykonywanie interfejsu API i dostęp do domyślnego uwierzytelniania aplikacji (ADC), utwórz dedykowane konto usługi laboratorium i włącz gcloud podszywanie się:

# 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

Tworzenie zbioru danych BigQuery

Utwórz zbiór danych BigQuery do przechowywania węzłów klientów, krawędzi, modeli wykresów i widoków oceny:

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

Dane wyjściowe powinny być podobne do tych:

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

Aby zapewnić dedykowaną moc obliczeniową na potrzeby wyszukiwania indeksów wektorowych, agregacji grafów i wykonywania funkcji zdalnych bez ograniczeń związanych z limitami procesora na żądanie lub limitami wspólnymi, utwórz rezerwację w wersji Enterprise z automatycznym skalowaniem w 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. Pozyskiwanie zbioru danych węzła klienta FEBRL3

Przed wdrożeniem funkcji zdalnej weryfikacji adresu i przeprowadzeniem rozpoznawania tożsamości wczytasz syntetyczny zbiór danych testowych FEBRL3 do rozpoznawania tożsamości (zawierający 5000 rekordów klientów z wieloma zduplikowanymi klastrami, w których na każdego klienta przypada do 5 duplikatów) za pomocą biblioteki recordlinkage w Pythonie i zapiszesz surowe węzły klientów (customer_nodes) oraz prawdziwe linki do dopasowań (ground_truth_links) w BigQuery za pomocą BigQuery DataFrames (bigframes).

Aby zainstalować zależności i uruchomić skrypt pozyskiwania, uruchom w Cloud Shell te polecenia:

# 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

W konsoli Google Cloud otwórz BigQuery Studio, otwórz nową kartę zapytania SQL (+) i uruchom poniższe zapytanie, aby sprawdzić wczytaną tabelę węzłów klientów:

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;

Dane wyjściowe powinny być podobne do tych:

rec_id

given_name

nazwisko

street_number

address_1

address_2

przedmieście

kod pocztowy

stan

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

Zwróć uwagę, jak zbiór danych testu porównawczego wprowadza realistyczne dane zanieczyszczone w zduplikowanych klastrach:

  • Warianty fonetyczne i pisowni: brentbrnt / bernt, woodwoode / wod oraz cliftoncliffton.
  • Skróty i literówki w adresach: girdlestone circuit vs. girdelstone circut / girdlestone cir / girdlestone crt oraz numer domu 11 vs. błąd OCR 15.
  • Transpozycje znaków i brakujące wartości: transpozycje daty urodzenia (19340706 vs. 19340760), brakujące stany ( ) i brakujące numery ubezpieczenia społecznego ( ).

W kolejnych krokach użyjesz SOUNDEXkodowania fonetycznego, funkcji zdefiniowanych przez użytkownika do normalizacji adresów, odległości edycji Levenshteina i AI.EMBED wyszukiwania wektorowego, aby wyeliminować te rozbieżności i dokładnie połączyć zduplikowane profile.

4. Wdrażanie funkcji zdalnej UDF weryfikacji adresu

Normalizacja adresu ujednolica nazwy ulic, granice przedmieść i kody pocztowe przed dopasowaniem. Interfejs Google Maps Address Validation API to usługa, która akceptuje adres, identyfikuje jego komponenty i je weryfikuje. W tym kroku wdrożysz w Cloud Shell funkcję w Pythonie w Cloud Functions, która udostępnia w BigQuery funkcję UDF do weryfikacji i normalizacji adresów.

Pisanie plików źródłowych funkcji w Cloud Functions

Uruchom w Cloud Shell to polecenie, aby utworzyć katalog źródłowy funkcji w Cloud Functions i zapisać pliki main.py i 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

Wdrażanie funkcji w Cloud Functions i konfigurowanie uprawnień

Aby wdrożyć funkcję Cloud Functions 2 generacji i skonfigurować połączenie z zasobem Cloud BigQuery, wykonaj w Cloud Shell te polecenia:

# 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

Powinny być widoczne dane wyjściowe wskazujące, że wdrożenie funkcji w Cloud Functions zostało zakończone, a powiązania IAM zostały zastosowane.

Rejestrowanie funkcji normalizacji adresu zdalnego

Teraz zarejestrujesz DDL funkcji zdalnej BigQuery (validate_address_udf), która łączy wiersze tabeli BigQuery z wdrożonym punktem końcowym funkcji Cloud (${FUNCTION_URL}).

Aby pobrać adres URL wdrożonej funkcji w Cloud Functions i automatycznie zarejestrować funkcję zdalną, uruchom w Cloud Shell to polecenie:

# 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. Wstępne przetwarzanie danych profilu i kodowanie fonetyczne

W tym kroku wykonasz zapytanie SQL w BigQuery, które przetworzy wstępnie zaimportowaną tabelę customer_nodes.

Wykonywanie oczyszczania danych i zapytań o cechy fonetyczne

W edytorze SQL BigQuery Studio uruchom to zapytanie, aby utworzyć tabelę customer_nodes_cleaned. To zapytanie:

  1. Wywołuje funkcję validate_address_udf, aby uzyskać znormalizowane adresy i wyniki weryfikacji.
  2. Generuje SOUNDEX kody fonetyczne dla given_namesurname, aby obsługiwać różne pisownie.
  3. Tworzy pole strukturalne profile_text.
CREATE OR REPLACE TABLE `identity_resolution.customer_nodes_cleaned` AS
WITH raw_data AS (
  SELECT 
    rec_id, dataset_source,
    TRIM(LOWER(given_name)) AS given_name_clean,
    TRIM(LOWER(surname)) AS surname_clean,
    `identity_resolution.validate_address_udf`(street_number, address_1, address_2, suburb, state, postcode) AS addr_json,
    TRIM(suburb) AS suburb, TRIM(state) AS state, TRIM(postcode) AS postcode,
    TRIM(date_of_birth) AS date_of_birth, TRIM(soc_sec_id) AS soc_sec_id
  FROM `identity_resolution.customer_nodes`
)
SELECT
  rec_id, dataset_source,
  given_name_clean AS given_name,
  SOUNDEX(given_name_clean) AS given_name_soundex,
  surname_clean AS surname,
  SOUNDEX(surname_clean) AS surname_soundex,
  CONCAT(given_name_clean, ' ', surname_clean) AS full_name,
  STRING(addr_json.formatted_address) AS formatted_address,
  BOOL(addr_json.address_is_valid) AS address_is_valid,
  STRING(addr_json.validation_granularity) AS validation_granularity,
  STRING(addr_json.possible_next_action) AS possible_next_action,
  suburb, state, postcode, date_of_birth, soc_sec_id,
  CONCAT('Name: ', CONCAT(given_name_clean, ' ', surname_clean), '; Address: ', STRING(addr_json.formatted_address), '; DOB: ', date_of_birth, '; SSN: ', soc_sec_id) AS profile_text
FROM raw_data;

Utwórz zapytanie do tabeli oczyszczonych węzłów:

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

Dane wyjściowe powinny być podobne do tych:

rec_id

given_name

given_name_soundex

nazwisko

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. Generowanie wektorów dystrybucyjnych profilu semantycznego i wyszukiwanie wektorowe

Oprócz normalizacji adresów, kluczy fonetycznych Soundex i odległości edycji Levenshteina BigQuery obsługuje wbudowane funkcje generatywnej AI do tworzenia wektorów za pomocą funkcji AI.EMBED.

Za pomocą AI.EMBED BigQuery generuje osadzanie tekstu bezpośrednio w SQL przy użyciu modeli podstawowych (takich jak text-embedding-005) bez konieczności ręcznego tworzenia instrukcji DDL indeksu wektorowego:

Uruchom poniższe zapytania w edytorze SQL 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. Ocena par kandydatów i łączenie funkcji hybrydowych

Ocena wszystkich możliwych par rekordów klientów (wzrost kwadratowy O(N²)) staje się obliczeniowo niemożliwa wraz ze wzrostem skali zbioru danych. W relacyjnych silnikach baz danych, takich jak BigQuery, próba wdrożenia blokowania opartego na regułach przy użyciu złożonych warunków OR łączenia w wielu kolumnach (np. łączenie na podstawie a.soc_sec_id = b.soc_sec_id OR a.given_name_soundex = b.given_name_soundex OR ...) uniemożliwia optymalizatorowi zapytań używanie skalowalnych łączeń typu hash join lub sort-merge join na pojedynczym kluczu łączenia równościowego. Zamiast tego silnik wraca do połączenia iloczynowego O(N²) i filtruje każdą parę, co nie sprawdza się w przypadku dużej skali.

W tym kroku weźmiesz pary kandydatów wygenerowane przez naszą tabelę wyszukiwania wektorowego (vector_candidate_edges) i połączysz je z tabelą customer_nodes_cleaned za pomocą szybkich, indeksowanych złączeń równościowych (ON c.source_id = a.rec_idON c.target_id = b.rec_id). Następnie obliczysz ważoną ocenę dopasowania, która będzie łączyć:

  • Wynik dopasowania numeru SSN (waga: 0.30)
  • Podobieństwo edycji nazwiska przy użyciu odległości Levenshteina EDIT_DISTANCE (waga: 0.20)
  • Podobieństwo edycji imienia (waga: 0.20)
  • Wynik dopasowania daty urodzenia (waga: 0.15)
  • Podobieństwo Jaccarda tokenów adresu (waga: 0.15) w przypadku SPLIT(LOWER(formatted_address), ' ')

Obliczanie krawędzi kandydatów i ważonych wyników podobieństwa

Aby wypełnić tabelę matched_edges, uruchom w edytorze SQL BigQuery Studio to zapytanie:

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;

Sprawdź dopasowania krawędzi kandydatów:

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

Dane wyjściowe powinny być podobne do tych:

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

Łączenie krawędzi wyszukiwania opartego na regułach i wyszukiwania wektorowego w ujednoliconej tabeli

Połącz krawędzie kandydatów z regułowego dopasowywania przybliżonego i semantycznego wyszukiwania wektorowego w jedną tabelę final_matched_edges bez duplikatów:

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. Tworzenie wykresów właściwości i przechodzenie ścieżek w języku GQL zgodnym z ISO

BigQuery obsługuje natywnie język ISO GQL (Graph Query Language) za pomocą grafów właściwości. Graf właściwości tworzy logiczny widok grafu na relacyjnych tabelach BigQuery bez powielania danych.

W tym kroku utworzysz wykres właściwości customer_identity_graph za pomocą ujednoliconej tabeli krawędzi kandydatów (final_matched_edges) i wykresu zapytań dotyczących połączeń między profilami klientów za pomocą {1, 2} przeszukiwania ścieżek k-hop.

Łączność przechodnia i przejścia K-hop

Tworzenie DDL grafu właściwości BigQuery

Wykonaj w edytorze SQL BigQuery Studio tę instrukcję DDL:

CREATE OR REPLACE PROPERTY GRAPH `identity_resolution.customer_identity_graph`
NODE TABLES (
  `identity_resolution.customer_nodes_cleaned` AS `Customer`
  KEY (rec_id)
)
EDGE TABLES (
  `identity_resolution.final_matched_edges`
  KEY (source_id, target_id)
  SOURCE KEY (source_id) REFERENCES `Customer`(rec_id)
  DESTINATION KEY (target_id) REFERENCES `Customer`(rec_id)
  LABEL MATCHED_TO
);

Wizualizacja klastrów grafu k-hop

Aby wizualizować pasujące klastry klientów w przypadku 1–2 relacji, uruchom to zapytanie:

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;

Wizualizacja klastrów grafu k-hopu

Rozwiązywanie problemów z kanonicznymi klastrami klientów

Aby przekształcić klastry encji w resolved_customers, wykonaj to zapytanie:

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;

Utwórz zapytanie do tabeli rozwiązanych klastrów:

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

Dane wyjściowe powinny być podobne do tych:

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

Tworzenie widoku wskaźników oceny

Aby obliczyć precyzję, czułość i wynik F1 w odniesieniu do tabeli ground_truth_links, uruchom to polecenie:

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;

Wysyłaj zapytania do widoku wskaźników oceny:

SELECT * FROM `identity_resolution.evaluation_metrics`;

Dane wyjściowe powinny być podobne do tych:

total_ground_truth

total_predictions

true_positives

false_positives

false_negatives

precyzja

wycofanie

f1_score

6538

5620

5608

13

930

0.9977

0.8579

0.9225

Rozwiązywanie problemów z grupami domowymi za pomocą ważenia grafów Adamic-Adar

Rozpoznawanie tożsamości poszczególnych osób pozwala łączyć rekordy należące do tej samej osoby, ale architektury Customer 360 dla przedsiębiorstw często wymagają grupowania na wyższym poziomie, czyli podmiotu gospodarstwa domowego, które obejmuje osoby mieszkające razem pod tym samym adresem.

Ponieważ syntetyczne testy porównawcze (np. FEBRL3) oceniają dane podstawowe na poziomie indywidualnym, rozdzielczość gospodarstwa domowego jest wykonywana jako krok końcowy. W przypadku braku historii przeniesień z sygnaturami czasowymi osoby powiązane z wieloma adresami mogą powodować nadmierne łączenie lub fragmentację klastrów. Aby rozwiązać ten problem, używamy ważenia grafu Adamic-Adar do tworzenia miękkich członkostw w gospodarstwach domowych.

Aby wypełnić tabelę household_clusters za pomocą ważenia grafu Adamic-Adar, uruchom w edytorze SQL BigQuery Studio to zapytanie:

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;

Utwórz zapytanie do tabeli rozwiązanych klastrów gospodarstw domowych:

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;

Wizualizacja pełnej hierarchii tożsamości za pomocą GQL

Aby wizualnie prześledzić całą 3-poziomową hierarchię tożsamości, łącząc niepogrupowanych surowych klientów z rozwiązanymi podmiotami klienta i dalej z rozwiązanymi podmiotami gospodarstwa domowego, wykonaj w BigQuery Studio to zapytanie DDL i ISO GQL:

-- 1. Create Household Node Table
CREATE OR REPLACE TABLE `identity_resolution.household_nodes` AS
SELECT DISTINCT 
  canonical_household_id, 
  household_address, 
  total_residents
FROM `identity_resolution.household_clusters`;

-- 2. Create Unresolved Record to Resolved Entity Edge Table
CREATE OR REPLACE TABLE `identity_resolution.customer_entity_edges` AS
SELECT DISTINCT
  rec_id,
  canonical_customer_id
FROM `identity_resolution.resolved_customers`,
UNNEST(customer_records) AS rec_id;

-- 3. Create Primary Household Edge Table (Highest Weighted Household Rank = 1)
CREATE OR REPLACE TABLE `identity_resolution.primary_household_edges` AS
SELECT 
  canonical_customer_id,
  canonical_household_id,
  household_membership_weight,
  household_rank
FROM `identity_resolution.household_clusters`
WHERE household_rank = 1;

-- 4. Update Unified Property Graph DDL
CREATE OR REPLACE PROPERTY GRAPH `identity_resolution.customer_identity_graph`
NODE TABLES (
  `identity_resolution.customer_nodes_cleaned` AS `RawCustomer`
    KEY (rec_id),
  `identity_resolution.resolved_customers` AS `ResolvedCustomer`
    KEY (canonical_customer_id),
  `identity_resolution.household_nodes` AS `ResolvedHousehold`
    KEY (canonical_household_id)
)
EDGE TABLES (
  `identity_resolution.final_matched_edges`
    KEY (source_id, target_id)
    SOURCE KEY (source_id) REFERENCES `RawCustomer`(rec_id)
    DESTINATION KEY (target_id) REFERENCES `RawCustomer`(rec_id)
    LABEL MATCHED_TO,
  `identity_resolution.customer_entity_edges`
    KEY (rec_id, canonical_customer_id)
    SOURCE KEY (rec_id) REFERENCES `RawCustomer`(rec_id)
    DESTINATION KEY (canonical_customer_id) REFERENCES `ResolvedCustomer`(canonical_customer_id)
    LABEL RESOLVED_TO,
  `identity_resolution.primary_household_edges`
    KEY (canonical_customer_id, canonical_household_id)
    SOURCE KEY (canonical_customer_id) REFERENCES `ResolvedCustomer`(canonical_customer_id)
    DESTINATION KEY (canonical_household_id) REFERENCES `ResolvedHousehold`(canonical_household_id)
    LABEL BELONGS_TO_HOUSEHOLD
);

-- 5. Execute 3-Tier GQL Query for Multi-Resident Household Visualization
GRAPH `identity_resolution.customer_identity_graph`
MATCH p = (raw:RawCustomer)-[e1:RESOLVED_TO]->(c:ResolvedCustomer)-[e2:BELONGS_TO_HOUSEHOLD]->(h:ResolvedHousehold)
WHERE h.total_residents > 1
RETURN TO_JSON(p) AS multi_resident_household_hierarchy_path
LIMIT 15;

Uruchomienie tego zapytania GQL w BigQuery Studio spowoduje wyrenderowanie 3-poziomowego interaktywnego wykresu wizualizacji przedstawiającego surowe rekordy profili klientów (RawCustomer) przekształcone w poszczególne jednostki kanoniczne (ResolvedCustomer), które są połączone ze wspólnymi jednostkami gospodarstw domowych z wieloma mieszkańcami (ResolvedHousehold).

Wizualizacja klastrów grafu gospodarstw domowych z wieloma mieszkańcami

9. Stopniowe zwiększanie rozdzielczości i stała stabilność

W przypadku zastosowań w przedsiębiorstwach w świecie rzeczywistym nowe rekordy klientów są przesyłane w sposób ciągły w ramach codziennych lub działających w czasie rzeczywistym partii danych. Zamiast ponownie uruchamiać pełną rozdzielczość wykresu w całym historycznym zbiorze danych, przyrostowy silnik dopasowywania różnic porównuje nowe przychodzące rekordy z istniejącymi rozwiązanymi klastrami bazowymi (resolved_customers).

Aby to osiągnąć, mechanizm korzysta z wyszukiwania wektorowego

VECTOR_SEARCH

) jako formę dynamicznego grupowania. Traktując każdy przychodzący rekord jako punkt zapytania, VECTOR_SEARCH pobiera z historycznego indeksu osadzania wartości bazowych zbiór K najbliższych sąsiadów. Jeśli przychodzący rekord pasuje do istniejącego profilu klienta powyżej progu podobieństwa, jest dynamicznie scalany z tym klastrem i dziedziczy wartość bazową canonical_customer_id (MATCHED_TO_EXISTING_CLUSTER). Jeśli nie zostanie znaleziony najbliższy sąsiad powyżej progu, zostanie utworzony nowy identyfikator UUID jednostki (NEW_CUSTOMER_ENTITY).

Przetwarzanie przykładowych rekordów przyrostowego wsadu

Wklej i wykonaj w edytorze SQL BigQuery Studio poniższy kod DDL, aby utworzyć 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
  )
]);

Wykonaj przyrostowe zapytanie o dopasowanie delta

Aby przeprowadzić dopasowywanie różnic w odniesieniu do rozwiązanego zbioru danych bazowych, uruchom w edytorze SQL BigQuery to zapytanie:

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;

Wykonywanie zapytań o wyniki zwiększania rozdzielczości:

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

Dane wyjściowe powinny być podobne do tych:

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. Konsolidacja klastrów i stabilność klastrów (nakładanie się 1-ε)

W systemach produkcyjnych klasy enterprise zespoły zwykle ponownie grupują cały wykres w sposób cykliczny (np. co tydzień lub co miesiąc), aby uwzględnić nowe krawędzie i źródła danych. W miarę tworzenia się nowych relacji pełne ponowne grupowanie w klastry może powodować dowolne przesuwanie lub odwracanie identyfikatorów klastrów w różnych wykonaniach potoku.

Aby zachować trwałe identyfikatory klientów w systemach CRM, CDP i rozliczeniowych, stabilność klastra ocenia nakładanie się węzłów między klastrami bieżącego uruchomienia (t) a klastrami poprzedniego uruchomienia (t-1) przy użyciu progu nakładania się (1 – ε), gdzie ε = 0, 30 (wymagane jest co najmniej 70% nakładanie się węzłów).

Jeśli nowo obliczony klaster w ramach działania (t) ma co najmniej 70% rekordów członkowskich wspólnych z klastrem z działania (t-1), dziedziczy historyczny trwały identyfikator klienta (STABLE_EVOLUTION). Nowe klastry otrzymują nowo wygenerowane identyfikatory UUID (NEW_CLUSTER_CREATED).

Wykonywanie zapytania o nakładanie się i stabilność klastrów

Aby wypełnić tabelę stable_resolved_customers, uruchom w edytorze SQL BigQuery Studio to zapytanie:

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;

Wyświetlanie odfiltrowanego podglądu rekordów przyrostowych

Uruchom to zapytanie, aby sprawdzić stan stabilności klastra w przypadku codziennych rekordów danych:

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;

Dane wyjściowe powinny być podobne do tych:

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. Czyszczenie danych

Aby uniknąć obciążenia konta Google Cloud bieżącymi opłatami, zwalniaj miejsce zajmowane przez wdrożone zasoby i zbiór danych BigQuery.

W Cloud Shell uruchom:

# 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

Jeśli na potrzeby tego laboratorium utworzono specjalny projekt w chmurze Google, możesz go usunąć:

gcloud projects delete ${GCP_PROJECT}

12. Gratulacje

Gratulacje! Udało Ci się utworzyć w Google Cloud BigQuery kompleksowy mechanizm rozwiązywania tożsamości klientów za pomocą wykresu właściwości BigQuery, zapytań ISO GQL, hybrydowego dopasowywania podobieństwa, przyrostowego dopasowywania różnic i gwarancji stabilności trwałych klastrów.

Czego się nauczysz

  • Jak wdrożyć funkcję Cloud Functions 2 generacji i udostępnić ją jako funkcję zdalną BigQuery.
  • Jak wstępnie przetwarzać dane demograficzne klientów za pomocą SOUNDEX kodowania fonetycznego i weryfikacji adresu.
  • Jak blokować kandydatów i obliczać hybrydowe wyniki podobieństwa za pomocą odległości Levenshteina (EDIT_DISTANCE) i podobieństwa Jaccarda tokenów.
  • Jak utworzyć wykres właściwości BigQuery (CREATE PROPERTY GRAPH) na podstawie tabel węzłów i krawędzi.
  • Jak wysyłać zapytania o ścieżki w grafie za pomocą języka GQL (GRAPH_TABLE) zgodnego z normą ISO z kwantyfikatorami {1, 2} k-hop.
  • Jak rozwiązywać problemy z kanonicznymi klastrami klientów i oceniać skuteczność modelu na podstawie wskaźników danych podstawowych.
  • Jak przeprowadzać przyrostowe dopasowywanie różnic w przypadku codziennych partii danych bez ponownego przetwarzania pełnego zbioru danych.
  • Jak zastosować gwarancję progu nakładania się (1 – ε), aby utrzymać stałą stabilność klastra w różnych przebiegach potoku.

Dalsze kroki

Dokumentacja