פתרון בעיות שקשורות לזהות לקוחות באמצעות BigQuery Graph

1. מבוא

ב-Codelab הזה תיצרו מנוע מודולרי מקצה לקצה של זיהוי זהויות לקוחות (התאמת ישויות) ישירות ב-Google Cloud BigQuery. תשלבו בין Google Cloud Shell לפריסת התשתית לבין כלי ה-SQL של BigQuery Studio לניקוי נתונים, לדירוג מועמדים, ליצירת גרף מאפיינים ולמעברים בנתיבים של GQL (Graph Query Language) בתקן ISO.

זיהוי זהויות הוא יכולת בסיסית לפתרון Customer 360 לארגונים, לזיהוי הונאות ולאיחוד נתונים ממערכות שונות. יש הרבה גישות תקפות לפתרון בעיות שקשורות לזהויות, בהתאם לרמת הבשלות של הנתונים ולצרכים העסקיים. לכן, כל השלבים בסדנת הקוד הזו הם מודולריים ואופציונליים. הצינור נועד להציג מגוון טכניקות נפוצות ברמת ייצור בתעשייה – כולל נורמליזציה של כתובות UDF מרוחקות, חסימה פונטית של Soundex, חיפוש וקטורי סמנטי (AI.EMBED), ניקוד תכונות היברידי וקיבוץ גרפים של מאפיינים ב-GQL – כדי שתוכלו לבחור את הדפוסים שמתאימים לארכיטקטורה שלכם.

צריך להתאים את שיטות ההתאמה ואת ספי ההתאמה בהתאם למידת הרצון של הארגון להשתמש בהתאמה דטרמיניסטית לעומת הסתברותית, שנקבעת על ידי תרחיש השימוש המטורגט. לדוגמה, פעולות שקשורות לתאימות, לחיוב או למימון בדרך כלל מסתמכות על כללים דטרמיניסטיים ברמת דיוק גבוהה (כמו התאמות מדויקות של מספר ביטוח לאומי או מספר מס) כדי למנוע קישור שגוי, בעוד שהתאמה אישית של שיווק, ניתוח נתונים ומנועי המלצות מסתמכים לעיתים קרובות על התאמה הסתברותית לא מדויקת ועל דמיון סמנטי של וקטורים כדי למקסם את ההיזכרות ולגלות קשרים עדינים.

ארכיטקטורה של BigQuery Customer Identity Resolution Engine

