1. परिचय
इस कोडलैब में, आपको Google Cloud BigQuery में सीधे तौर पर मॉड्यूलर, एंड-टू-एंड कस्टमर आइडेंटिटी रिज़ॉल्यूशन (इकाई मिलान) इंजन बनाने का तरीका बताया जाएगा. इन्फ़्रास्ट्रक्चर डिप्लॉयमेंट के लिए, Google Cloud Shell को डेटा को साफ़ करने, उम्मीदवार को स्कोर देने, प्रॉपर्टी ग्राफ़ बनाने, और आईएसओ GQL (ग्राफ़ क्वेरी लैंग्वेज) पाथ ट्रैवर्सल के लिए, BigQuery Studio SQL Editor के साथ जोड़ा जाएगा.
पहचान हल करने की सुविधा, एंटरप्राइज़ Customer 360, धोखाधड़ी का पता लगाने, और कई सिस्टम से डेटा इकट्ठा करने के लिए ज़रूरी है. डेटा मैच करने के कई मान्य तरीके हैं. ये तरीके, डेटा की मैचिंग की क्वालिटी और कारोबार की ज़रूरतों के हिसाब से तय होते हैं. इसलिए, इस कोडलैब में दिए गए सभी चरण मॉड्यूलर हैं और इन्हें करना ज़रूरी नहीं है. इस पाइपलाइन को, इंडस्ट्री में इस्तेमाल होने वाली कई सामान्य और प्रोडक्शन-ग्रेड तकनीकों को दिखाने के लिए डिज़ाइन किया गया है. इनमें रिमोट यूडीएफ़ पता सामान्यीकरण, साउंडेक्स फ़ोनेटिक ब्लॉकिंग, सिमैंटिक वेक्टर सर्च (AI.EMBED), हाइब्रिड फ़ीचर स्कोरिंग, और जीक्यूएल प्रॉपर्टी ग्राफ़ क्लस्टरिंग शामिल हैं. इससे, अपनी ज़रूरत के हिसाब से पैटर्न को अपनाया जा सकता है.
मैचिंग के तरीकों और स्कोरिंग थ्रेशोल्ड को, डिटरमिनिस्टिक बनाम प्रॉबेबिलिस्टिक मैचिंग के लिए, आपके संगठन की ज़रूरत के हिसाब से ट्यून किया जाना चाहिए. यह ज़रूरत, टारगेट किए गए इस्तेमाल के उदाहरण से तय होती है. उदाहरण के लिए, नियमों का सख्ती से पालन करने, बिलिंग या वित्तीय कार्रवाइयों के लिए, आम तौर पर सटीक नियमों (जैसे, एसएसएन या टैक्स आईडी का सटीक मिलान) का इस्तेमाल किया जाता है. इससे गलत लिंक करने से बचा जा सकता है. वहीं, मार्केटिंग को ज़्यादा निजी बनाने, आंकड़ों का विश्लेषण करने, और सुझाव देने वाले इंजन के लिए, अक्सर संभावित फ़ज़ी मैचिंग और सिमैंटिक वेक्टर समानता का इस्तेमाल किया जाता है. इससे ज़्यादा से ज़्यादा लोगों को याद रखने और बारीकी से कनेक्शन का पता लगाने में मदद मिलती है.

