Penyelesaian Identitas Pelanggan dengan BigQuery Graph

1. Pengantar

Dalam codelab ini, Anda akan membangun mesin Penyelesaian Identitas Pelanggan (Pencocokan Entitas) end-to-end yang modular langsung di dalam Google Cloud BigQuery. Anda akan menggabungkan Google Cloud Shell untuk deployment infrastruktur dengan Editor SQL BigQuery Studio untuk pembersihan data, pemberian skor kandidat, pembuatan grafik properti, dan penelusuran jalur GQL (Graph Query Language) ISO.

Penyelesaian identitas adalah kemampuan mendasar untuk Customer 360 perusahaan, deteksi penipuan, dan konsolidasi data multi-sistem. Karena ada banyak pendekatan yang valid untuk penyelesaian identitas, bergantung pada kematangan data dan kebutuhan bisnis, semua langkah dalam codelab ini bersifat modular dan opsional. Pipeline ini dirancang untuk menampilkan berbagai teknik umum tingkat produksi yang umum digunakan di industri—termasuk normalisasi alamat UDF Jarak Jauh, pemblokiran fonetik Soundex, penelusuran vektor semantik (AI.EMBED), pemberian skor fitur hibrida, dan pengelompokan grafik properti GQL—sehingga Anda dapat secara selektif mengadopsi pola yang sesuai dengan arsitektur Anda.

Metode pencocokan dan nilai minimum skor harus disesuaikan berdasarkan keinginan organisasi Anda untuk pencocokan deterministik vs. probabilistik, yang ditentukan oleh kasus penggunaan target. Misalnya, kepatuhan, penagihan, atau operasi keuangan yang ketat biasanya lebih memilih aturan deterministik presisi tinggi (seperti kecocokan SSN atau ID Pajak yang tepat) untuk mencegah penautan palsu, sedangkan personalisasi pemasaran, analisis, dan mesin rekomendasi sering kali mengandalkan pencocokan fuzzy probabilistik dan kesamaan vektor semantik untuk memaksimalkan perolehan dan mengungkap koneksi yang tidak terlihat.

Arsitektur Customer Identity Resolution Engine BigQuery

Yang akan Anda lakukan

  • Menyerap FEBRL3 Benchmark Dataset: Muat pasangan kecocokan kebenaran dasar dan data pelanggan sintetis ke BigQuery.
  • Deploy Address Validation Remote UDF: Deploy Cloud Function Python dan daftarkan Fungsi Jarak Jauh BigQuery untuk menormalisasi alamat.
  • Memproses Awal Data Profil & Encoding Fonetik: Jalankan pembersihan data SQL, panggil UDF alamat, dan hitung jarak edit Levenshtein dan kunci fonetik SOUNDEX:
    • Encoding Fonetik Soundex: Algoritma fonetik untuk mengindeks nama berdasarkan bunyi seperti yang diucapkan dalam bahasa Inggris. Fitur ini mengonversi nama menjadi kode 4 karakter (huruf awal diikuti dengan tiga digit) yang merepresentasikan kelompok bunyi konsonan (misalnya, "John" dan "Jon" dipetakan ke J500, sedangkan "Smith" dan "Smyth" dipetakan ke S530), sehingga memberikan sinyal pencocokan fonetik untuk penskoran fitur dan pemblokiran delta inkremental real-time.
    • Jarak Levenshtein (EDIT_DISTANCE): Metrik string yang mengukur jumlah minimum pengeditan satu karakter (penyisipan, penghapusan, atau penggantian) yang diperlukan untuk mengubah satu string menjadi string lain, sehingga memungkinkan pencocokan nama dan alamat fuzzy yang akurat.
  • Membuat Embedding Profil Semantik & Penelusuran Vektor: Buat embedding teks langsung di SQL menggunakan AI.EMBED (text-embedding-005) dan temukan K tetangga terdekat teratas menggunakan VECTOR_SEARCH untuk berfungsi sebagai lapisan pembuatan kandidat sub-linear.
  • Penskoran Pasangan Kandidat & Penggabungan Fitur Edge Hibrida: Memanfaatkan pasangan kandidat penelusuran vektor untuk menghilangkan kompleksitas cross-join O(N²), menghitung skor kemiripan berbobot multi-fitur (SSN, jarak edit Levenshtein, DOB, Jaccard alamat), dan menggabungkan edge ke dalam tabel kandidat terpadu.
  • Konstruksi Grafik Properti & Traversal Jalur GQL ISO: Buat PROPERTY GRAPH BigQuery, jalankan kueri jalur {1, 2} GQL ISO (GRAPH_TABLE) untuk menyelesaikan pengelompokan pelanggan yang terhubung, hitung metrik evaluasi individual, dan lakukan pengelompokan rumah tangga ringan menggunakan pembobotan grafik Adamic-Adar.
  • Resolusi Bertahap & Stabilitas Persisten: Memproses penyerapan batch harian dengan pencocokan delta inkremental.
  • Penggabungan Pengelompokan & Stabilitas Cluster (Tumpang-Tindih 1-ε): Terapkan stabilitas cluster persisten di seluruh proses pipeline menggunakan jaminan nilai minimum tumpang-tindih (1-ε).