הפעולות שתבצעו:

  • הוספת מערך נתונים של FEBRL3 Benchmark: טעינת רשומות סינתטיות של לקוחות וזוגות תואמים של נתוני אמת ב-BigQuery.
  • פריסת Address Validation Remote UDF: פריסת פונקציה ב-Python Cloud ורישום פונקציה מרוחקת ב-BigQuery כדי לבצע נורמליזציה של כתובות.
  • עיבוד מוקדם של נתוני הפרופיל וקידודים פונטיים: מריצים ניקוי נתונים באמצעות SQL, מפעילים את פונקציית הכתובת המוגדרת על ידי המשתמש ומחשבים SOUNDEX מפתחות פונטיים ומרחקי עריכה של לבנשטיין:
    • Soundex Phonetic Encoding: אלגוריתם פונטי לאינדוקס שמות לפי הצליל שלהם באנגלית. הוא ממיר שמות לקוד בן 4 תווים (אות ראשונה ואחריה שלוש ספרות) שמייצג קבוצות של צלילי עיצורים (לדוגמה, גם "John" וגם "Jon" ממופים ל-J500, וגם "Smith" וגם "Smyth" ממופים ל-S530). כך הוא מספק אותות של התאמה פונטית לצורך מתן ציונים לתכונות וחסימה מצטברת של דלתא בזמן אמת.
    • מרחק לבנשטיין (EDIT_DISTANCE): מדד מחרוזת שמודד את המספר המינימלי של עריכות בתו אחד (הוספות, מחיקות או החלפות) שנדרשות כדי לשנות מחרוזת אחת למחרוזת אחרת, ומאפשר התאמה מדויקת של שמות וכתובות לא מדויקים.
  • יצירת הטמעות של פרופילים סמנטיים וחיפוש וקטורי: יצירת הטמעות של טקסט ישירות ב-SQL באמצעות AI.EMBED (text-embedding-005) ומציאת K השכנים הקרובים ביותר באמצעות VECTOR_SEARCH כדי לשמש כשכבת יצירת מועמדים תת-לינארית.
  • ניקוד של זוגות מועמדים ושילוב תכונות היברידיות של קצוות: שימוש בזוגות מועמדים של חיפוש וקטורי כדי לבטל את המורכבות של הצלבה מסדר גודל O(N²), חישוב של ציוני דמיון משוקללים של תכונות מרובות (מספר ביטוח לאומי, מרחק עריכה של Levenshtein, תאריך לידה, כתובת Jaccard) ושילוב של קצוות בטבלת מועמדים מאוחדת.
  • יצירה של גרף נכסים ומעברים בנתיבים של ISO GQL: יצירה של PROPERTY GRAPH ב-BigQuery, הפעלה של שאילתות נתיבים של ISO GQL‏ (GRAPH_TABLE) כדי לזהות אשכולות של לקוחות מקושרים, חישוב של מדדי הערכה פרטניים וביצוע אשכול רך של משקי בית באמצעות שקלול גרף Adamic-Adar.{1, 2}
  • רזולוציה מצטברת ויציבות מתמשכת: עיבוד של קליטת נתונים יומית עם התאמה מצטברת של דלתא.
  • איחוד אשכולות ויציבות אשכולות (חפיפה של 1-ε): שמירה על יציבות אשכולות לאורך זמן בכל הפעלות הצינור באמצעות ערך סף של חפיפה של (1-ε).

הדרישות

  • דפדפן אינטרנט כמו Chrome.
  • פרויקט ב-Google Cloud שהחיוב בו מופעל.

שיעור ה-Codelab הזה מיועד למהנדסי נתונים, למפתחי מסדי נתונים ולמומחי AI/ML בכל הרמות, כולל מתחילים.

משך הזמן המשוער: 45 דקות
העלות המשוערת: פחות מ-2.00 דולר ארה"ב (השימוש הוא בפונקציות Cloud Functions ובעיבוד שאילתות ב-BigQuery בתשלום לפי שימוש).

‫2. לפני שמתחילים

יצירת פרויקט ב-Google Cloud

  1. במסוף Google Cloud, בדף לבחירת הפרויקט, בוחרים פרויקט ב-Google Cloud או יוצרים פרויקט.
  2. הקפידו לוודא שהחיוב מופעל בפרויקט שלכם ב-Cloud. כך בודקים אם החיוב מופעל בפרויקט

הפעלת Cloud Shell

Cloud Shell היא סביבת שורת פקודה שפועלת ב-Google Cloud וכוללת מראש את הכלים הנדרשים.

  1. לוחצים על Activate Cloud Shell בחלק העליון של מסוף Google Cloud.
  2. מאמתים את האימות:
gcloud auth list
  1. מגדירים משתני סביבה ב-Cloud Shell:
export GCP_PROJECT=$(gcloud config get-value project)
export REGION="us-central1"
export DATASET_ID="identity_resolution"

הפעלת ממשקי ה-API הנדרשים

מריצים את הפקודה הבאה ב-Cloud Shell באמצעות חשבון המשתמש כדי להפעיל את כל שירותי Google Cloud הנדרשים:

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

כדי להבטיח הפעלה חלקה של ה-API וגישה ל-Application Default Credentials (‏ADC), צריך ליצור חשבון שירות ייעודי למעבדה ולהפעיל התחזות לחשבון: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

יצירת מערך נתונים ב-BigQuery

יוצרים את מערך הנתונים ב-BigQuery לאחסון של צמתי הלקוחות, הקשתות, מודלי הגרפים ותצוגות ההערכה:

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

הפלט אמור להיראות כך:

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

