BigQuery Graph ile Müşteri Kimliği Çözümleme

1. Giriş

Bu codelab'de, doğrudan Google Cloud BigQuery'de modüler ve uçtan uca bir Müşteri Kimliği Çözümü (Varlık Eşleştirme) motoru oluşturacaksınız. Altyapı dağıtımı için Google Cloud Shell'i, veri temizleme, aday puanlama, özellik grafiği oluşturma ve ISO GQL (Graph Query Language) yol geçişleri için BigQuery Studio SQL Düzenleyici ile birlikte kullanacaksınız.

Kimlik çözümü; kurumsal Customer 360, sahtekarlık tespiti ve çok sistemli veri birleştirme için temel bir özelliktir. Veri olgunluğuna ve işletme ihtiyaçlarına bağlı olarak kimlik çözümlemeye yönelik birçok geçerli yaklaşım olduğundan bu codelab'deki tüm adımlar modüler ve isteğe bağlıdır. Bu işlem hattı, mimarinize uygun kalıpları seçerek kullanabilmeniz için uzaktan UDF adres normalleştirme, Soundex fonetik engelleme, semantik vektör araması (AI.EMBED), karma özellik puanlaması ve GQL özellik grafiği kümeleme gibi çeşitli yaygın ve üretime uygun endüstri tekniklerini sergilemek üzere tasarlanmıştır.

Eşleştirme yöntemleri ve puanlama eşikleri, hedef kullanım alanına göre belirlenen kuruluşunuzun belirleyici ve olasılıksal eşleştirme tercihine göre ayarlanmalıdır. Örneğin, katı uyumluluk, faturalandırma veya finansal işlemler genellikle yanlış bağlantıyı önlemek için yüksek hassasiyetli deterministik kuralları (ör. tam SSN veya Vergi Numarası eşleşmeleri) tercih ederken pazarlama kişiselleştirme, analiz ve öneri motorları, geri çağırmayı en üst düzeye çıkarmak ve ince bağlantıları ortaya çıkarmak için genellikle olasılıksal bulanık eşleşme ve anlamsal vektör benzerliğine yönelir.

BigQuery Müşteri Kimliği Çözümleme Motoru Mimarisi

Yapacaklarınız

  • FEBRL3 Benchmark Veri Kümesini Alma: Sentetik müşteri kayıtlarını ve kesin referans eşleşme çiftlerini BigQuery'ye yükleyin.
  • Adres Doğrulama Uzak UDF'sini dağıtma: Bir Python Cloud Function dağıtın ve sokak adreslerini normalleştirmek için bir BigQuery Uzak İşlevi kaydedin.
  • Profil Verilerini ve Fonetik Kodlamaları Ön İşleme: SQL veri temizleme işlemini yürütün, adres UDF'sini çağırın ve SOUNDEX fonetik anahtarlarını ve Levenshtein düzenleme mesafelerini hesaplayın:
    • Soundex Fonetik Kodlama: İngilizce'de telaffuz edildiği şekliyle adları sese göre indekslemek için kullanılan bir fonetik algoritma. Bu işlev, adları 4 karakterlik bir koda (bir ilk harf ve ardından üç rakam) dönüştürerek ünsüz ses gruplarını temsil eder (ör.hem "John" hem de "Jon", J500 ile eşlenirken "Smith" ve "Smyth", S530 ile eşlenir). Böylece, özellik puanlama ve anlık artımlı delta engelleme için fonetik eşleşme sinyalleri sağlanır.
    • Levenshtein mesafesi (EDIT_DISTANCE): Bir dizeyi diğerine dönüştürmek için gereken tek karakterlik düzenlemelerin (ekleme, silme veya değiştirme) minimum sayısını ölçen bir dize metriğidir. Bu metrik, ad ve adreslerin yaklaşık olarak eşleştirilmesini sağlar.
  • Semantik profil yerleştirmeleri ve vektör araması oluşturma: AI.EMBED (text-embedding-005) kullanarak doğrudan SQL'de metin yerleştirmeleri oluşturun ve alt doğrusal aday oluşturma katmanı olarak hizmet vermek üzere VECTOR_SEARCH kullanarak en yakın K komşuyu bulun.
  • Aday Çifti Puanlama ve Hibrit Kenar Özelliği Birleştirme: O(N²) çapraz birleştirme karmaşıklığını ortadan kaldırmak, çok özellikli ağırlıklı benzerlik puanlarını (SSN, Levenshtein düzenleme mesafesi, doğum tarihi, adres Jaccard) hesaplamak ve kenarları birleşik bir aday tablosunda birleştirmek için vektör arama aday çiftlerinden yararlanın.
  • Özellik Grafiği Oluşturma ve ISO GQL Yol Geçişleri: Bir BigQuery PROPERTY GRAPH oluşturun, bağlı müşteri kümelerini çözmek, bireysel değerlendirme metriklerini hesaplamak ve Adamic-Adar grafiği ağırlıklandırmasını kullanarak yumuşak hane halkı kümeleme gerçekleştirmek için ISO GQL {1, 2} yol sorguları (GRAPH_TABLE) çalıştırın.
  • Artımlı Çözüm ve Sürekli Kararlılık: Artımlı delta eşleştirme ile günlük toplu alımları işleyin.
  • Kümeleme Birleştirme ve Küme Kararlılığı (1-ε Çakışma): (1-ε) çakışma eşiği garantisi kullanarak ardışık düzen çalıştırmaları genelinde kalıcı küme kararlılığını zorunlu kılın.