Yang Anda butuhkan

  • Browser web seperti Chrome.
  • Project Google Cloud yang mengaktifkan penagihan.

Codelab ini dirancang untuk engineer data, developer database, dan praktisi AI/ML dari semua tingkat, termasuk pemula.

Perkiraan Durasi: 45 menit
Perkiraan Biaya: Kurang dari $2,00 USD (menggunakan pemrosesan kueri BigQuery dan Cloud Functions bayar sesuai penggunaan).

2. Sebelum memulai

Buat Project Google Cloud

  1. Di Konsol Google Cloud, di halaman pemilih project, pilih atau buat project Google Cloud.
  2. Pastikan penagihan diaktifkan untuk project Cloud Anda. Pelajari cara memeriksa apakah penagihan telah diaktifkan pada suatu project.

Mulai Cloud Shell

Cloud Shell adalah lingkungan command line yang berjalan di Google Cloud yang telah dilengkapi dengan alat yang diperlukan.

  1. Klik Activate Cloud Shell di bagian atas konsol Google Cloud.
  2. Verifikasi autentikasi Anda:
gcloud auth list
  1. Konfigurasi variabel lingkungan di Cloud Shell:
export GCP_PROJECT=$(gcloud config get-value project)
export REGION="us-central1"
export DATASET_ID="identity_resolution"

Mengaktifkan API yang Diperlukan

Jalankan perintah berikut di Cloud Shell menggunakan akun pengguna Anda untuk mengaktifkan semua layanan Google Cloud yang diperlukan:

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

Untuk memastikan eksekusi API yang lancar dan akses Kredensial Default Aplikasi (ADC), buat Akun Layanan lab khusus dan aktifkan peniruan identitas 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

Buat Set Data BigQuery

Buat set data BigQuery untuk menyimpan node pelanggan, tepi, model grafik, dan tampilan evaluasi:

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

Anda akan melihat output yang mirip dengan:

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

Untuk memastikan kapasitas komputasi khusus untuk penelusuran indeks vektor, agregasi grafik, dan eksekusi Fungsi Jarak Jauh tanpa dibatasi oleh batas CPU Sesuai Permintaan atau kuota bersama, buat reservasi Enterprise Edition dengan penskalaan otomatis di 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. Menyerap Set Data Node Pelanggan FEBRL3

Sebelum men-deploy fungsi jarak jauh validasi alamat dan melakukan penyelesaian identitas, Anda akan memuat kumpulan data tolok ukur penyelesaian entitas FEBRL3 sintetis (yang berisi 5.000 data pelanggan dengan beberapa duplikat hingga 5 duplikat per pelanggan) menggunakan library recordlinkage Python dan menulis node pelanggan mentah (customer_nodes) dan link kecocokan kebenaran dasar (ground_truth_links) ke BigQuery menggunakan BigQuery DataFrames (bigframes).

Jalankan perintah berikut di Cloud Shell untuk menginstal dependensi dan menjalankan skrip penyerapan:

# 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

Di Konsol Google Cloud, buka BigQuery Studio, buka tab Kueri SQL baru (+), lalu jalankan kueri di bawah untuk memeriksa tabel node pelanggan yang di-ingest:

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;

Anda akan melihat output yang mirip dengan:

rec_id

given_name

marga

street_number

address_1

address_2

pinggiran kota

kode pos

dengan status tersembunyi akhir

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