כדי להבטיח קיבולת מחשוב ייעודית לחיפושים במדד וקטורי, לצבירות גרפים ולביצוע פונקציות מרחוק, בלי להיות מוגבלים על ידי מגבלות CPU לפי דרישה או מכסות משותפות, צריך ליצור הזמנה של Enterprise Edition עם שינוי גודל אוטומטי ב-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. הטמעה של מערך נתונים של צומת לקוח ב-FEBRL3

לפני שמפעילים את הפונקציה המרוחקת לאימות כתובות ומבצעים זיהוי ישויות, צריך לטעון את מערך הנתונים הסינתטי של FEBRL3 להשוואה בין ישויות (שכולל 5,000 רשומות של לקוחות עם אשכולות של כפילויות מרובות, עד 5 כפילויות לכל לקוח) באמצעות הספרייה recordlinkage של Python, ולכתוב את צמתי הלקוחות הגולמיים (customer_nodes) ואת קישורי ההתאמה של נתוני האמת (ground_truth_links) ל-BigQuery באמצעות BigQuery DataFrames‏ (bigframes).

מריצים את הפקודות הבאות ב-Cloud Shell כדי להתקין את יחסי התלות ולהפעיל את סקריפט ההטמעה:

# 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, עוברים אל BigQuery Studio, פותחים כרטיסייה חדשה של שאילתת SQL (+) ומריצים את השאילתה שבהמשך כדי לבדוק את הטבלה של צמתי הלקוחות שהועברו:

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;

הפלט אמור להיראות כך:

rec_id

given_name

surname

street_number

address_1

address_2

פרבר

מיקוד

הסמוי הסופי

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

שימו לב איך קבוצת הנתונים של ההשוואה לשוק מציגה נתונים מלוכלכים ריאליסטיים באוספים כפולים:

  • וריאציות פונטיות ואיותיות: brent לעומת brnt / bernt,‏ wood לעומת woode / wod ו-clifton לעומת cliffton.
  • קיצורים וטעויות הקלדה בכתובת: girdlestone circuit לעומת girdelstone circut / girdlestone cir / girdlestone crt, ומספר הבית 11 לעומת שגיאת OCR 15.
  • החלפת מיקום של תווים וערכים חסרים: החלפת מיקום של תווים בתאריך הלידה (19340706 לעומת 19340760), מצבים חסרים ( ) ומספרי ביטוח לאומי חסרים ( ).

בשלבים הבאים תשתמשו בSOUNDEX קידודים פונטיים, בפונקציות מוגדרות על ידי המשתמש (UDF) לנרמול כתובות, במרחק לבנשטיין לעריכה ובAI.EMBEDחיפוש וקטורי כדי לגשר על הפערים האלה ולקשר בצורה מדויקת בין פרופילים כפולים.

4. פריסת פונקציית UDF של Address Validation משירות חיצוני

נירמול הכתובות מתבצע על ידי סטנדרטיזציה של שמות רחובות, גבולות פרברים ומיקודים לפני ביצוע ההתאמה. ‫Google Maps Address Validation API הוא שירות שמקבל כתובת, מזהה את רכיבי הכתובת ומאמת אותם. בשלב הזה תפרסו ב-Cloud Shell פונקציה ב-Cloud Functions ב-Python, שחושפת ל-BigQuery פונקציה בהגדרת המשתמש (UDF) לאימות ולנרמול של כתובות.

כתיבת קובצי מקור של Cloud Functions

מריצים את הפקודה הבאה ב-Cloud Shell כדי ליצור את ספריית קובצי המקור של Cloud Functions ולכתוב את main.py ואת 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

פריסת Cloud Function והגדרת הרשאות IAM

מריצים את הפקודות האלה ב-Cloud Shell כדי לפרוס את פונקציית Cloud מהדור השני ולהגדיר קישור למשאבים ב-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

אמורות להופיע תוצאות שמציינות שהפריסה של Cloud Functions הושלמה והרשאות ה-IAM הוחלו בהצלחה.

הרשמה של פונקציה לנרמול כתובות מרוחקות

