حلّ هوية العميل باستخدام BigQuery Graph

1. مقدمة

في هذا الدرس التطبيقي حول الترميز، ستنشئ محرّكًا معياريًا متكاملاً لتحديد هوية العملاء (مطابقة الكيانات) مباشرةً داخل Google Cloud BigQuery. ستجمع بين Google Cloud Shell لنشر البنية الأساسية وأداة تعديل SQL في BigQuery Studio لتنظيف البيانات وتسجيل المرشّحين وإنشاء الرسم البياني للخصائص وعمليات اجتياز مسار GQL (لغة طلبات البحث في الرسوم البيانية) وفقًا لمعيار ISO.

تُعدّ عملية تحديد الهوية إحدى الإمكانات الأساسية لخدمة Customer 360 للمؤسسات، ورصد الاحتيال، ودمج البيانات من أنظمة متعددة. بما أنّ هناك العديد من الطرق الصالحة لتحديد الهوية استنادًا إلى مستوى اكتمال البيانات واحتياجات النشاط التجاري، فإنّ جميع الخطوات في هذا الدرس التطبيقي حول الترميز مقسّمة إلى وحدات واختيارية. تم تصميم خطوة المعالجة لعرض مجموعة متنوعة من التقنيات الشائعة والمناسبة للاستخدام في بيئة الإنتاج، بما في ذلك تسوية عناوين UDF البعيدة، والحظر الصوتي باستخدام Soundex، والبحث عن المتجهات الدلالية (AI.EMBED)، وتسجيل الميزات المختلطة، وتجميع الرسومات البيانية للخصائص في GQL، ما يتيح لك اعتماد الأنماط التي تناسب تصميمك بشكل انتقائي.

يجب ضبط طرق المطابقة وحدود التسجيل استنادًا إلى مدى رغبة مؤسستك في استخدام المطابقة المحدّدة مقابل الاحتمالية، وهو ما تحدّده حالة الاستخدام المستهدَفة. على سبيل المثال، تتطلّب عمليات الامتثال أو الفوترة أو العمليات المالية عادةً قواعد حتمية عالية الدقة (مثل مطابقة رقم التأمين الاجتماعي أو رقم التعريف الضريبي) لمنع الربط الخاطئ، بينما تعتمد عمليات تخصيص التسويق والإحصاءات ومحركات الاقتراحات غالبًا على المطابقة التقريبية الاحتمالية وتشابه المتجهات الدلالية لتحقيق أقصى قدر من الاسترجاع والكشف عن الروابط الدقيقة.

بنية BigQuery Customer Identity Resolution Engine

الإجراءات التي ستنفذّها

  • استيعاب مجموعة بيانات FEBRL3 Benchmark: حمِّل سجلّات العملاء الاصطناعية وأزواج المطابقة الصحيحة إلى BigQuery.
  • نشر دالة معرَّفة من قِبل المستخدم عن بُعد لخدمة "التحقّق من صحة العناوين": يمكنك نشر دالة Python Cloud Function وتسجيل دالة BigQuery عن بُعد لتسوية عناوين الشوارع.
  • المعالجة المُسبقة لبيانات الملفات الشخصية والترميز الصوتي: نفِّذ عملية تنظيف بيانات SQL، واستدعِ دالة المستخدم المحدّدة (UDF) الخاصة بالعنوان، واحتسِب مفاتيح SOUNDEX الصوتية ومسافات التعديل في Levenshtein:
    • ترميز Soundex الصوتي: هي خوارزمية صوتية لفهرسة الأسماء حسب الصوت كما يُنطق باللغة الإنجليزية. يحوّل هذا النظام الأسماء إلى رمز مكوّن من 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 {1, 2} (GRAPH_TABLE) لحلّ مجموعات العملاء المرتبطة، واحتساب مقاييس التقييم الفردية، وإجراء تجميع مرن للعائلات باستخدام ترجيح المخطط البياني Adamic-Adar.
  • التحسين التدريجي للدقة والاستقرار المستمر: يمكنك معالجة عمليات استيعاب الدفعات اليومية من خلال مطابقة دلتا التدريجية.
  • تجميع التجميع وثبات المجموعة (تداخل 1-ε): فرض ثبات المجموعة المستمر على مستوى عمليات تشغيل خطوط الأنابيب باستخدام ضمان حدّ التداخل (1-ε).