आपको क्या करना होगा
- FEBRL3 बेंचमार्क डेटासेट इंपोर्ट करें: सिंथेटिक ग्राहक रिकॉर्ड और ग्राउंड ट्रुथ मैच पेयर को BigQuery में लोड करें.
- Address Validation रिमोट यूडीएफ़ डिप्लॉय करना: सड़क के पतों को सामान्य बनाने के लिए, Python Cloud Function डिप्लॉय करें और BigQuery Remote Function रजिस्टर करें.
- प्रोफ़ाइल के डेटा और फ़ोनेटिक एन्कोडिंग को पहले से प्रोसेस करना: SQL डेटा क्लीनिंग को लागू करें, पते के यूडीएफ़ को लागू करें, और
SOUNDEXफ़ोनेटिक कुंजियों और लेवेंश्टाइन एडिट दूरी का हिसाब लगाएं:- Soundex फ़ोनेटिक एन्कोडिंग: यह एक फ़ोनेटिक एल्गोरिदम है. इसका इस्तेमाल, अंग्रेज़ी में बोले गए शब्दों के आधार पर नामों को इंडेक्स करने के लिए किया जाता है. यह नामों को चार वर्णों वाले कोड में बदलता है. इसमें पहला अक्षर व्यंजन की आवाज़ वाले ग्रुप को दिखाता है और इसके बाद तीन अंक होते हैं. उदाहरण के लिए,
"John"और"Jon", दोनों कोJ500के तौर पर मैप किया जाता है. वहीं,"Smith"और"Smyth"कोS530के तौर पर मैप किया जाता है. इससे, फ़ीचर स्कोरिंग और रीयल-टाइम में डेल्टा को धीरे-धीरे ब्लॉक करने के लिए, फ़ोनेटिक मैचिंग सिग्नल मिलते हैं. - लेवेंश्टाइन दूरी (
EDIT_DISTANCE): यह एक स्ट्रिंग मेट्रिक है. इससे यह पता चलता है कि एक स्ट्रिंग को दूसरी स्ट्रिंग में बदलने के लिए, कम से कम कितने वर्णों में बदलाव (जोड़ना, मिटाना या बदलना) करना होगा. इससे नाम और पते को सटीक तरीके से मैच करने में मदद मिलती है.
- Soundex फ़ोनेटिक एन्कोडिंग: यह एक फ़ोनेटिक एल्गोरिदम है. इसका इस्तेमाल, अंग्रेज़ी में बोले गए शब्दों के आधार पर नामों को इंडेक्स करने के लिए किया जाता है. यह नामों को चार वर्णों वाले कोड में बदलता है. इसमें पहला अक्षर व्यंजन की आवाज़ वाले ग्रुप को दिखाता है और इसके बाद तीन अंक होते हैं. उदाहरण के लिए,
- सिमैंटिक प्रोफ़ाइल एम्बेडिंग और वेक्टर सर्च जनरेट करना:
AI.EMBED(text-embedding-005) का इस्तेमाल करके, सीधे तौर पर SQL में टेक्स्ट एम्बेडिंग जनरेट करें. साथ ही, सब-लीनियर कैंडिडेट जनरेशन लेयर के तौर पर काम करने के लिए,VECTOR_SEARCHका इस्तेमाल करके टॉप-K सबसे नज़दीकी पड़ोसी ढूंढें. - उम्मीदवारों के जोड़े को स्कोर करना और हाइब्रिड एज फ़ीचर फ़्यूज़न: क्रॉस-जॉइन की O(N²) जटिलता को खत्म करने के लिए, वेक्टर सर्च के उम्मीदवारों के जोड़े का फ़ायदा उठाएं. साथ ही, कई फ़ीचर के वज़न के हिसाब से समानता स्कोर (एसएसएन, लेवेंश्टाइन एडिट डिस्टेंस, जन्म की तारीख, पता जैकार्ड) का हिसाब लगाएं और किनारों को एक ही उम्मीदवार टेबल में मर्ज करें.
- प्रॉपर्टी ग्राफ़ बनाना और आईएसओ जीक्यूएल पाथ ट्रैवर्सल: BigQuery
PROPERTY GRAPHबनाएं. साथ ही, कनेक्टेड ग्राहक क्लस्टर को हल करने, अलग-अलग आकलन मेट्रिक का हिसाब लगाने, और Adamic-Adar ग्राफ़ वेटिंग का इस्तेमाल करके सॉफ्ट हाउसहोल्ड क्लस्टरिंग करने के लिए, आईएसओ जीक्यूएल{1, 2}पाथ क्वेरी (GRAPH_TABLE) चलाएं. - बेहतर रिज़ॉल्यूशन और लगातार स्थिरता: हर दिन के बैच के डेटा को प्रोसेस करें. साथ ही, डेल्टा मैचिंग को बेहतर बनाएं.
- क्लस्टरिंग कंसोलिडेशन और क्लस्टर स्टैबिलिटी (1-ε ओवरलैप): (1-ε) ओवरलैप थ्रेशोल्ड गारंटी का इस्तेमाल करके, पाइपलाइन रन में क्लस्टर स्टैबिलिटी को लागू करें.
आपको किन चीज़ों की ज़रूरत होगी
- कोई वेब ब्राउज़र, जैसे कि Chrome.
- बिलिंग की सुविधा वाला Google क्लाउड प्रोजेक्ट.
यह कोडलैब, डेटा इंजीनियर, डेटाबेस डेवलपर, और एआई/एमएल के सभी लेवल के प्रैक्टिशनर के लिए बनाया गया है. इसमें शुरुआती लेवल के लोग भी शामिल हैं.
अनुमानित अवधि: 45 मिनट
अनुमानित लागत: 2.00 डॉलर से कम (इसमें इस्तेमाल के हिसाब से Cloud Functions और BigQuery क्वेरी प्रोसेसिंग के लिए शुल्क लिया जाता है).
2. शुरू करने से पहले
Google Cloud प्रोजेक्ट बनाना
- Google Cloud Console में, प्रोजेक्ट चुनने वाले पेज पर, Google Cloud प्रोजेक्ट चुनें या बनाएं.
- पक्का करें कि आपके क्लाउड प्रोजेक्ट के लिए बिलिंग की सुविधा चालू हो. किसी प्रोजेक्ट के लिए बिलिंग चालू है या नहीं, यह देखने का तरीका जानें.
Cloud Shell शुरू करना
Cloud Shell, Google Cloud में चलने वाला एक कमांड-लाइन एनवायरमेंट है. इसमें ज़रूरी टूल पहले से लोड होते हैं.
- Google Cloud कंसोल में सबसे ऊपर मौजूद, Cloud Shell चालू करें पर क्लिक करें.
- पुष्टि की पुष्टि करें:
gcloud auth list
- Cloud Shell में एनवायरमेंट वैरिएबल कॉन्फ़िगर करें:
export GCP_PROJECT=$(gcloud config get-value project)
export REGION="us-central1"
export DATASET_ID="identity_resolution"
ज़रूरी एपीआई चालू करना
ज़रूरी सभी Google Cloud सेवाएं चालू करने के लिए, अपने उपयोगकर्ता खाते का इस्तेमाल करके Cloud Shell में यह कमांड चलाएं:
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
सेवा खाता और उपयोगकर्ता के तौर पर काम करने की सुविधा सेट अप करना (सुझाया गया)
एपीआई को आसानी से लागू करने और ऐप्लिकेशन के डिफ़ॉल्ट क्रेडेंशियल (एडीसी) ऐक्सेस करने के लिए, लैब का एक सेवा खाता बनाएं और 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.
BigQuery रिज़र्वेशन और असाइनमेंट बनाएं (ज़रूरी नहीं / सुझाव दिया गया)
वेक्टर इंडेक्स सर्च, ग्राफ़ एग्रीगेशन, और रिमोट फ़ंक्शन को बिना किसी रुकावट के एक्ज़ीक्यूट करने के लिए, Cloud Shell में Enterprise Edition के लिए आरक्षण बनाएं. इससे, आपको ऑन-डिमांड सीपीयू की सीमाओं या शेयर किए गए कोटे की वजह से कोई समस्या नहीं आएगी. साथ ही, ऑटोस्केलिंग की सुविधा भी मिलेगी:
# 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 के ग्राहक नोड डेटासेट को इनजेस्ट करना
पते की पुष्टि करने वाले रिमोट फ़ंक्शन को डिप्लॉय करने और पहचान हल करने की प्रोसेस को पूरा करने से पहले, आपको Python की recordlinkage लाइब्रेरी का इस्तेमाल करके, सिंथेटिक FEBRL3 इकाई के समाधान का बेंचमार्क डेटासेट लोड करना होगा. इसमें 5,000 ग्राहक रिकॉर्ड होते हैं. हर ग्राहक के लिए, डुप्लीकेट क्लस्टर की संख्या पांच तक होती है. इसके बाद, आपको BigQuery DataFrames (bigframes) का इस्तेमाल करके, ग्राहकों के रॉ नोड (customer_nodes) और ग्राउंड ट्रुथ मैच लिंक (ground_truth_links) को BigQuery में लिखना होगा.
डेटा ट्रांसफ़र करने की स्क्रिप्ट को चलाने और ज़रूरी सॉफ़्टवेयर इंस्टॉल करने के लिए, 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 | सरनेम | street_number | address_1 | address_2 | उपनगर | पिन कोड | राज्य | date_of_birth | soc_sec_id | dataset_source |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| |
|
|
|
|
|
|
|
|
|
|
| |
|
|
|
|
|
|
|
|
|
| ||
|
|
|
|
|
|
|
|
|
|
|
|
ध्यान दें कि बेंचमार्क डेटासेट, डुप्लीकेट क्लस्टर में किस तरह का अशुद्ध डेटा शामिल करता है:
- उच्चारण और स्पेलिंग में अंतर:
brentबनामbrnt/bernt,woodबनामwoode/wod, औरcliftonबनामcliffton. - पते में इस्तेमाल किए गए छोटे शब्दों और टाइपिंग की गड़बड़ियों को ठीक करना:
girdlestone circuitबनामgirdelstone circut/girdlestone cir/girdlestone crtऔर घर का नंबर11बनाम ओसीआर की गड़बड़ी15. - वर्णों की जगह बदलना और वैल्यू मौजूद न होना: जन्म की तारीख में वर्णों की जगह बदलना (
19340706बनाम19340760), राज्यों की जानकारी मौजूद न होना ( ), और सामाजिक सुरक्षा आईडी मौजूद न होना ( ).
आने वाले चरणों में, इन अंतरों को कम करने और डुप्लीकेट प्रोफ़ाइलों को सटीक तरीके से लिंक करने के लिए, SOUNDEX फ़ोनेटिक एन्कोडिंग, पते को सामान्य बनाने वाले यूडीएफ़, लेवेंश्टाइन एडिट डिस्टेंस, और AI.EMBED वेक्टर सर्च का इस्तेमाल किया जाएगा.
4. Address Validation Remote Function UDF को डिप्लॉय करना
पते को सामान्य बनाने की प्रोसेस में, मैचिंग करने से पहले सड़कों के नाम, उपनगरीय सीमाओं, और पिन कोड को स्टैंडर्ड बनाया जाता है. Google Maps Address Validation API एक ऐसी सेवा है जो पते को स्वीकार करती है, पते के कॉम्पोनेंट की पहचान करती है, और उनकी पुष्टि करती है. इस चरण में, Cloud Shell में एक Python Cloud फ़ंक्शन डिप्लॉय किया जाएगा. यह BigQuery को पते की पुष्टि करने और उसे सामान्य बनाने वाला यूडीएफ़ उपलब्ध कराता है.
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 Functions को डिप्लॉय करना और IAM अनुमतियां कॉन्फ़िगर करना
दूसरी जनरेशन के Cloud फ़ंक्शन को डिप्लॉय करने और BigQuery Cloud Resource Connection को कॉन्फ़िगर करने के लिए, Cloud Shell में ये कमांड चलाएं:
# 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 बाइंडिंग लागू हो गई हैं.
रिमोट पते को सामान्य बनाने वाले फ़ंक्शन को रजिस्टर करना
अब आपको BigQuery रिमोट फ़ंक्शन DDL (validate_address_udf) रजिस्टर करना होगा. यह BigQuery टेबल की लाइनों को आपके डिप्लॉय किए गए Cloud Functions एंडपॉइंट (${FUNCTION_URL}) से कनेक्ट करता है.
डिप्लॉय किए गए Cloud Function के यूआरएल को वापस पाने और रिमोट फ़ंक्शन को अपने-आप रजिस्टर करने के लिए, Cloud Shell में यह कमांड चलाएं:
# 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. प्रोफ़ाइल डेटा और फ़ोनेटिक एन्कोडिंग को पहले से प्रोसेस करना
इस चरण में, आपको अपनी customer_nodes टेबल में शामिल किए गए डेटा पर, BigQuery एसक्यूएल प्रीप्रोसेसिंग क्वेरी को लागू करना होगा.
डेटा को साफ़ करना और फ़ोनेटिक फ़ीचर क्वेरी को लागू करना
customer_nodes_cleaned बनाने के लिए, BigQuery Studio के SQL एडिटर में नीचे दी गई क्वेरी चलाएं. इस क्वेरी में:
- सामान्य किए गए पते और पुष्टि के नतीजों को पाने के लिए,
validate_address_udfको कॉल करता है. - यह कुकी, स्पेलिंग में होने वाले बदलावों को मैनेज करने के लिए,
given_nameऔरsurnameके लिएSOUNDEXफ़ोनेटिक एन्कोडिंग जनरेट करती है. - यह स्ट्रक्चर्ड
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_soundex | formatted_address |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
6. सिमैंटिक प्रोफ़ाइल एम्बेडिंग और वेक्टर सर्च जनरेट करना
BigQuery, पते को सामान्य बनाने, साउंडेक्स फ़ोनेटिक कुंजियों, और लेवेंश्टाइन एडिट डिस्टेंस के साथ-साथ, AI.EMBED के ज़रिए जनरेटिव एआई के पहले से मौजूद एम्बेडिंग फ़ंक्शन का इस्तेमाल करने की सुविधा देता है.
AI.EMBED का इस्तेमाल करके, BigQuery सीधे तौर पर एसक्यूएल में टेक्स्ट एम्बेडिंग जनरेट करता है. इसके लिए, फ़ाउंडेशन मॉडल (जैसे कि text-embedding-005) का इस्तेमाल किया जाता है. इसमें, मैन्युअल वेक्टर इंडेक्स DDL की ज़रूरत नहीं होती:
प्रोफ़ाइल एम्बेडिंग जनरेट करना और टॉप-के वेक्टर सर्च को लागू करना
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. उम्मीदवारों के जोड़े को स्कोर करना और हाइब्रिड एज फ़ीचर फ़्यूज़न
डेटासेट का साइज़ बढ़ने पर, सभी संभावित ग्राहक रिकॉर्ड पेयर (O(N²) क्वाड्रेटिक ग्रोथ) का आकलन करना मुश्किल हो जाता है. BigQuery जैसे रिलेशनल डेटाबेस इंजन में, कई कॉलम में जटिल OR जॉइन की शर्तों का इस्तेमाल करके, नियम के आधार पर ब्लॉक करने की सुविधा लागू करने की कोशिश की जाती है. इससे क्वेरी ऑप्टिमाइज़र, एक ही इक्वी-जॉइन कुंजी पर स्केलेबल हैश जॉइन या सॉर्ट-मर्ज जॉइन का इस्तेमाल नहीं कर पाता.a.soc_sec_id = b.soc_sec_id OR a.given_name_soundex = b.given_name_soundex OR ... इसके बजाय, इंजन O(N²) क्रॉस जॉइन का इस्तेमाल करता है और हर पेयर को फ़िल्टर करता है. हालांकि, यह बड़े डेटासेट के लिए सही नहीं है.
इस चरण में, आपको हमारी वेक्टर सर्च टेबल (vector_candidate_edges) से जनरेट किए गए उम्मीदवार जोड़े लेने होंगे. इसके बाद, आपको उन्हें customer_nodes_cleaned के साथ जोड़ना होगा. इसके लिए, इंडेक्स किए गए इक्वी-जॉइन (ON c.source_id = a.rec_id और ON c.target_id = b.rec_id) का इस्तेमाल करें. इसके बाद, आपको इन दोनों को मिलाकर, वज़न के हिसाब से मैच स्कोर का हिसाब लगाना होगा:
- एसएसएन मैच स्कोर (वज़न:
0.30) - लेवेंश्टाइन दूरी
EDIT_DISTANCEका इस्तेमाल करके, सरनेम में बदलाव करने पर समानता (वज़न:0.20) - दिए गए नाम में बदलाव करने पर समानता (वज़न:
0.20) - जन्म की तारीख के मैच होने का स्कोर (वज़न:
0.15) - पते के टोकन की जैकार्ड सिमिलैरिटी (वज़न:
0.15)SPLIT(LOWER(formatted_address), ' ')से ज़्यादा है
उम्मीदवार के किनारों और वज़न के हिसाब से समानता के स्कोर का हिसाब लगाना
matched_edges में डेटा भरने के लिए, BigQuery Studio के एसक्यूएल एडिटर में यह क्वेरी चलाएं:
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 |
|
|
|
|
|
|
|
|
|
नियमों पर आधारित और वेक्टर सर्च की सुविधाओं को एक ही टेबल में मर्ज करना
नियमों पर आधारित फ़ज़ी मैचिंग और सेमैंटिक वेक्टर सर्च से मिले उम्मीदवार के एज को एक ही, डुप्लीकेट हटाए गए 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. प्रॉपर्टी ग्राफ़ बनाना और आईएसओ GQL पाथ ट्रैवर्सल
BigQuery, प्रॉपर्टी ग्राफ़ के ज़रिए आईएसओ जीक्यूएल (ग्राफ़ क्वेरी लैंग्वेज) के साथ नेटिव तौर पर काम करता है. प्रॉपर्टी ग्राफ़, डेटा को डुप्लीकेट किए बिना, रिलेशनल BigQuery टेबल पर लॉजिकल ग्राफ़ व्यू बनाता है.
इस चरण में, आपको अपनी यूनीफ़ाइड कैंडिडेट एज टेबल (final_matched_edges) का इस्तेमाल करके, एक प्रॉपर्टी ग्राफ़ customer_identity_graph बनाना होगा. साथ ही, {1, 2} k-हॉप पाथ ट्रैवर्सल का इस्तेमाल करके, ग्राहक प्रोफ़ाइलों के बीच क्वेरी ग्राफ़ कनेक्शन बनाने होंगे.