İhtiyacınız olanlar

  • Chrome gibi bir web tarayıcısı
  • Faturalandırmanın etkin olduğu bir Google Cloud projesi.

Bu codelab, yeni başlayanlar da dahil olmak üzere her seviyeden veri mühendisleri, veritabanı geliştiricileri ve yapay zeka/makine öğrenimi uzmanları için tasarlanmıştır.

Tahmini Süre: 45 dakika
Tahmini Maliyet: 2,00 ABD dolarından az (kullandıkça öde Cloud Functions ve BigQuery sorgu işleme kullanılır).

2. Başlamadan önce

Google Cloud projesi oluşturma

  1. Google Cloud Console'daki proje seçici sayfasında bir Google Cloud projesi seçin veya oluşturun.
  2. Cloud projeniz için faturalandırmanın etkinleştirildiğinden emin olun. Bir projede faturalandırmanın etkin olup olmadığını kontrol etmeyi öğrenin.

Cloud Shell'i Başlatma

Cloud Shell, Google Cloud'da çalışan ve gerekli araçların önceden yüklendiği bir komut satırı ortamıdır.

  1. Google Cloud Console'un üst kısmında Activate Cloud Shell'i (Cloud Shell'i Etkinleştir) tıklayın.
  2. Kimlik doğrulamayı onaylayın:
gcloud auth list
  1. Cloud Shell'de ortam değişkenlerini yapılandırma:
export GCP_PROJECT=$(gcloud config get-value project)
export REGION="us-central1"
export DATASET_ID="identity_resolution"

Gerekli API'leri etkinleştirme

Gerekli tüm Google Cloud hizmetlerini etkinleştirmek için Cloud Shell'de kullanıcı hesabınızı kullanarak aşağıdaki komutu çalıştırın:

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

API'nin sorunsuz bir şekilde yürütülmesini ve Uygulama Varsayılan Kimlik Bilgileri'ne (ADC) erişimi sağlamak için özel bir laboratuvar hizmet hesabı oluşturun ve gcloud kimliğe bürünme özelliğini etkinleştirin:

# 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

BigQuery veri kümesi oluşturma

Müşteri düğümlerinizi, kenarlarınızı, grafik modellerinizi ve değerlendirme görünümlerinizi depolamak için BigQuery veri kümesini oluşturun:

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

Şuna benzer bir çıkış alırsınız:

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

İsteğe bağlı CPU sınırları veya paylaşılan kotalarla kısıtlanmadan vektör dizini aramaları, grafik toplamaları ve uzak işlev yürütme için özel bilgi işlem kapasitesi sağlamak üzere Cloud Shell'de otomatik ölçeklendirme ile Enterprise Edition rezervasyonu oluşturun:

# 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. FEBRL3 Müşteri Düğümü Veri Kümesini Kullanma

Adres doğrulama uzak işlevini dağıtmadan ve kimlik çözümleme işlemini gerçekleştirmeden önce, Python'un recordlinkage kitaplığını kullanarak sentetik FEBRL3 kimlik çözümleme karşılaştırma veri kümesini (müşteri başına en fazla 5 kopya olmak üzere çoklu kopya kümeleri içeren 5.000 müşteri kaydı içerir) yükleyecek ve BigQuery DataFrames'i (bigframes) kullanarak ham müşteri düğümlerini (customer_nodes) ve temel gerçek eşleşme bağlantılarını (ground_truth_links) BigQuery'ye yazacaksınız.

Bağımlılıkları yüklemek ve alım komut dosyasını yürütmek için Cloud Shell'de aşağıdaki komutları çalıştırın:

# 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

Google Cloud Console'da BigQuery Studio'ya gidin, yeni bir SQL sorgusu sekmesi (+) açın ve alınan müşteri düğümleri tablosunu incelemek için aşağıdaki sorguyu çalıştırın:

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;

Şuna benzer bir çıkış alırsınız:

rec_id

given_name

surname

street_number

address_1

address_2

banliyö

posta kodu

durum

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

Karşılaştırma veri kümesinin, yinelenen kümelerde gerçekçi kirli verileri nasıl sunduğuna dikkat edin:

  • Fonetik ve yazım farklılıkları: brent ve brnt / bernt, wood ve woode / wod, clifton ve cliffton.
  • Adres kısaltmaları ve yazım hataları: girdlestone circuit ve girdelstone circut / girdlestone cir / girdlestone crt, bina numarası 11 ve OCR hatası 15.
  • Karakter yer değiştirmeleri ve eksik değerler: Doğum tarihi yer değiştirmeleri (19340706 ve 19340760), eksik eyaletler ( ) ve eksik vatandaşlık numaraları ( ).

İlerleyen adımlarda, bu tutarsızlıkları gidermek ve yinelenen profilleri doğru şekilde bağlamak için SOUNDEX fonetik kodlamalar, adres normalleştirme UDF'leri, Levenshtein düzenleme mesafesi ve AI.EMBED vektör araması kullanacaksınız.

4. Address Validation Uzak İşlev UDF'sini dağıtma

Adres normalleştirme, eşleştirme yapmadan önce sokak adlarını, banliyö sınırlarını ve posta kodlarını standartlaştırır. Google Haritalar Adres Doğrulama API'si, bir adresi kabul eden, adres bileşenlerini tanımlayan ve bunları doğrulayan bir hizmettir. Bu adımda, Cloud Shell'de bir Python Cloud Functions işlevi dağıtacak ve bu işlev, BigQuery'ye bir adres doğrulama ve normalleştirme UDF'si sunacak.

Cloud Functions Kaynak Dosyalarını Yazma

Cloud Shell'de aşağıdaki komutu çalıştırarak Cloud Functions kaynak dizinini oluşturun ve main.py ile requirements.txt yazın:

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

Cloud Functions işlevini dağıtma ve IAM izinlerini yapılandırma

2. nesil Cloud Function'ı dağıtmak ve BigQuery Cloud Resource Connection'ı yapılandırmak için bu komutları Cloud Shell'de çalıştırın:

# 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

Cloud Functions dağıtımının tamamlandığını ve IAM bağlamalarının başarıyla uygulandığını belirten bir çıkış görmeniz gerekir.

Uzak Adres Normalleştirme İşlevini Kaydetme

Şimdi BigQuery tablo satırlarını dağıtılan Cloud Functions uç noktanıza (${FUNCTION_URL}) bağlayan BigQuery Remote Function DDL'yi (validate_address_udf) kaydedeceksiniz.

Dağıtılan Cloud Function URL'nizi almak ve uzak işlevi otomatik olarak kaydetmek için Cloud Shell'de aşağıdaki komutu çalıştırın:

# 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. Profil Verilerini ve Fonetik Kodlamaları Önceden İşleme

Bu adımda, alınan customer_nodes tablonuzda bir BigQuery SQL ön işleme sorgusu yürüteceksiniz.

Veri Temizleme ve Fonetik Özellik Sorgusu Yürütme

BigQuery Studio SQL Düzenleyici'de, customer_nodes_cleaned oluşturmak için aşağıdaki sorguyu çalıştırın. Bu sorgu:

  1. Normalleştirilmiş adresler ve doğrulama kararları almak için yapılan validate_address_udf çağrıları.
  2. Yazım varyasyonlarını işlemek için SOUNDEX ve surname için given_name fonetik kodlamalar oluşturur.
  3. Yapılandırılmış bir profile_text alanı oluşturur.
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;

Temizlenmiş düğümler tablosunu sorgulayın:

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

Şuna benzer bir çıkış alırsınız:

rec_id

given_name

given_name_soundex

surname

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. Semantik profil yerleştirmeleri ve Vector Search oluşturma

BigQuery, adres normalleştirme, Soundex fonetik anahtarları ve Levenshtein düzenleme mesafesine ek olarak AI.EMBED aracılığıyla yerleşik üretken yapay zeka yerleştirme işlevlerini destekler.

AI.EMBED kullanarak BigQuery, temel modelleri (ör. text-embedding-005) kullanarak doğrudan SQL'de metin yerleştirmeleri oluşturur. Bunun için manuel vektör dizini DDL'si gerekmez:

Aşağıdaki sorguları BigQuery Studio SQL Düzenleyici'de çalıştırın:

-- 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. Aday Çift Puanlama ve Hibrit Edge Özellik Birleştirme

Olası tüm müşteri kaydı çiftlerinin değerlendirilmesi (O(N²) kuadratik büyüme), veri kümesi ölçeği büyüdükçe hesaplama açısından yasaklayıcı hale gelir. BigQuery gibi ilişkisel veritabanı motorlarında, birden fazla sütunda karmaşık OR birleştirme koşulları (ör. a.soc_sec_id = b.soc_sec_id OR a.given_name_soundex = b.given_name_soundex OR ... üzerinde birleştirme) kullanarak kural tabanlı engellemeyi uygulamaya çalışmak, sorgu optimize edicinin tek bir eşitlik birleştirme anahtarında ölçeklenebilir karma birleştirmeleri veya sıralama-birleştirme birleştirmelerini kullanmasını engeller. Bunun yerine, motor O(N²) çapraz birleşime geri döner ve her çifti filtreler. Bu da ölçeklendirme sırasında başarısız olur.

Bu adımda, vektör arama tablomuz (vector_candidate_edges) tarafından oluşturulan aday çiftleri alıp hızlı, dizine eklenmiş eşitlik birleştirmeleri (ON c.source_id = a.rec_id ve ON c.target_id = b.rec_id) aracılığıyla customer_nodes_cleaned ile birleştireceksiniz. Ardından, aşağıdakileri birleştiren ağırlıklı bir eşleşme puanı hesaplayacaksınız:

  • SSN Eşleşme Puanı (ağırlık: 0.30)
  • Levenshtein mesafesi kullanılarak Soyadı Düzenleme Benzerliği EDIT_DISTANCE (ağırlık: 0.20)
  • Ad Düzenleme Benzerliği (ağırlık: 0.20)
  • DOB Match Score (weight: 0.15)
  • SPLIT(LOWER(formatted_address), ' ') üzerinde Address Token Jaccard Similarity (Adres Jetonu Jaccard Benzerliği) (ağırlık: 0.15)

Aday Kenarlarını ve Ağırlıklı Benzerlik Puanlarını Hesaplama

matched_edges değerini doldurmak için BigQuery Studio SQL Düzenleyici'de aşağıdaki sorguyu çalıştırın:

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;

Aday kenar eşleşmelerini inceleyin:

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

Şuna benzer bir çıkış alırsınız:

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

Kural tabanlı ve vektör arama kenarlarını birleştirilmiş tabloda birleştirme

Kural tabanlı yaklaşık eşleme ve semantik vektör aramasından elde edilen aday kenarlarını tek bir tekilleştirilmiş final_matched_edges tablosunda birleştirin:

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. Özellik grafiği oluşturma ve ISO GQL yolu geçişleri

BigQuery, Özellik Grafikleri aracılığıyla ISO GQL'yi (Graph Query Language) yerel olarak destekler. Mülk grafiği, veri tekilleştirme olmadan ilişkisel BigQuery tabloları üzerinde mantıksal bir grafik görünümü oluşturur.

Bu adımda, birleştirilmiş aday kenar tablonuzu (final_matched_edges) kullanarak bir özellik grafiği customer_identity_graph oluşturacak ve {1, 2} k-hop yol geçişini kullanarak müşteri profilleri arasında sorgu grafiği bağlantıları oluşturacaksınız.

Geçişli Bağlantı ve K-Hop Geçişleri

BigQuery Property Graph DDL'si oluşturma

BigQuery Studio SQL Düzenleyici'de aşağıdaki DDL ifadesini yürütün:

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
);