Perhatikan bagaimana set data tolok ukur memperkenalkan data kotor yang realistis di seluruh cluster duplikat:

  • Variasi fonetik & ejaan: brent vs. brnt / bernt, wood vs. woode / wod, dan clifton vs. cliffton.
  • Singkatan & salah ketik alamat: girdlestone circuit vs. girdelstone circut / girdlestone cir / girdlestone crt, dan nomor rumah 11 vs. kesalahan OCR 15.
  • Transposisi karakter & nilai yang tidak ada: Transposisi tanggal lahir (19340706 vs. 19340760), negara bagian yang tidak ada ( ), dan ID Jaminan Sosial yang tidak ada ( ).

Pada langkah-langkah berikutnya, Anda akan menggunakan SOUNDEX encoding fonetik, UDF normalisasi alamat, jarak pengeditan Levenshtein, dan penelusuran vektor AI.EMBED untuk menjembatani perbedaan ini dan menautkan profil duplikat secara akurat.

4. Men-deploy UDF Fungsi Jarak Jauh Address Validation

Normalisasi alamat menstandarkan nama jalan, batas pinggiran kota, dan kode pos sebelum melakukan pencocokan. Google Maps Address Validation API adalah layanan yang menerima alamat, mengidentifikasi komponen alamat, dan memvalidasinya. Pada langkah ini, Anda akan men-deploy Cloud Function Python di Cloud Shell yang mengekspos UDF validasi dan normalisasi alamat ke BigQuery.

Menulis File Sumber Cloud Function

Jalankan perintah berikut di Cloud Shell untuk membuat direktori sumber Cloud Function dan menulis main.py dan 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

Deploy Cloud Function & Mengonfigurasi Izin IAM

Jalankan perintah ini di Cloud Shell untuk men-deploy Cloud Function generasi ke-2 dan mengonfigurasi Koneksi Resource Cloud BigQuery:

# 1. Deploy 2nd-Gen Cloud Function
gcloud functions deploy validate_address_udf \
  --gen2 \
  --runtime=python311 \
  --region=${REGION} \
  --source=. \
  --entry-point=validate_address_udf \
  --trigger-http \
  --no-allow-unauthenticated \
  --memory=512Mi \
  --cpu=1 \
  --concurrency=80 \
  --quiet

# 2. Extract Function Endpoint URI
export FUNCTION_URL=$(gcloud functions describe validate_address_udf --region=${REGION} --gen2 --format="value(serviceConfig.uri)")

# 3. Create BigQuery Cloud Resource Connection
bq mk --connection --location=US --project_id=${GCP_PROJECT} --connection_type=CLOUD_RESOURCE address_val_conn || true

# 4. Extract Connection Service Account Email
export BQ_SA_EMAIL=$(bq show --format=prettyjson --connection US.address_val_conn | grep -o '"serviceAccountId": "[^"]*"' | cut -d'"' -f4)

# 5. Bind Cloud Run Invoker and Vertex AI User IAM Roles to BigQuery Connection Service Account
gcloud run services add-iam-policy-binding validate-address-udf \
  --region=${REGION} \
  --member="serviceAccount:${BQ_SA_EMAIL}" \
  --role="roles/run.invoker" --quiet

gcloud projects add-iam-policy-binding ${GCP_PROJECT} \
  --member="serviceAccount:${BQ_SA_EMAIL}" \
  --role="roles/aiplatform.user" --quiet

# 6. Wait for connection IAM policy propagation
echo "Waiting 60 seconds for BigQuery connection IAM policy to propagate..."
sleep 60

Anda akan melihat output yang menunjukkan bahwa deployment Cloud Function telah selesai dan binding IAM berhasil diterapkan.

Mendaftarkan Fungsi Normalisasi Alamat Jarak Jauh

Sekarang Anda akan mendaftarkan DDL Fungsi Jarak Jauh BigQuery (validate_address_udf) yang menghubungkan baris tabel BigQuery ke endpoint Cloud Function yang di-deploy (${FUNCTION_URL}).

Jalankan perintah berikut di Cloud Shell untuk mengambil URL Cloud Function yang di-deploy dan mendaftarkan Remote Function secara otomatis:

# 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. Memproses Awal Data Profil & Encoding Fonetik