עכשיו צריך לרשום את ה-DDL של הפונקציה המרוחקת ב-BigQuery‏ (validate_address_udf) שמקשרת בין שורות בטבלת BigQuery לבין נקודת הקצה של Cloud Functions שפרסתם (${FUNCTION_URL}).

מריצים את הפקודה הבאה ב-Cloud Shell כדי לאחזר את כתובת ה-URL של פונקציית Cloud Functions שפרסתם ולרשום את הפונקציה המרוחקת באופן אוטומטי:

# 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. עיבוד מוקדם של נתוני פרופיל וקידודים פונטיים

בשלב הזה, תריצו שאילתת SQL לעיבוד מקדים ב-BigQuery על הטבלה customer_nodes שהועברה.

הפעלת ניקוי נתונים ושאילתת תכונות פונטיות

בעורך ה-SQL של BigQuery Studio, מריצים את השאילתה שלמטה כדי ליצור את customer_nodes_cleaned. השאילתה הזו:

  1. ‫Calls validate_address_udf כדי לקבל כתובות מנורמלות ופסקי דין של אימות.
  2. יוצר SOUNDEX קידודים פונטיים ל-given_name ול-surname כדי לטפל בווריאציות איות.
  3. יוצרת שדה מובנה profile_text.
CREATE OR REPLACE TABLE `identity_resolution.customer_nodes_cleaned` AS
WITH raw_data AS (
  SELECT 
    rec_id, dataset_source,
    TRIM(LOWER(given_name)) AS given_name_clean,
    TRIM(LOWER(surname)) AS surname_clean,
    `identity_resolution.validate_address_udf`(street_number, address_1, address_2, suburb, state, postcode) AS addr_json,
    TRIM(suburb) AS suburb, TRIM(state) AS state, TRIM(postcode) AS postcode,
    TRIM(date_of_birth) AS date_of_birth, TRIM(soc_sec_id) AS soc_sec_id
  FROM `identity_resolution.customer_nodes`
)
SELECT
  rec_id, dataset_source,
  given_name_clean AS given_name,
  SOUNDEX(given_name_clean) AS given_name_soundex,
  surname_clean AS surname,
  SOUNDEX(surname_clean) AS surname_soundex,
  CONCAT(given_name_clean, ' ', surname_clean) AS full_name,
  STRING(addr_json.formatted_address) AS formatted_address,
  BOOL(addr_json.address_is_valid) AS address_is_valid,
  STRING(addr_json.validation_granularity) AS validation_granularity,
  STRING(addr_json.possible_next_action) AS possible_next_action,
  suburb, state, postcode, date_of_birth, soc_sec_id,
  CONCAT('Name: ', CONCAT(given_name_clean, ' ', surname_clean), '; Address: ', STRING(addr_json.formatted_address), '; DOB: ', date_of_birth, '; SSN: ', soc_sec_id) AS profile_text
FROM raw_data;

מריצים שאילתה על טבלת הצמתים הנקייה:

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

הפלט אמור להיראות כך:

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. יצירת הטמעות של פרופילים סמנטיים וחיפוש וקטורי

בנוסף לנורמליזציה של כתובות, למפתחות פונטיים של Soundex ולמרחק העריכה של Levenshtein,‏ BigQuery תומך בפונקציות מובנות של הטמעה של AI גנרטיבי באמצעות AI.EMBED.

באמצעות AI.EMBED, ‏ BigQuery יוצר הטבעות טקסט ישירות ב-SQL באמצעות מודלים בסיסיים (כמו text-embedding-005) בלי לדרוש DDL של אינדקס וקטורי ידני:

מריצים את השאילתות הבאות בעורך ה-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. דירוג של זוגות מועמדים ושילוב תכונות היברידיות ב-Edge