K-Hop Grafik Kümelerini Görselleştirme

1-2 ilişki atlaması arasında eşleşen müşteri kümelerini görselleştirmek için aşağıdaki sorguyu çalıştırın:

GRAPH `identity_resolution.customer_identity_graph`
MATCH p = (c1:Customer)-[e:MATCHED_TO]->{1, 2}(c2:Customer)
RETURN TO_JSON(p) AS graph_cluster_path
LIMIT 10;

K-Hop Grafik Kümeleri Görselleştirme

Canonical Müşteri Kümelerini Çözme

Öğe kümelerini resolved_customers olarak çözmek için aşağıdaki sorguyu çalıştırın:

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;

Çözümlenmiş kümeler tablosunu sorgulayın:

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

Şuna benzer bir çıkış alırsınız:

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

Değerlendirme metrikleri görünümü oluşturma

Hassasiyet, Geri Çağırma ve F1 puanını ground_truth_links tablosuna göre hesaplamak için şunu çalıştırın:

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;

Değerlendirme metrikleri görünümünü sorgulayın:

SELECT * FROM `identity_resolution.evaluation_metrics`;

Şuna benzer bir çıkış alırsınız:

total_ground_truth

total_predictions

true_positives

false_positives

false_negatives

precision

hatırlanabilirlik