BigQuery प्रॉपर्टी ग्राफ़ DDL बनाना
BigQuery Studio के SQL एडिटर में, यह DDL स्टेटमेंट लागू करें:
CREATE OR REPLACE PROPERTY GRAPH `identity_resolution.customer_identity_graph`
NODE TABLES (
`identity_resolution.customer_nodes_cleaned` AS `Customer`
KEY (rec_id)
)
EDGE TABLES (
`identity_resolution.final_matched_edges`
KEY (source_id, target_id)
SOURCE KEY (source_id) REFERENCES `Customer`(rec_id)
DESTINATION KEY (target_id) REFERENCES `Customer`(rec_id)
LABEL MATCHED_TO
);
K-हॉप ग्राफ़ क्लस्टर को विज़ुअलाइज़ करना
एक से दो रिलेशनशिप हॉप में, मैच करने वाले ग्राहक क्लस्टर को विज़ुअलाइज़ करने के लिए, यहां दी गई क्वेरी चलाएं:
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;

कैननिकल कस्टमर क्लस्टर की समस्या हल करना
इकाई क्लस्टर को 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 |
|
|
|
|
|
|
|
|
|
इवैलुएशन मेट्रिक व्यू बनाना
ground_truth_links टेबल के हिसाब से, सटीक अनुमान, रिकॉल, और F1-स्कोर का हिसाब लगाने के लिए, यह कमांड चलाएं:
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 | प्रीसिज़न | रीकॉल | f1_score |
|
|
|
|
|
|
|
|
एडमिक-अदार ग्राफ़ वेटिंग की मदद से, घर के सदस्यों के क्लस्टर की पहचान करना
पहचान की पुष्टि करने की सुविधा, एक ही व्यक्ति से जुड़े रिकॉर्ड को हल करती है. हालांकि, एंटरप्राइज़ Customer 360 आर्किटेक्चर को अक्सर परिवार के सदस्य के हिसाब से ग्रुप बनाने की ज़रूरत होती है. इसमें एक ही पते पर रहने वाले लोगों को शामिल किया जाता है.
सिंथेटिक बेंचमार्क (जैसे, FEBRL3) में, व्यक्ति के लेवल पर ग्राउंड ट्रुथ का आकलन किया जाता है. इसलिए, घर के हिसाब से रिज़ॉल्यूशन को बाद में किया जाता है. जगह बदलने के टाइमस्टैंप वाले इतिहास के न होने पर, एक से ज़्यादा पतों से जुड़े लोग, एक से ज़्यादा कारोबारों के एक साथ मर्ज होने या क्लस्टर के फ़्रैगमेंटेशन की वजह बन सकते हैं. इस समस्या को हल करने के लिए, हम Adamic-Adar Graph Weighting का इस्तेमाल करते हैं, ताकि घर के सदस्यों की सदस्यताएं बनाई जा सकें.
household_clusters को भरने के लिए, BigQuery Studio के एसक्यूएल एडिटर में नीचे दी गई क्वेरी चलाएं. इसके लिए, 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 की मदद से, एंड-टू-एंड आइडेंटिटी के पदानुक्रम को विज़ुअलाइज़ करना
तीन टियर वाली आइडेंटिटी के पूरे क्रम को विज़ुअली ट्रेस करने के लिए, BigQuery Studio में यहां दी गई DDL और आईएसओ GQL क्वेरी को एक्ज़ीक्यूट करें:
-- 1. Create Household Node Table
CREATE OR REPLACE TABLE `identity_resolution.household_nodes` AS
SELECT DISTINCT
canonical_household_id,
household_address,
total_residents
FROM `identity_resolution.household_clusters`;
-- 2. Create Unresolved Record to Resolved Entity Edge Table
CREATE OR REPLACE TABLE `identity_resolution.customer_entity_edges` AS
SELECT DISTINCT
rec_id,
canonical_customer_id
FROM `identity_resolution.resolved_customers`,
UNNEST(customer_records) AS rec_id;
-- 3. Create Primary Household Edge Table (Highest Weighted Household Rank = 1)
CREATE OR REPLACE TABLE `identity_resolution.primary_household_edges` AS
SELECT
canonical_customer_id,
canonical_household_id,
household_membership_weight,
household_rank
FROM `identity_resolution.household_clusters`
WHERE household_rank = 1;
-- 4. Update Unified Property Graph DDL
CREATE OR REPLACE PROPERTY GRAPH `identity_resolution.customer_identity_graph`
NODE TABLES (
`identity_resolution.customer_nodes_cleaned` AS `RawCustomer`
KEY (rec_id),
`identity_resolution.resolved_customers` AS `ResolvedCustomer`
KEY (canonical_customer_id),
`identity_resolution.household_nodes` AS `ResolvedHousehold`
KEY (canonical_household_id)
)
EDGE TABLES (
`identity_resolution.final_matched_edges`
KEY (source_id, target_id)
SOURCE KEY (source_id) REFERENCES `RawCustomer`(rec_id)
DESTINATION KEY (target_id) REFERENCES `RawCustomer`(rec_id)
LABEL MATCHED_TO,
`identity_resolution.customer_entity_edges`
KEY (rec_id, canonical_customer_id)
SOURCE KEY (rec_id) REFERENCES `RawCustomer`(rec_id)
DESTINATION KEY (canonical_customer_id) REFERENCES `ResolvedCustomer`(canonical_customer_id)
LABEL RESOLVED_TO,
`identity_resolution.primary_household_edges`
KEY (canonical_customer_id, canonical_household_id)
SOURCE KEY (canonical_customer_id) REFERENCES `ResolvedCustomer`(canonical_customer_id)
DESTINATION KEY (canonical_household_id) REFERENCES `ResolvedHousehold`(canonical_household_id)
LABEL BELONGS_TO_HOUSEHOLD
);
-- 5. Execute 3-Tier GQL Query for Multi-Resident Household Visualization
GRAPH `identity_resolution.customer_identity_graph`
MATCH p = (raw:RawCustomer)-[e1:RESOLVED_TO]->(c:ResolvedCustomer)-[e2:BELONGS_TO_HOUSEHOLD]->(h:ResolvedHousehold)
WHERE h.total_residents > 1
RETURN TO_JSON(p) AS multi_resident_household_hierarchy_path
LIMIT 15;
BigQuery Studio में इस GQL क्वेरी को चलाने पर, तीन लेयर वाला इंटरैक्टिव ग्राफ़ विज़ुअलाइज़ेशन कैनवस रेंडर होता है. इसमें, ग्राहक की प्रोफ़ाइल के रॉ रिकॉर्ड (RawCustomer) दिखाए जाते हैं. इन्हें अलग-अलग कैननिकल इकाइयों (ResolvedCustomer) के तौर पर दिखाया जाता है. ये इकाइयां, एक से ज़्यादा लोगों के साथ शेयर की गई घरेलू इकाइयों (ResolvedHousehold) से जुड़ी होती हैं.