הערכה של כל זוגות הרשומות האפשריים של הלקוחות (גידול ריבועי O(N²)) הופכת לבעיה חישובית ככל שהנתונים גדלים. במנועי מסדי נתונים רלציוניים כמו BigQuery, ניסיון להטמיע חסימה מבוססת-כללים באמצעות OR תנאי איחוד מורכבים בכמה עמודות (למשל איחוד לפי a.soc_sec_id = b.soc_sec_id OR a.given_name_soundex = b.given_name_soundex OR ...) מונע ממייעל השאילתות להשתמש באיחודים מבוססי-גיבוב או באיחודים מבוססי-מיון שניתנים להרחבה במפתח איחוד יחיד. במקום זאת, המנוע חוזר ל-cross join של O(N²) ומסנן כל זוג, מה שגורם לכשל בהרחבה.

בשלב הזה, תיקחו זוגות של מועמדים שנוצרו על ידי טבלת החיפוש הווקטורי שלנו (vector_candidate_edges) ותצטרפו אליהם מול customer_nodes_cleaned באמצעות צירופים מהירים עם אינדקסים (ON c.source_id = a.rec_id ו-ON c.target_id = b.rec_id). לאחר מכן, תחשבו ציון התאמה משוקלל שמשלב:

  • ציון התאמה של מספר תעודת זהות (SSN) (משקל: 0.30)
  • Surname Edit Similarity באמצעות מרחק לבנשטיין EDIT_DISTANCE (משקל: 0.20)
  • דמיון בעריכת שם פרטי (משקל: 0.20)
  • ציון ההתאמה של תאריך הלידה (משקל: 0.15)
  • Address Token Jaccard Similarity (משקל: 0.15) over SPLIT(LOWER(formatted_address), ' ')

חישוב קצוות של מועמדים וציוני דמיון משוקללים

מריצים את השאילתה הבאה בעורך ה-SQL של BigQuery Studio כדי לאכלס את 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;

בודקים את ההתאמות האפשריות של קצוות:

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

הפלט אמור להיראות כך:

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

מיזוג של קצוות מבוססי-כללים וקצוות של חיפוש וקטורי לטבלה מאוחדת

שילוב של קצוות מועמדים מחיפוש וקטורי סמנטי והתאמה משוערת מבוססת-כללים לטבלה אחת של final_matched_edges ללא כפילויות:

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. יצירת גרף מאפיינים ומעברים בנתיבים של ISO GQL

‫BigQuery תומך ב-ISO GQL (Graph Query Language) באופן מקורי דרך גרפים של נכסים. תרשים מאפיינים יוצר תצוגת תרשים לוגית על טבלאות רלציוניות ב-BigQuery בלי לשכפל את הנתונים.

בשלב הזה, תיצרו גרף נכסים customer_identity_graph באמצעות טבלת הקצוות המאוחדת של המועמדים (final_matched_edges) וחיבורי גרף השאילתות בפרופילי הלקוחות באמצעות {1, 2} k-hop path traversal.

קישוריות טרנזיטיבית ומעברים של K-Hop

יצירת DDL של גרף מאפיינים ב-BigQuery

מריצים את הצהרת ה-DDL הבאה בעורך ה-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
);

הדמיה של אשכולות בתרשים K-Hop

מריצים את השאילתה הבאה כדי להציג באופן חזותי את אשכולות הלקוחות התואמים בטווח של 1 עד 2 קפיצות בקשר:

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

פתרון של אשכולות לקוחות קנוניים

מריצים את השאילתה הבאה כדי לפתור את הבעיה של אשכולות ישויות ב-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;

מריצים שאילתה על הטבלה של האשכולות שנפתרו:

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

הפלט אמור להיראות כך:

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

יצירת תצוגה של מדדי הערכה

כדי לחשב את מדדי הדיוק, ההחזרה ו-F1 ביחס לטבלה ground_truth_links, מריצים את הפקודה:

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;

שליחת שאילתה לתצוגת מדדי ההערכה:

SELECT * FROM `identity_resolution.evaluation_metrics`;

הפלט אמור להיראות כך:

total_ground_truth

total_predictions

true_positives

false_positives

false_negatives

דיוק

recall

f1_score

6538

5620

5608

13

930

0.9977

0.8579

0.9225

פתרון של אשכולות משקי בית באמצעות שקלול גרפים של Adamic-Adar