f1_score

6538

5620

5608

13

930

0.9977

0.8579

0.9225

Adamic-Adar grafik ağırlıklandırması aracılığıyla hane grubu kümelerini çözme

Bireysel kimlik çözümü, aynı kişiye ait kayıtları çözerken kurumsal Customer 360 mimarileri genellikle aynı adresi paylaşan, birlikte yaşayan kişileri gruplandıran daha üst düzey bir hane halkı varlığı gerektirir.

FEBRL3 gibi sentetik karşılaştırma testleri, bireysel düzeyde kesin referansı değerlendirdiğinden hane halkı çözümü sonraki bir adım olarak gerçekleştirilir. Zaman damgalı yeniden konumlandırma geçmişi olmadığında, birden fazla adrese bağlı kişiler aşırı birleştirme veya küme parçalanmasına neden olabilir. Bu sorunu çözmek için Adamic-Adar grafiği ağırlıklandırması kullanarak esnek hane halkı üyelikleri oluşturuyoruz.

household_clusters değerini Adamic-Adar grafiği ağırlıklandırması kullanarak doldurmak için BigQuery Studio SQL Düzenleyici'de aşağıdaki sorguyu çalıştırın:

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;

Çözümlenmiş hane kümeleri tablosunu sorgulayın:

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;