Pada langkah ini, Anda akan menjalankan kueri pra-pemrosesan SQL BigQuery pada tabel customer_nodes yang telah di-ingest.

Menjalankan Pembersihan Data & Kueri Fitur Fonetik

Di Editor SQL BigQuery Studio, jalankan kueri di bawah untuk membuat customer_nodes_cleaned. Kueri ini:

  1. Memanggil validate_address_udf untuk mendapatkan alamat yang dinormalisasi dan hasil validasi.
  2. Membuat encoding fonetik SOUNDEX untuk given_name dan surname guna menangani variasi ejaan.
  3. Membuat kolom profile_text terstruktur.
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;

Buat kueri tabel node yang sudah dibersihkan:

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

Anda akan melihat output yang mirip dengan:

rec_id

given_name

given_name_soundex

marga

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. Membuat Embedding Profil Semantik & Penelusuran Vektor

Selain normalisasi alamat, kunci fonetik Soundex, dan jarak pengeditan Levenshtein, BigQuery mendukung fungsi embedding AI generatif bawaan melalui AI.EMBED.

Dengan AI.EMBED, BigQuery membuat embedding teks langsung di SQL menggunakan model dasar (seperti text-embedding-005) tanpa memerlukan DDL indeks vektor manual:

Jalankan kueri di bawah di Editor 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. Pemberian Skor Pasangan Kandidat & Penggabungan Fitur Edge Hybrid

Mengevaluasi semua kemungkinan pasangan data pelanggan (pertumbuhan kuadrat O(N²)) menjadi tidak mungkin dilakukan secara komputasi seiring skala set data bertambah. Di mesin database relasional seperti BigQuery, upaya untuk menerapkan pemblokiran berbasis aturan menggunakan kondisi gabungan OR yang kompleks di beberapa kolom (seperti menggabungkan a.soc_sec_id = b.soc_sec_id OR a.given_name_soundex = b.given_name_soundex OR ...) mencegah pengoptimal kueri menggunakan gabungan hash yang skalabel atau gabungan penggabungan pengurutan pada satu kunci gabungan yang sama. Sebagai gantinya, mesin kembali ke cross join O(N²) dan memfilter setiap pasangan, yang gagal dalam skala besar.

Pada langkah ini, Anda akan mengambil pasangan kandidat yang dihasilkan oleh tabel penelusuran vektor kami (vector_candidate_edges) dan menggabungkannya dengan customer_nodes_cleaned melalui gabungan yang sama dan diindeks dengan cepat (ON c.source_id = a.rec_id dan ON c.target_id = b.rec_id). Kemudian, Anda akan menghitung skor kecocokan berbobot yang menggabungkan:

  • Skor Kecocokan SSN (bobot: 0.30)
  • Kemiripan Edit Nama Keluarga menggunakan jarak Levenshtein EDIT_DISTANCE (bobot: 0.20)
  • Kesamaan Edit Nama Depan (bobot: 0.20)
  • Skor Kecocokan TGL LHR (bobot: 0.15)
  • Kesamaan Jaccard Token Alamat (bobot: 0.15) selama SPLIT(LOWER(formatted_address), ' ')

Menghitung Tepi Kandidat & Skor Kesamaan Berbobot

Jalankan kueri berikut di Editor SQL BigQuery Studio untuk mengisi 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;

Periksa kecocokan tepi kandidat:

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

Anda akan melihat output yang mirip dengan:

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

Menggabungkan Edge Berbasis Aturan dan Penelusuran Vektor ke dalam Tabel Terpadu

Gabungkan tepi kandidat dari pencocokan fuzzy berbasis aturan dan penelusuran vektor semantik ke dalam satu tabel final_matched_edges yang sudah dihapus duplikatnya:

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. Konstruksi Grafik Properti & Traversal Jalur GQL ISO

BigQuery mendukung ISO GQL (Graph Query Language) secara native melalui Grafik Properti. Grafik Properti membuat tampilan grafik logis di atas tabel BigQuery relasional tanpa duplikasi data.

Pada langkah ini, Anda akan membuat grafik properti customer_identity_graph menggunakan tabel edge kandidat terpadu (final_matched_edges) dan kueri koneksi grafik di seluruh profil pelanggan menggunakan traversal jalur k-hop {1, 2}.

Konektivitas Transitif & Traversals K-Hop

Membuat DDL Grafik Properti BigQuery