المتطلبات

  • متصفّح ويب، مثل Chrome
  • مشروع Google Cloud تم تفعيل الفوترة فيه

تم تصميم هذا الدرس التطبيقي حول الترميز لمهندسي البيانات ومطوّري قواعد البيانات وممارسي الذكاء الاصطناعي/تعلُّم الآلة من جميع المستويات، بما في ذلك المبتدئين.

المدة المقدَّرة: 45 دقيقة
التكلفة المقدَّرة: أقل من 2.00 دولار أمريكي (تستخدِم وظائف Cloud Functions ومعالجة طلبات البحث في BigQuery بنظام الدفع حسب الاستخدام).

2. قبل البدء

إنشاء مشروع على Google Cloud

  1. في Google Cloud Console، في صفحة اختيار المشروع، اختَر مشروعًا على السحابة الإلكترونية أو أنشِئ مشروعًا على السحابة الإلكترونية.
  2. تأكَّد من تفعيل الفوترة لمشروعك على السحابة الإلكترونية. كيفية التحقّق من تفعيل الفوترة في مشروع

بدء Cloud Shell

Cloud Shell هي بيئة سطر أوامر تعمل في Google Cloud ومحمّلة مسبقًا بالأدوات اللازمة.

  1. انقر على تفعيل 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"

تفعيل واجهات برمجة التطبيقات المطلوبة

نفِّذ الأمر التالي في 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

لضمان تنفيذ واجهة برمجة التطبيقات بسلاسة والوصول إلى "بيانات الاعتماد التلقائية للتطبيق" (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.

لضمان توفّر سعة حوسبة مخصّصة لعمليات البحث في فهارس المتجهات وعمليات تجميع الرسومات وتنفيذ الدوال عن بُعد بدون أن تكون مقيّدًا بحدود وحدة المعالجة المركزية عند الطلب أو الحصص المشترَكة، أنشئ حجزًا في 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 مقابل خطأ التعرّف البصري على الأحرف 15
  • تبديل الأحرف والقيم غير المتوفّرة: تبديل تاريخ الميلاد (19340706 مقابل 19340760)، والولايات غير المتوفّرة ( )، ومعرّفات التأمين الاجتماعي غير المتوفّرة ( ).

في الخطوات القادمة، ستستخدم SOUNDEX عمليات الترميز الصوتية، ودوال المستخدم المحدّدة (UDF) لتسوية العناوين، ومسافة التعديل في Levenshtein، وAI.EMBED البحث المتّجه لسدّ هذه التناقضات وربط الملفات المكرّرة بدقة.

4. نشر دالة Address Validation Remote Function المعرَّفة من قِبل المستخدم

تؤدّي تسوية العناوين إلى توحيد أسماء الشوارع وحدود الضواحي والرموز البريدية قبل إجراء عملية المطابقة. Google Maps Address Validation API هي خدمة تقبل عنوانًا وتحدّد مكوّناته وتتحقّق من صحتها. في هذه الخطوة، ستنشئ دالة Cloud Function بلغة Python في Cloud Shell تعرض دالة معرَّفة من قِبل المستخدم للتحقّق من صحة العناوين وتوحيد تنسيقها في 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 Function وضبط أذونات "إدارة الهوية والوصول"

نفِّذ الأوامر التالية في Cloud Shell لنشر الجيل الثاني من Cloud Functions وإعداد عملية ربط بين Cloud Resource و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 Function وتطبيق عمليات ربط إدارة الهوية وإمكانية الوصول (IAM) بنجاح.

تسجيل دالة تسوية العنوان البعيد

عليك الآن تسجيل لغة تعريف البيانات (DDL) الخاصة بـ "الدالة البعيدة" في BigQuery (validate_address_udf) التي تربط صفوف جدول BigQuery بنقطة نهاية Cloud Functions التي تم نشرها (${FUNCTION_URL}).

نفِّذ الأمر التالي في Cloud Shell لاسترداد عنوان URL الخاص بدالة Cloud التي تم نشرها وتسجيل الدالة البعيدة تلقائيًا:

# 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. طلبات validate_address_udf للحصول على عناوين موحّدة ونتائج التحقّق من صحة العناوين
  2. تنشئ هذه السمة ترميزات صوتية لـ given_name وsurname للتعامل مع الاختلافات الإملائية.SOUNDEX
  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. إنشاء تضمينات الملفات الشخصية الدلالية وVector Search

بالإضافة إلى تسوية العناوين ومفاتيح Soundex الصوتية ومسافة التعديل في Levenshtein، يتيح BigQuery استخدام دوال التضمين المضمّنة المستندة إلى الذكاء الاصطناعي التوليدي من خلال 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. تسجيل أزواج المرشّحين ودمج الميزات الهجينة على الحافة

يصبح تقييم كل أزواج سجلات العملاء المحتملة (النمو التربيعي 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)
  • تشابه تعديل الاسم الأخير باستخدام مسافة Levenshtein EDIT_DISTANCE (الوزن: 0.20)
  • تشابه الاسم الأول (الوزن: 0.20)
  • نتيجة تطابق تاريخ الميلاد (الوزن: 0.15)
  • تشابه جاكارد لرموز العناوين (الوزن: 0.15) على 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 مع لغة طلبات الرسم البياني (GQL) وفقًا لمعيار ISO بشكلٍ أصلي من خلال رسومات بيانية للعناصر. ينشئ "رسم بياني للعلاقات" طريقة عرض منطقية للرسم البياني على جداول 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
);