GQL aracılığıyla uçtan uca kimlik hiyerarşisini görselleştirme

Kümelenmemiş Ham Müşteriler'i çözümlenmiş Müşteri Varlıkları'na ve oradan da çözümlenmiş Hane Halkı Varlıkları'na bağlayarak 3 katmanlı kimlik hiyerarşisinin tamamını görsel olarak izlemek için BigQuery Studio'da aşağıdaki DDL ve ISO GQL sorgusunu yürütün:

-- 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;

Bu GQL sorgusunu BigQuery Studio'da çalıştırmak, ham müşteri profili kayıtlarını (RawCustomer) bireysel kanonik varlıklara (ResolvedCustomer) çözümlenmiş şekilde gösteren 3 katmanlı etkileşimli bir grafik görselleştirme tuvali oluşturur. Bu kayıtlar, paylaşılan çok konutlu hane varlıklarına (ResolvedHousehold) bağlanır.

Birden Fazla Sakinin Bulunduğu Hane Grafiği Kümeleri Görselleştirmesi

9. Artımlı Çözüm ve Sürekli Kararlılık

Gerçek dünyadaki kurumsal uygulamalarda, yeni müşteri kayıtları günlük veya gerçek zamanlı toplu alımlar yoluyla sürekli olarak gelir. Artımlı Delta Matching Engine, tüm geçmiş veri kümesinde tam grafik çözünürlüğünü yeniden çalıştırmak yerine yeni gelen kayıtları mevcut çözümlenmiş temel kümelere (resolved_customers) göre karşılaştırır.

Motor, bunu verimli bir şekilde yapmak için Vector Search'ü (

VECTOR_SEARCH

) dinamik kümeleme biçimi olarak. Her gelen kaydı sorgu noktası olarak değerlendiren VECTOR_SEARCH, geçmişteki temel yerleştirme dizininden en yakın K komşusunun kümesini alır. Gelen bir kayıt, benzerlik eşiğinin üzerindeki mevcut bir müşteri profiliyle eşleşirse dinamik olarak bu kümeyle birleştirilir ve temel canonical_customer_id (MATCHED_TO_EXISTING_CLUSTER) değerini devralır. Eşiğin üzerinde en yakın temel komşu bulunamazsa yeni bir öğe UUID'si oluşturulur (NEW_CUSTOMER_ENTITY).

Örnek Artımlı Toplu Alma Kayıtlarını Alma