Jalankan pernyataan DDL berikut di Editor SQL 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
);

Memvisualisasikan Cluster Grafik K-Hop

Jalankan kueri di bawah untuk memvisualisasikan kecocokan kelompok pelanggan di seluruh 1 hingga 2 lompatan hubungan:

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;

Visualisasi Kelompok Grafik K-Hop

Menyelesaikan Kelompok Pelanggan Kanonis

Jalankan kueri berikut untuk menyelesaikan pengelompokan entitas ke dalam 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;

Buat kueri tabel cluster yang diselesaikan:

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

Anda akan melihat output yang mirip dengan:

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

Membuat Tampilan Metrik Evaluasi

Untuk menghitung Presisi, Perolehan, dan skor F1 terhadap tabel ground_truth_links, jalankan:

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;

Buat kueri tampilan metrik evaluasi:

SELECT * FROM `identity_resolution.evaluation_metrics`;

Anda akan melihat output yang mirip dengan:

total_ground_truth

total_predictions

true_positives

false_positives

false_negatives

presisi

ingatan

f1_score

6538

5620

5608

13

930

0.9977

0.8579

0.9225

Menyelesaikan Kelompok Rumah Tangga melalui Pembobotan Grafik Adamic-Adar

Meskipun penyelesaian identitas individu menyelesaikan data milik orang yang sama, arsitektur Customer 360 perusahaan sering kali memerlukan pengelompokan Entitas Rumah Tangga tingkat yang lebih tinggi untuk individu yang tinggal di alamat yang sama.

Karena tolok ukur sintetis (seperti FEBRL3) mengevaluasi kebenaran nyata tingkat individu, resolusi rumah tangga dilakukan sebagai langkah hilir. Jika tidak ada histori perpindahan yang diberi stempel waktu, individu yang ditautkan ke beberapa alamat dapat menyebabkan penggabungan berlebih atau fragmentasi cluster. Untuk mengatasinya, kami menggunakan Pembobotan Grafik Adamic-Adar untuk membuat keanggotaan keluarga yang tidak tetap.

Jalankan kueri di bawah di Editor SQL BigQuery Studio untuk mengisi household_clusters menggunakan pembobotan grafik 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;

Buat kueri tabel cluster rumah tangga yang telah diselesaikan:

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;

Memvisualisasikan Hierarki Identitas End-to-End melalui GQL

Untuk melacak secara visual hierarki identitas 3 tingkat yang lengkap—menghubungkan Pelanggan Mentah yang tidak dikelompokkan ke Entitas Pelanggan yang telah diselesaikan, dan selanjutnya ke Entitas Rumah Tangga yang telah diselesaikan—jalankan kueri DDL dan ISO GQL berikut di 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;

Menjalankan kueri GQL ini di BigQuery Studio akan merender kanvas visualisasi grafik interaktif 3 tingkat yang menampilkan rekaman profil pelanggan mentah (RawCustomer) yang diselesaikan ke entitas kanonis individual (ResolvedCustomer), yang ditautkan ke entitas rumah tangga multi-penghuni bersama (ResolvedHousehold).

Visualisasi Cluster Grafik Rumah Tangga Multi-Penghuni

9. Resolusi Bertahap & Stabilitas Persisten

Dalam aplikasi perusahaan di dunia nyata, data pelanggan baru tiba secara berkelanjutan melalui penyerapan batch harian atau real-time. Daripada menjalankan ulang resolusi grafik lengkap di seluruh set data historis, Mesin Pencocokan Delta Inkremental membandingkan rekaman masuk baru dengan cluster dasar yang telah diselesaikan (resolved_customers).

Untuk mencapainya secara efisien, mesin ini menggunakan Vector Search (

VECTOR_SEARCH

) sebagai bentuk pengelompokan dinamis. Dengan memperlakukan setiap rekaman masuk sebagai titik kueri, VECTOR_SEARCH mengambil set tetangga terdekat K teratas dari indeks embedding dasar historis. Jika data masuk cocok dengan profil pelanggan lama yang ada di atas nilai minimum kesamaan, data tersebut akan digabungkan secara dinamis ke dalam cluster tersebut dan mewarisi canonical_customer_id dasar pengukuran (MATCHED_TO_EXISTING_CLUSTER). Jika tidak ada tetangga terdekat dasar pengukuran yang ditemukan di atas nilai minimum, UUID entity baru akan dibuat (NEW_CUSTOMER_ENTITY).