تصوُّر مجموعات الرسم البياني لموسيقى الهوب الكورية

نفِّذ طلب البحث أدناه لتصوّر مجموعات العملاء المتطابقة على مستوى خطوتَين أو خطوة واحدة من خطوات العلاقة:

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

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

إنشاء "عرض مقاييس التقييم"

لاحتساب مقياس صحة النموذج ومقياس المراجعة ومقياس دقة الاختبار مقابل الجدول 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

الدقة

تذكُّر الإعلان

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 جيران من فهرس التضمين الأساسي السابق. إذا تطابق سجلّ وارد مع ملف عميل حالي أعلى من مستوى التشابه، يتم دمجه ديناميكيًا في مجموعة الملفات هذه ويرث خط الأساس canonical_customer_id (MATCHED_TO_EXISTING_CLUSTER). وإذا لم يتم العثور على أقرب جار لخط الأساس أعلى من الحد، يتم إنشاء معرّف فريد عالمي جديد للكيان (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).

إذا كانت مجموعة تم احتسابها حديثًا في عملية التنفيذ (t) تتشارك% 70 على الأقل من سجلّات الأعضاء مع مجموعة من عملية التنفيذ (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 لهذا المختبر، يمكنك حذف المشروع باتّباع الخطوات التالية:

gcloud projects delete ${GCP_PROJECT}

12. تهانينا

تهانينا! لقد أنشأت بنجاح محرّكًا شاملاً لحلّ هوية العميل داخل Google Cloud BigQuery باستخدام BigQuery Property Graph، وطلبات بحث ISO GQL، ومطابقة التشابه المختلط، ومطابقة دلتا التزايدية، وضمانات ثبات المجموعات الدائمة.

ما تعلّمته

  • كيفية نشر "دالة Cloud" من الجيل الثاني وإتاحتها كدالة BigQuery عن بُعد
  • كيفية المعالجة المُسبقة للخصائص الديمغرافية للعملاء باستخدام SOUNDEX عمليات الترميز الصوتية والتحقّق من صحة العناوين
  • كيفية تنفيذ حظر المرشحين واحتساب نتائج التشابه المختلط باستخدام مسافة Levenshtein (EDIT_DISTANCE) وتشابه Jaccard المميز.
  • كيفية إنشاء BigQuery Property Graph (CREATE PROPERTY GRAPH) على جداول العُقد والحواف
  • كيفية طلب مسارات الرسم البياني باستخدام GQL (GRAPH_TABLE) وفقًا لمعيار ISO مع أدوات تحديد الكمية {1, 2} k-hop
  • كيفية حلّ المجموعات العنقودية للعملاء الأساسيين وتقييم أداء النموذج مقارنةً بمقاييس الواقع
  • كيفية تنفيذ عملية مطابقة دلتا تدريجية لعمليات استيعاب الدفعات اليومية بدون إعادة معالجة مجموعة البيانات الكاملة
  • كيفية تطبيق ضمان حدّ التداخل (1 - ε) للحفاظ على ثبات المجموعات المتكرّرة في جميع عمليات تنفيذ المسار

الخطوات التالية

المستندات المرجعية