בעוד שפתרון זהות אישי מאפשר לפתור רשומות ששייכות לאותו אדם, ארכיטקטורות של Customer 360 בארגונים דורשות לעיתים קרובות קיבוץ ברמה גבוהה יותר של ישות משפחתית של אנשים שגרים באותה כתובת.

מכיוון שבדיקות השוואה סינתטיות (כמו FEBRL3) מעריכות את נתוני האמת ברמה האישית, ההתאמה לרמת משק הבית מתבצעת כשלב המשך. אם אין היסטוריית מיקומים עם חותמות זמן, אנשים שמקושרים לכמה כתובות עלולים לגרום למיזוג יתר או לפיצול של אשכולות. כדי לפתור את הבעיה הזו, אנחנו משתמשים בשקלול גרפים של Adamic-Adar כדי ליצור חברות רכה במשקי בית.

מריצים את השאילתה שלמטה בעורך ה-SQL של BigQuery Studio כדי לאכלס את household_clusters באמצעות שקלול גרף 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;

מריצים שאילתה על הטבלה של אשכולות משקי הבית שזוהו:

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

כדי לעקוב באופן ויזואלי אחר ההיררכיה המלאה של הזהויות ברמת 3 השכבות – מלקוחות גולמיים לא מקובצים ועד ישויות לקוח וישויות משפחה מפוענחות – מריצים את שאילתת ה-DDL וה-ISO GQL הבאה ב-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;

הרצת שאילתת ה-GQL הזו ב-BigQuery Studio מציגה בד ציור של המחשה ויזואלית אינטראקטיבית בת 3 רמות, שבו מוצגים רשומות גולמיות של פרופיל לקוח (RawCustomer) שנפתרו לישויות קנוניות נפרדות (ResolvedCustomer), שמקושרות לישויות משותפות של משקי בית עם כמה דיירים (ResolvedHousehold).

המחשה של אשכולות בתרשים של משקי בית עם כמה דיירים

9. שיפור הדרגתי של הרזולוציה ויציבות מתמשכת

ביישומים ארגוניים בעולם האמיתי, רשומות חדשות של לקוחות מגיעות באופן רציף באמצעות קליטה יומית או בזמן אמת של נתונים. במקום להריץ מחדש את הפתרון המלא של הגרף על כל מערך הנתונים ההיסטורי, מנוע התאמה מצטבר של דלתא משווה רשומות נכנסות חדשות לאשכולות בסיסיים קיימים שנפתרו (resolved_customers).

כדי לעשות את זה ביעילות, המנוע משתמש בחיפוש וקטורי (

VECTOR_SEARCH

) כסוג של אשכול דינמי. המערכת מתייחסת לכל רשומה נכנסת כנקודת שאילתה, ואז VECTOR_SEARCH מאחזרת את קבוצת השכנים הקרובים ביותר (K) ממדד ההטמעה של בסיס הנתונים ההיסטורי. אם רשומה נכנסת תואמת לפרופיל לקוח קיים מעל מידת דמיון מינימלית בין המשתמשים, היא מתמזגת באופן דינמי לאותו אשכול ומקבלת בירושה את ה-Baseline canonical_customer_id (MATCHED_TO_EXISTING_CLUSTER). אם לא נמצא שכן Baseline מעל הסף, נוצר מזהה ייחודי אוניברסלי (UUID) חדש של ישות (NEW_CUSTOMER_ENTITY).

הטמעת רשומות לדוגמה של נתונים מצטברים באצווה

מדביקים את ה-DDL הבא ומריצים אותו בעורך ה-SQL של BigQuery Studio כדי ליצור את 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
  )
]);

הפעלת שאילתה מצטברת של התאמת דלתא

מריצים את השאילתה הבאה בעורך ה-SQL של BigQuery כדי לבצע התאמה של שינויים ביחס למערך נתוני הבסיס שנוצר:

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;

שליחת שאילתה לגבי תוצאות ההגדרה המצטברת:

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

הפלט אמור להיראות כך:

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. איחוד אשכולות ויציבות האשכולות (חפיפה של 1-ε)