Menyerap Contoh Data Tambahan Batch

Tempel dan jalankan DDL berikut di Editor SQL BigQuery Studio untuk membuat 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
  )
]);

Jalankan Kueri Pencocokan Delta Tambahan

Jalankan kueri berikut di Editor SQL BigQuery untuk melakukan pencocokan delta terhadap set data dasar yang telah diselesaikan:

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;

Membuat kueri hasil resolusi inkremental:

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

Anda akan melihat output yang mirip dengan:

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. Penggabungan Pengelompokan & Stabilitas Cluster (Tumpang-Tindih 1-ε)

Dalam sistem perusahaan produksi, tim biasanya mengelompokkan ulang seluruh grafik secara berulang (misalnya, mingguan atau bulanan) untuk menggabungkan sumber data dan tepi baru. Saat hubungan baru terbentuk, pengelompokan ulang grafik penuh dapat menyebabkan ID cluster bergeser atau berubah secara acak di seluruh eksekusi pipeline.

Untuk mempertahankan ID pelanggan yang persisten untuk sistem CRM, CDP, dan penagihan hilir, Stabilitas Cluster mengevaluasi tumpang-tindih node antara cluster Run saat ini (t) dan cluster Run sebelumnya (t-1) menggunakan nilai minimum tumpang-tindih (1 - ε) (dengan ε = 0,30, yang memerlukan tumpang-tindih node minimum 70%).

Jika cluster yang baru dihitung dalam Run (t) memiliki setidaknya 70% data anggotanya yang sama dengan cluster dari Run (t-1), cluster tersebut akan mewarisi ID pelanggan persisten historis (STABLE_EVOLUTION). Cluster baru akan menerima UUID yang baru dibuat (NEW_CLUSTER_CREATED).

Menjalankan Kueri Tumpang-Tindih & Stabilitas Cluster

Jalankan kueri berikut di Editor SQL BigQuery Studio untuk mengisi 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;

Melihat Pratinjau yang Difilter untuk Kumpulan Data Inkremental

Jalankan kueri ini untuk memverifikasi status stabilitas cluster untuk data penyerapan harian Anda:

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;

Anda akan melihat output yang mirip dengan:

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

Untuk menghindari biaya berkelanjutan pada akun Google Cloud Anda, bersihkan resource yang di-deploy dan set data BigQuery.

Di Cloud Shell, jalankan:

# 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

Jika Anda membuat project Google Cloud khusus untuk lab ini, Anda dapat menghapus project tersebut:

gcloud projects delete ${GCP_PROJECT}

12. Selamat

Selamat! Anda telah berhasil membangun mesin Customer Identity Resolution end-to-end di dalam Google Cloud BigQuery menggunakan BigQuery Property Graph, kueri ISO GQL, pencocokan kemiripan hibrida, pencocokan delta inkremental, dan jaminan stabilitas cluster persisten.

Yang telah Anda pelajari

  • Cara men-deploy Cloud Function generasi ke-2 dan mengeksposnya sebagai Fungsi Jarak Jauh BigQuery.
  • Cara memproses awal demografi pelanggan menggunakan SOUNDEX encoding fonetik dan validasi alamat.
  • Cara menjalankan pemblokiran kandidat dan menghitung skor kemiripan hybrid menggunakan jarak Levenshtein (EDIT_DISTANCE) dan kemiripan Jaccard token.
  • Cara membuat Grafik Properti BigQuery (CREATE PROPERTY GRAPH) melalui tabel node dan edge.
  • Cara membuat kueri jalur grafik menggunakan GQL (GRAPH_TABLE) ISO dengan pengukur k-hop {1, 2}.
  • Cara menyelesaikan pengelompokan pelanggan kanonis dan mengevaluasi performa model terhadap metrik kebenaran nyata.
  • Cara menjalankan pencocokan delta inkremental untuk penyerapan batch harian tanpa pemrosesan ulang set data lengkap.
  • Cara menerapkan jaminan nilai minimum tumpang-tindih (1 - ε) untuk mempertahankan stabilitas cluster persisten di seluruh operasi pipeline.

Langkah berikutnya

Dokumen referensi