incremental_daily_intake oluşturmak için aşağıdaki DDL'yi BigQuery Studio SQL Düzenleyici'ye yapıştırıp çalıştırın:

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
  )
]);

Artımlı Delta Eşleme Sorgusu Yürütme

Çözümlenmiş temel veri kümenize karşı delta eşleştirme işlemi gerçekleştirmek için BigQuery SQL Düzenleyici'de aşağıdaki sorguyu çalıştırın:

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;

Artımlı çözünürlük sonuçlarını sorgulayın:

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

Şuna benzer bir çıkış alırsınız:

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. Kümeleme Konsolidasyonu ve Küme Kararlılığı (1-ε Çakışma)

Üretimdeki kurumsal sistemlerde ekipler, yeni kenarları ve veri kaynaklarını dahil etmek için genellikle tüm grafiği düzenli olarak (ör. haftalık veya aylık) yeniden kümelendirir. Yeni ilişkiler oluştuğunda, tam grafik yeniden kümeleme işlemi, küme tanımlayıcılarının işlem hattı yürütmeleri arasında rastgele kaymasına veya ters dönmesine neden olabilir.

Aşağı akış CRM, CDP ve faturalandırma sistemleri için kalıcı müşteri kimliklerini korumak amacıyla Küme Kararlılığı, (1 - ε) çakışma eşiğini kullanarak mevcut Çalıştırma (t) kümeleri ile önceki Çalıştırma (t-1) kümeleri arasındaki düğüm çakışmasını değerlendirir (burada ε = 0,30 olup minimum% 70 düğüm çakışması gerekir).

Run'da (t) yeni hesaplanan bir küme, Run'daki (t-1) bir küme ile üye kayıtlarının en az% 70'ini paylaşıyorsa geçmişteki kalıcı müşteri kimliğini (STABLE_EVOLUTION) devralır. Yepyeni kümeler, yeni oluşturulan UUID'leri (NEW_CLUSTER_CREATED) alır.

Küme Çakışması ve Kararlılık Sorgusu Yürütme

stable_resolved_customers değerini doldurmak için BigQuery Studio SQL Düzenleyici'de aşağıdaki sorguyu çalıştırın:

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;

Artımlı Kayıtlar İçin Filtrelenmiş Önizlemeyi Görüntüleme

Günlük alım kayıtlarınız için küme kararlılığı durumunu doğrulamak üzere bu sorguyu yürütün:

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;

Şuna benzer bir çıkış alırsınız:

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

Google Cloud hesabınızın sürekli olarak ücretlendirilmesini önlemek için dağıtılan kaynakları ve BigQuery veri kümesini temizleyin.

Cloud Shell'de şunu çalıştırın:

# 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

Bu laboratuvar için özel bir Google Cloud projesi oluşturduysanız projeyi silebilirsiniz:

gcloud projects delete ${GCP_PROJECT}

12. Tebrikler

Tebrikler! BigQuery Property Graph, ISO GQL sorguları, hibrit benzerlik eşleştirme, artımlı delta eşleştirme ve kalıcı küme kararlılığı garantilerini kullanarak Google Cloud BigQuery'de uçtan uca bir Müşteri Kimliği Çözümleme motoru oluşturmayı başardınız.

Öğrendikleriniz

  • 2. nesil Cloud Functions işlevini dağıtma ve BigQuery Remote Function olarak kullanıma sunma
  • SOUNDEX Fonetik kodlamalar ve adres doğrulamayı kullanarak müşteri demografik bilgilerine nasıl ön işlem uygulanır?
  • Levenshtein mesafesini (EDIT_DISTANCE) ve Jaccard belirteç benzerliğini kullanarak aday engelleme ve karma benzerlik puanlarını hesaplama
  • Düğüm ve kenar tabloları üzerinde BigQuery Özellik Grafiği (CREATE PROPERTY GRAPH) oluşturma
  • {1, 2} k-hop niceleyicileriyle ISO GQL (GRAPH_TABLE) kullanarak grafik yollarını sorgulama
  • Kanonik müşteri kümeleri nasıl çözülür ve model performansı, kesin referans metriklerine göre nasıl değerlendirilir?
  • Tam veri kümesi yeniden işlenmeden günlük toplu alımlar için artımlı delta eşleştirme nasıl yürütülür?
  • Ardışık düzen çalıştırmaları genelinde kalıcı küme kararlılığını korumak için (1 - ε) çakışma eşiği garantisi nasıl uygulanır?

Sonraki adımlar

Referans belgeleri