9. इंक्रीमेंटल रिज़ॉल्यूशन और लगातार स्थिरता
असल दुनिया में एंटरप्राइज़ ऐप्लिकेशन के लिए, नए ग्राहक रिकॉर्ड हर दिन या रीयल-टाइम में बैच के तौर पर लगातार मिलते रहते हैं. इंक्रीमेंटल डेल्टा मैचिंग इंजन, पूरे पुराने डेटासेट पर पूरे ग्राफ़ के रिज़ॉल्यूशन को फिर से चलाने के बजाय, नई आने वाली रिकॉर्ड की तुलना मौजूदा रिज़ॉल्व किए गए बेसलाइन क्लस्टर (resolved_customers) से करता है.
इस काम को बेहतर तरीके से करने के लिए, इंजन वेक्टर सर्च (
VECTOR_SEARCH
) को डाइनैमिक क्लस्टरिंग के तौर पर इस्तेमाल किया जाता है. हर नए रिकॉर्ड को क्वेरी पॉइंट के तौर पर इस्तेमाल करके, VECTOR_SEARCH, हिस्टोरिकल बेसलाइन एम्बेडिंग इंडेक्स से टॉप-K सबसे मिलते-जुलते पड़ोसी रिकॉर्ड का सेट वापस पाता है. अगर कोई नया रिकॉर्ड, समानता थ्रेशोल्ड से ऊपर मौजूद किसी ग्राहक प्रोफ़ाइल से मेल खाता है, तो वह डाइनैमिक तरीके से उस क्लस्टर में मर्ज हो जाता है. साथ ही, उसे बेसलाइन canonical_customer_id (MATCHED_TO_EXISTING_CLUSTER) मिलती है. अगर थ्रेशोल्ड से ऊपर कोई बेसलाइन नज़दीकी पड़ोसी नहीं मिलता है, तो एक नया इकाई यूयूआईडी बनाया जाता है (NEW_CUSTOMER_ENTITY).
इंक्रीमेंटल बैच इंटेक के सैंपल रिकॉर्ड डालना
incremental_daily_intake बनाने के लिए, BigQuery Studio के एसक्यूएल एडिटर में यह डीडीएल चिपकाएं और इसे लागू करें:
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
)
]);
इंक्रीमेंटल डेल्टा मैच क्वेरी को लागू करना
अपने हल किए गए बेसलाइन डेटासेट के साथ डेल्टा मैचिंग करने के लिए, BigQuery SQL एडिटर में यह क्वेरी चलाएं:
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 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
10. क्लस्टरिंग कंसोलिडेशन और क्लस्टर स्टैबिलिटी (1-ε ओवरलैप)
आम तौर पर, प्रोडक्शन एंटरप्राइज़ सिस्टम में टीमें, नए एज और डेटा सोर्स को शामिल करने के लिए, पूरे ग्राफ़ को बार-बार (जैसे, हर हफ़्ते या हर महीने) फिर से क्लस्टर करती हैं. नए संबंध बनने पर, पूरे ग्राफ़ को फिर से क्लस्टर करने से, क्लस्टर आइडेंटिफ़ायर में बदलाव हो सकता है या पाइपलाइन के एक्ज़ीक्यूशन के दौरान वे मनमाने तरीके से फ़्लिप हो सकते हैं.
डाउनस्ट्रीम सीआरएम, सीडीपी, और बिलिंग सिस्टम के लिए, ग्राहक आईडी को बनाए रखने के लिए, क्लस्टर स्टेबिलिटी, मौजूदा रन (t) क्लस्टर और पिछले रन (t-1) क्लस्टर के बीच नोड ओवरलैप का आकलन करता है. इसके लिए, (1 - ε) ओवरलैप थ्रेशोल्ड का इस्तेमाल किया जाता है. यहां ε = 0.30 है, जिसके लिए कम से कम 70% नोड ओवरलैप की ज़रूरत होती है.
अगर रन (t) में कंप्यूट किया गया नया क्लस्टर, रन (t-1) के क्लस्टर के साथ कम से कम 70% सदस्य रिकॉर्ड शेयर करता है, तो उसे पुराने परसिस्टेंट ग्राहक आईडी (STABLE_EVOLUTION) मिलते हैं. बिलकुल नए क्लस्टर को नए जनरेट किए गए यूयूआईडी (NEW_CLUSTER_CREATED) मिलते हैं.
क्लस्टर ओवरलैप और स्थिरता क्वेरी को लागू करना
stable_resolved_customers में डेटा भरने के लिए, BigQuery Studio के एसक्यूएल एडिटर में यह क्वेरी चलाएं:
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 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
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 प्रोजेक्ट बनाया है, तो उसे मिटाया जा सकता है:
gcloud projects delete ${GCP_PROJECT}
12. बधाई हो
बधाई हो! आपने BigQuery प्रॉपर्टी ग्राफ़, आईएसओ GQL क्वेरी, हाइब्रिड सिमिलैरिटी मैचिंग, इंक्रीमेंटल डेल्टा मैचिंग, और लगातार क्लस्टर स्टेबल रहने की गारंटी का इस्तेमाल करके, Google Cloud BigQuery में एंड-टू-एंड कस्टमर आइडेंटिटी रिज़ॉल्यूशन इंजन बनाया है.
आपने क्या सीखा
- दूसरी जनरेशन के Cloud Functions को डिप्लॉय करने और उन्हें BigQuery रिमोट फ़ंक्शन के तौर पर इस्तेमाल करने का तरीका.
SOUNDEXफ़ोनेटिक एन्कोडिंग और पते की पुष्टि करने की सुविधा का इस्तेमाल करके, ग्राहक की जनसांख्यिकी की जानकारी को पहले से प्रोसेस करने का तरीका.- उम्मीदवार को ब्लॉक करने और लेवेंश्टाइन दूरी (
EDIT_DISTANCE) और टोकन जैकार्ड समानता का इस्तेमाल करके, हाइब्रिड समानता स्कोर का हिसाब लगाने का तरीका. - नोड और एज टेबल पर BigQuery प्रॉपर्टी ग्राफ़ (
CREATE PROPERTY GRAPH) कैसे बनाया जाता है. {1, 2}k-हॉप क्वांटीफ़ायर के साथ ISO GQL (GRAPH_TABLE) का इस्तेमाल करके, ग्राफ़ पाथ के बारे में क्वेरी कैसे करें.- कैननिकल कस्टमर क्लस्टर की समस्या हल करने और सटीक मेट्रिक के हिसाब से मॉडल की परफ़ॉर्मेंस का आकलन करने का तरीका.
- पूरे डेटासेट को फिर से प्रोसेस किए बिना, हर दिन के बैच के लिए इंक्रीमेंटल डेल्टा मैचिंग को कैसे लागू करें.
- पाइपलाइन रन के दौरान, क्लस्टर की स्थिरता बनाए रखने के लिए, (1 - ε) ओवरलैप थ्रेशोल्ड गारंटी कैसे लागू करें.
अगले चरण
- BigQuery प्रॉपर्टी ग्राफ़ के दस्तावेज़ देखें.
- सेमैंटिक कैंडिडेट जनरेट करने के लिए, BigQuery Vector Search और Vertex AI टेक्स्ट एम्बेडिंग का इस्तेमाल करें.
- स्केल किए जा सकने वाले बाहरी एपीआई इंटिग्रेशन के लिए, BigQuery के रिमोट फ़ंक्शन के बारे में पढ़ें.