במערכות ארגוניות בייצור, צוותים בדרך כלל יוצרים מחדש את כל הגרף על בסיס חוזר (למשל, שבועי או חודשי) כדי לשלב קצוות ומקורות נתונים חדשים. ככל שנוצרים קשרים חדשים, שינוי מלא של אשכולות בגרף עלול לגרום למזהי האשכולות להשתנות או להתהפך באופן שרירותי במהלך הפעלות של צינורות.

כדי לשמור על מזהי לקוחות קבועים במערכות CRM,‏ CDP וחיוב בהמשך הדרך, יציבות האשכול בודקת את החפיפה בין הצמתים באשכולות של הריצה הנוכחית (t) לבין הצמתים באשכולות של הריצה הקודמת (t-1) באמצעות סף חפיפה של (1 – ε) (כאשר ε = 0.30, כלומר נדרשת חפיפה של 70% לפחות בין הצמתים).

אם אשכול חדש שחושב ב-Run ‏ (t) חולק לפחות 70% מהרשומות החברות שלו עם אשכול מ-Run ‏ (t-1), הוא מקבל בירושה את מזהה הלקוח ההיסטורי הקבוע (STABLE_EVOLUTION). אשכולות חדשים לגמרי מקבלים מזהי UUID חדשים שנוצרו (NEW_CLUSTER_CREATED).

הפעלת שאילתה של חפיפה ויציבות בין אשכולות

מריצים את השאילתה הבאה בעורך ה-SQL של BigQuery Studio כדי לאכלס את 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;

הצגת תצוגה מקדימה מסוננת של רשומות מצטברות

מריצים את השאילתה הזו כדי לוודא מהו סטטוס היציבות של האשכול עבור רשומות הקליטה היומיות:

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;

הפלט אמור להיראות כך:

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. הסרת המשאבים

כדי למנוע חיובים שוטפים בחשבון Google Cloud, צריך לנקות את המשאבים שנפרסו ואת מערך הנתונים ב-BigQuery.

ב-Cloud Shell, מריצים את הפקודה:

# 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

אם יצרתם פרויקט בענן ייעודי ב-Google Cloud לשיעור ה-Lab הזה, אתם יכולים למחוק את הפרויקט:

gcloud projects delete ${GCP_PROJECT}

12. מזל טוב

מעולה! הצלחתם לבנות מנוע מקיף לזיהוי זהויות לקוחות ב-Google Cloud BigQuery באמצעות גרף מאפיינים של BigQuery, שאילתות ISO GQL, התאמה היברידית של דמיון, התאמה מצטברת של דלתא והבטחות יציבות מתמשכות של אשכולות.

מה למדתם

  • איך פורסים פונקציה של Cloud Functions מדור שני וחושפים אותה כפונקציה מרוחקת של BigQuery.
  • איך לבצע עיבוד מקדים של נתונים דמוגרפיים של לקוחות באמצעות SOUNDEX קידודים פונטיים ואימות כתובות.
  • איך מבצעים חסימה של מועמדים וחישוב של ציוני דמיון היברידיים באמצעות מרחק לוינשטיין (EDIT_DISTANCE) ודמיון Jaccard של טוקנים.
  • איך יוצרים גרף מאפיינים של BigQuery (CREATE PROPERTY GRAPH) על טבלאות של צמתים וקשתות.
  • איך שולחים שאילתות על נתיבי גרף באמצעות ISO GQL (GRAPH_TABLE) עם {1, 2} כמתייחסים למרחק של k צעדים.
  • איך לפתור בעיות שקשורות לאשכולות לקוחות קנוניים ולהעריך את ביצועי המודל בהשוואה למדדי אמת קרקעית (ground truth).
  • איך מבצעים התאמה מצטברת של דלתא לצריכת נתונים יומיים בחבילות בלי לעבד מחדש את מערך הנתונים המלא.
  • איך להחיל ערך סף של חפיפה (1 - ε) כדי לשמור על יציבות מתמשכת של האשכול במהלך הרצת צינורות.

השלבים הבאים

מסמכים לדוגמה