1. Введение
В этом практическом занятии вы создадите модульный, комплексный механизм разрешения личности клиента (сопоставления сущностей) непосредственно в Google Cloud BigQuery. Вы объедините Google Cloud Shell для развертывания инфраструктуры с редактором SQL BigQuery Studio для очистки данных, оценки кандидатов, построения графа свойств и обхода путей в соответствии со стандартом ISO GQL (Graph Query Language) .
Разрешение идентификационных данных — это основополагающая возможность для корпоративного учета «Клиент 360», обнаружения мошенничества и консолидации данных в нескольких системах. Поскольку существует множество допустимых подходов к разрешению идентификационных данных в зависимости от зрелости данных и потребностей бизнеса, все шаги в этом практическом руководстве являются модульными и необязательными. Конвейер разработан для демонстрации различных распространенных, используемых в производственной среде отраслевых методов, включая нормализацию адресов удаленных пользовательских функций (UDF), фонетическое блокирование Soundex, семантический векторный поиск ( AI.EMBED ), гибридную оценку признаков и кластеризацию графов свойств GQL, чтобы вы могли выборочно использовать шаблоны, подходящие для вашей архитектуры.
Методы сопоставления и пороговые значения оценки следует настраивать в зависимости от предпочтений вашей организации в отношении детерминированного или вероятностного сопоставления, что определяется целевым сценарием использования. Например, в системах строгого соответствия требованиям, выставления счетов или финансовых операциях обычно предпочтение отдается высокоточным детерминированным правилам (таким как точное совпадение номера социального страхования или налогового идентификатора) для предотвращения ложных связей, в то время как персонализация маркетинга, аналитика и системы рекомендаций часто используют вероятностное нечеткое сопоставление и сходство семантических векторов для максимизации полноты охвата и выявления тонких связей.

Что вы будете делать
- Загрузка эталонного набора данных FEBRL3 : загрузка синтетических записей о клиентах и эталонных пар соответствия в BigQuery.
- Внедрение удаленной пользовательской функции проверки адресов : разверните облачную функцию Python и зарегистрируйте удаленную функцию BigQuery для нормализации адресов улиц.
- Предварительная обработка данных профиля и фонетической кодировки : выполнить очистку данных SQL, вызвать пользовательскую функцию для работы с адресами и вычислить фонетические ключи
SOUNDEXи расстояния редактирования Левенштейна:- Фонетическая кодировка Soundex : фонетический алгоритм для индексации имен по звучанию, как оно произносится в английском языке. Он преобразует имена в 4-символьный код (начальная буква, за которой следуют три цифры), представляющий группы согласных звуков (например,
"John"и"Jon"соответствуютJ500, а"Smith"и"Smyth"—S530), предоставляя фонетические сигналы для оценки характеристик и инкрементального дельта-блокирования в реальном времени. - Расстояние Левенштейна (
EDIT_DISTANCE) : строковая метрика, измеряющая минимальное количество изменений одного символа (вставок, удалений или замен), необходимых для преобразования одной строки в другую, что позволяет точно и нечетко сопоставлять имена и адреса.
- Фонетическая кодировка Soundex : фонетический алгоритм для индексации имен по звучанию, как оно произносится в английском языке. Он преобразует имена в 4-символьный код (начальная буква, за которой следуют три цифры), представляющий группы согласных звуков (например,
- Генерация встраивания семантического профиля и векторного поиска : Генерируйте текстовые встраивания непосредственно в SQL с помощью
AI.EMBED(text-embedding-005) и находите K ближайших соседей с помощьюVECTOR_SEARCH, которые будут служить в качестве слоя для генерации сублинейных кандидатов. - Оценка пар-кандидатов и гибридное слияние признаков ребер : Использование пар-кандидатов, полученных с помощью векторного поиска, позволяет устранить сложность перекрестного соединения O(N²), вычислить взвешенные оценки сходства по нескольким признакам (SSN, расстояние редактирования Левенштейна, дата рождения, коэффициент Жаккара для адреса) и объединить ребра в единую таблицу кандидатов.
- Построение графа свойств и обход путей ISO GQL : Постройте
PROPERTY GRAPHBigQuery, выполните запросы путей ISO GQL{1, 2}(GRAPH_TABLE) для определения связанных кластеров клиентов, вычислите индивидуальные метрики оценки и выполните мягкую кластеризацию домохозяйств с использованием взвешивания графа Адамика-Адара. - Постепенное повышение разрешения и постоянная стабильность : обработка ежедневных партий данных о поступлении с постепенным сопоставлением дельта-значений.
- Консолидация кластеризации и стабильность кластеров (перекрытие 1-ε) : Обеспечение постоянной стабильности кластеров на протяжении всего выполнения конвейера с использованием гарантии порогового значения перекрытия (1-ε).
Что вам понадобится
- Веб-браузер, например Chrome .
- Проект Google Cloud с включенной функцией выставления счетов.
Данный практический семинар предназначен для инженеров данных, разработчиков баз данных и специалистов по искусственному интеллекту и машинному обучению всех уровней, включая начинающих.
Примерное время: 45 минут
Ориентировочная стоимость: менее 2,00 долларов США (используется облачные функции с оплатой по мере использования и обработка запросов BigQuery).
2. Прежде чем начать
Создайте проект в Google Cloud.
- В консоли Google Cloud на странице выбора проекта выберите или создайте проект 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"
Включить необходимые API
Выполните следующую команду в Cloud Shell, используя свою учетную запись пользователя, чтобы включить все необходимые сервисы Google Cloud:
gcloud services enable \
addressvalidation.googleapis.com \
cloudbuild.googleapis.com \
cloudfunctions.googleapis.com \
cloudresourcemanager.googleapis.com \
artifactregistry.googleapis.com \
aiplatform.googleapis.com \
run.googleapis.com \
bigqueryconnection.googleapis.com \
bigqueryreservation.googleapis.com \
bigquery.googleapis.com
Настройка служебной учетной записи и имитации пользователя (рекомендуется)
Для обеспечения бесперебойного выполнения API и доступа к учетным данным приложения по умолчанию (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.
Создание резервирования и назначения в BigQuery
Для выполнения GQL-запросов необходимо иметь резервирование, использующее версию Enterprise или Enterprise Plus. Создайте резервирование версии Enterprise с автомасштабированием в 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 для разрешения сущностей (содержащий 5000 записей о клиентах с кластерами множественных дубликатов, содержащими до 5 дубликатов на клиента) с помощью библиотеки recordlinkage на Python и записать исходные узлы клиентов ( customer_nodes ) и ссылки на истинные совпадения ( ground_truth_links ) в BigQuery, используя DataFrames BigQuery ( 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 перейдите в 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 | собственное имя | фамилия | номер_улицы | адрес_1 | адрес_2 | пригород | почтовый индекс | состояние | Дата рождения | 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. Разверните удаленную функцию проверки адресов (UDF).
Нормализация адресов стандартизирует названия улиц, границы пригородов и почтовые индексы перед выполнением сопоставления. API проверки адресов Google Maps — это сервис, который принимает адрес, идентифицирует его компоненты и проверяет их. На этом шаге вы развернете облачную функцию Python в Cloud Shell, которая предоставит BigQuery пользовательскую функцию для проверки и нормализации адресов.
Запись исходных файлов облачных функций
Выполните следующую команду в Cloud Shell, чтобы создать каталог с исходным кодом облачной функции и записать в него 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
Развертывание облачной функции и настройка разрешений IAM.
Выполните следующие команды в Cloud Shell, чтобы развернуть облачную функцию второго поколения и настроить подключение к облачному ресурсу 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
Вы должны увидеть сообщение о завершении развертывания облачной функции и успешном применении привязок IAM.
Зарегистрировать функцию нормализации удаленного адреса
Теперь вам нужно зарегистрировать DDL удаленной функции BigQuery ( validate_address_udf ), которая связывает строки таблицы BigQuery с развернутой конечной точкой облачной функции ( ${FUNCTION_URL} ).
Выполните следующую команду в Cloud Shell , чтобы получить URL-адрес развернутой облачной функции и автоматически зарегистрировать удаленную функцию:
# 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 . Этот запрос:
- Вызывает
validate_address_udfдля получения нормализованных адресов и результатов проверки. - Генерирует фонетические кодировки
SOUNDEXдляgiven_nameиsurnameдля обработки вариантов написания. - Создаёт структурированное поле
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_soundex | фамилия | surname_soundex | отформатированный_адрес |
| | | | | |
| | | | | |
| | | | | |
6. Генерация векторных представлений семантического профиля и векторный поиск.
В дополнение к нормализации адресов, фонетическим ключам Soundex и расстоянию редактирования Левенштейна, BigQuery поддерживает встроенные функции встраивания генеративного ИИ через AI.EMBED .
Используя AI.EMBED , BigQuery генерирует текстовые встраивания непосредственно в SQL, используя базовые модели (например, text-embedding-005 ):
Сгенерировать векторные представления профилей и выполнить векторный поиск Top-K.
Выполните указанные ниже запросы в редакторе 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 DISTINCT
LEAST(query.rec_id, base.rec_id) AS source_id,
GREATEST(query.rec_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 => 10,
distance_type => 'COSINE'
)
WHERE query.rec_id != base.rec_id AND distance <= 0.25;
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 ). Затем вы вычислите взвешенную оценку соответствия, объединив:
- Рейтинг соответствия SSN (вес:
0.30) - Поиск сходства фамилий с использованием расстояния Левенштейна
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 | счет матча |
| | |
| | |
| | |
Объедините данные о ребрах, полученные с помощью правил и векторного поиска, в единую таблицу.
Объедините потенциальные ребра, полученные с помощью нечеткого сопоставления на основе правил и семантического векторного поиска, в единую, дедуплицированную таблицу 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.08
)
GROUP BY source_id, target_id;
8. Построение графа свойств и обход путей ISO GQL.
BigQuery поддерживает ISO GQL (Graph Query Language) нативно через Property Graphs . Property Graph создает логическое графовое представление реляционных таблиц BigQuery без дублирования данных.
На этом этапе вы построите граф свойств customer_identity_graph используя вашу унифицированную таблицу потенциальных ребер ( final_matched_edges ), и выполните запрос к графу связей между профилями клиентов, используя обход пути k-hop {1, 2} .

Создание DDL-скрипта для графа свойств BigQuery
Выполните следующее DDL-выражение в редакторе SQL BigQuery Studio:
CREATE OR REPLACE PROPERTY GRAPH `identity_resolution.customer_identity_graph`
NODE TABLES (
`identity_resolution.customer_nodes_cleaned` AS `Customer`
KEY (rec_id)
)
EDGE TABLES (
`identity_resolution.final_matched_edges`
KEY (source_id, target_id)
SOURCE KEY (source_id) REFERENCES `Customer`(rec_id)
DESTINATION KEY (target_id) REFERENCES `Customer`(rec_id)
LABEL MATCHED_TO
);
Визуализация кластеров графа K-Hop
Выполните приведенный ниже запрос, чтобы визуализировать совпадающие кластеры клиентов на расстоянии от 1 до 2 шагов связи:
GRAPH `identity_resolution.customer_identity_graph`
MATCH p = (c1:Customer)-[e:MATCHED_TO]->{1, 2}(c2:Customer)
RETURN TO_JSON(p) AS graph_cluster_path
LIMIT 10;

Разрешение проблем с клиентскими кластерами Canonical
Выполните следующий запрос для преобразования кластеров сущностей в список 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 | количество записей | записи клиентов |
| | |
| | |
| | |
Создать представление метрик оценки
Отступление: Оценка на уровне кластера против прямых связей :
Приведенные ниже метрики оценки позволяют оценить точность определения сущностей клиентов ( resolved_customers ) путем генерации всех пар записей внутри кластера и сравнения их с эталонными данными на индивидуальном уровне ( ground_truth_links ). Оценка на уровне кластера позволяет в полной мере использовать преимущества разрешения путей графа ISO GQL (транзитивные многошаговые связи), что обеспечивает точное отражение качества сквозного разрешения сущностей.
Для вычисления точности, полноты и F1-меры по отношению к таблице ground_truth_links выполните следующую команду:
CREATE OR REPLACE VIEW `identity_resolution.evaluation_metrics` AS
WITH predictions AS (
-- Generate all pairwise record combinations within each resolved canonical customer cluster
SELECT
r1 AS source_id,
r2 AS target_id
FROM `identity_resolution.resolved_customers`,
UNNEST(customer_records) AS r1,
UNNEST(customer_records) AS r2
WHERE r1 < r2
),
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 | ложноположительные результаты | ложноотрицательные | точность | отзывать | f1_score |
| | | | | | | |
Разрешение кластеров домохозяйств с помощью взвешивания по графу Адамика-Адара
В то время как идентификация отдельных лиц позволяет определить, к какому именно человеку относятся записи, принадлежащие одному и тому же человеку, корпоративные архитектуры Customer 360 часто требуют группировки домохозяйств на более высоком уровне, объединяющей совместно проживающих по одному адресу лиц.
Поскольку синтетические эталонные тесты (например, FEBRL3) оценивают истинные данные на индивидуальном уровне, определение принадлежности к домохозяйству выполняется на последующем этапе. В отсутствие истории перемещений с указанием времени, лица, связанные с несколькими адресами, могут привести к чрезмерному объединению или фрагментации кластеров. Для решения этой проблемы мы используем взвешивание графов Адамика-Адара для построения условной принадлежности к домохозяйству.
Выполните приведенный ниже запрос в редакторе SQL BigQuery Studio, чтобы заполнить поле household_clusters , используя взвешивание по графу Адамика-Адара:
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 AS raw_household_affinity
FROM customer_addresses ca
JOIN address_degrees ad ON ca.formatted_address = ad.formatted_address
),
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
Чтобы визуально отследить полную трехуровневую иерархию идентификации — соединение некластеризованных исходных клиентов с разрешенными сущностями клиентов и далее с разрешенными сущностями домохозяйств — сначала создайте вспомогательные таблицы узлов и ребер и обновите DDL-скрипт графа свойств:
-- 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
);
Выполните приведенный ниже трехуровневый GQL-запрос в BigQuery Studio, чтобы визуализировать иерархию домохозяйств, состоящих из нескольких жильцов:
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 20;
Выполнение этого GQL-запроса в BigQuery Studio отображает интерактивную трехуровневую визуализацию графа, показывающую исходные записи профилей клиентов ( RawCustomer ), преобразованные в отдельные канонические сущности ( ResolvedCustomer ), которые связаны с общими сущностями домохозяйств, в которых проживает несколько человек ( ResolvedHousehold ).

9. Постепенное разрешение и устойчивая стабильность
В реальных корпоративных приложениях новые записи о клиентах поступают непрерывно в виде ежедневных или пакетных данных в режиме реального времени. Вместо повторного выполнения полного разрешения графа по всему историческому набору данных, механизм инкрементального дельта-сопоставления сравнивает новые поступающие записи с существующими разрешенными базовыми кластерами ( resolved_customers ).
Для эффективного достижения этой цели механизм использует векторный поиск ( VECTOR_SEARCH ) в качестве динамической кластеризации . Рассматривая каждую входящую запись как точку запроса , VECTOR_SEARCH извлекает набор из K ближайших соседей из исторического базового индекса встраивания. Если входящая запись соответствует существующему профилю клиента выше порогового значения сходства, она динамически объединяется с этим кластером и наследует базовый canonical_customer_id ( MATCHED_TO_EXISTING_CLUSTER ). Если выше порогового значения не найдено ни одного ближайшего соседа из базового списка, создается новый UUID сущности ( NEW_CUSTOMER_ENTITY ).
Записи о приеме образцов и поэтапном приеме партий.
Вставьте и выполните следующий DDL-скрипт в редакторе SQL BigQuery Studio, чтобы создать incremental_daily_intake :
CREATE OR REPLACE TABLE `identity_resolution.incremental_daily_intake` AS
SELECT * FROM UNNEST([
STRUCT(
'rec-9999-new-1' AS rec_id, 'erin' AS given_name, 'donaldson' AS surname,
'E650' AS given_name_soundex, 'D543' AS surname_soundex,
'19810427' AS date_of_birth, '2955815' AS soc_sec_id,
'13 hawkesbury crescent aralee lewiston 7018' AS formatted_address, '7018' AS postcode
),
STRUCT(
'rec-9999-new-2' AS rec_id, 'hollie' AS given_name, 'lillie-hinrichs' AS surname,
'H400' AS given_name_soundex, 'L446' AS surname_soundex,
'19251130' AS date_of_birth, '4920253' AS soc_sec_id,
'27 hemmings crescent kilvinton village banyo 4030' AS formatted_address, '4030' AS postcode
),
STRUCT(
'rec-9999-new-3' AS rec_id, 'sarah' AS given_name, 'ryan' AS surname,
'S600' AS given_name_soundex, 'R500' AS surname_soundex,
'20010101' AS date_of_birth, '999999999' AS soc_sec_id,
'500 market st melbourne vic 3000' AS formatted_address, '3000' AS postcode
),
STRUCT(
'rec-9999-new-4' AS rec_id, 'zzyzx' AS given_name, 'qx-vonderland' AS surname,
'Z220' AS given_name_soundex, 'Q215' AS surname_soundex,
'19991231' AS date_of_birth, '999887766' AS soc_sec_id,
'9999 zulu orbit station moon-base alpha 9999' AS formatted_address, '9999' AS postcode
)
]);
Выполнить запрос с инкрементальным дельта-сопоставлением
Выполните следующий запрос в редакторе SQL BigQuery, чтобы сопоставить данные с базовым набором данных, полученным в результате обработки:
CREATE OR REPLACE TABLE `identity_resolution.incremental_resolved_customers` AS
WITH historical_resolved_base AS (
SELECT c.rec_id, c.given_name, c.surname, c.given_name_soundex, c.surname_soundex, c.date_of_birth, c.soc_sec_id, c.formatted_address, c.postcode, r.canonical_customer_id
FROM `identity_resolution.customer_nodes_cleaned` c
JOIN (
SELECT canonical_customer_id, node_id
FROM `identity_resolution.resolved_customers`, UNNEST(customer_records) AS node_id
) r ON c.rec_id = r.node_id
),
incremental_intake AS (
SELECT
rec_id AS new_rec_id, given_name, surname, given_name_soundex, surname_soundex, date_of_birth, soc_sec_id, formatted_address, postcode,
CONCAT('Name: ', CONCAT(given_name, ' ', surname), '; Address: ', formatted_address, '; DOB: ', date_of_birth, '; SSN: ', soc_sec_id) AS profile_text
FROM `identity_resolution.incremental_daily_intake`
),
rule_delta_matches AS (
SELECT
i.new_rec_id,
h.canonical_customer_id AS matched_canonical_id,
h.rec_id AS matched_baseline_rec_id,
(
0.30 * (CASE WHEN i.soc_sec_id = h.soc_sec_id AND i.soc_sec_id != '' THEN 1.0 ELSE 0.0 END) +
0.20 * (1.0 - (EDIT_DISTANCE(i.surname, h.surname) / GREATEST(LENGTH(i.surname), LENGTH(h.surname), 1))) +
0.20 * (1.0 - (EDIT_DISTANCE(i.given_name, h.given_name) / GREATEST(LENGTH(i.given_name), LENGTH(h.given_name), 1))) +
0.15 * (CASE WHEN i.date_of_birth = h.date_of_birth THEN 1.0 ELSE 0.0 END) +
0.15 * (1.0 - (EDIT_DISTANCE(i.formatted_address, h.formatted_address) / GREATEST(LENGTH(i.formatted_address), LENGTH(h.formatted_address), 1)))
) AS match_score
FROM incremental_intake i
JOIN historical_resolved_base h
ON (i.soc_sec_id = h.soc_sec_id AND i.soc_sec_id != '')
OR (i.date_of_birth = h.date_of_birth AND i.given_name_soundex = h.given_name_soundex)
),
vector_intake_embeddings AS (
SELECT new_rec_id, AI.EMBED(profile_text, connection_id => 'us.address_val_conn', endpoint => 'text-embedding-005').result AS text_embedding
FROM incremental_intake
),
vector_delta_matches AS (
SELECT
v.query.new_rec_id,
h.canonical_customer_id AS matched_canonical_id,
h.rec_id AS matched_baseline_rec_id,
ROUND(1.0 - v.distance, 4) AS match_score
FROM VECTOR_SEARCH(
TABLE `identity_resolution.customer_embeddings`,
'text_embedding',
TABLE vector_intake_embeddings,
top_k => 3,
distance_type => 'COSINE'
) v
JOIN historical_resolved_base h ON v.base.rec_id = h.rec_id
WHERE v.distance <= 0.20
),
combined_delta AS (
SELECT
new_rec_id, matched_canonical_id, matched_baseline_rec_id, match_score,
'RULE_BASED' AS match_strategy
FROM rule_delta_matches WHERE match_score >= 0.55
UNION ALL
SELECT
new_rec_id, matched_canonical_id, matched_baseline_rec_id, match_score,
'VECTOR_SEARCH' AS match_strategy
FROM vector_delta_matches
),
aggregated_delta AS (
SELECT
new_rec_id,
matched_canonical_id,
ARRAY_AGG(matched_baseline_rec_id ORDER BY match_score DESC LIMIT 1)[OFFSET(0)] AS 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
),
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 | тип_задания | счет матча | стратегия матча |
| | | | | |
| | | | | |
| | | | | |
| | | | | |
10. Консолидация кластеров и стабильность кластеров (1-ε перекрытие)
В корпоративных системах, используемых в производственной среде, команды обычно регулярно (например, еженедельно или ежемесячно) перегруппировывают весь граф, чтобы включить новые ребра и источники данных. По мере формирования новых связей полная перегруппировка графа может привести к произвольному изменению или перестановке идентификаторов кластеров в ходе выполнения конвейера.
Для поддержания постоянных идентификаторов клиентов в нижестоящих системах CRM, CDP и выставления счетов, стабильность кластера оценивает перекрытие узлов между текущими кластерами Run( t ) и предыдущими кластерами Run( t -1) с использованием порогового значения перекрытия (1 - ε) (где ε = 0,30, что требует минимального перекрытия узлов в 70%).
Если вновь созданный кластер в Run( t ) имеет как минимум 70% общих записей с кластером из Run( t -1), он наследует исторический постоянный идентификатор клиента ( STABLE_EVOLUTION ). Совершенно новые кластеры получают вновь сгенерированные UUID ( NEW_CLUSTER_CREATED ).
Выполнить запрос на определение перекрытия и стабильности кластера.
Выполните следующий запрос в редакторе SQL BigQuery Studio, чтобы заполнить поле stable_resolved_customers :
CREATE OR REPLACE TABLE `identity_resolution.stable_resolved_customers` AS
WITH previous_run_clusters AS (
SELECT
canonical_customer_id AS previous_persistent_id,
node_id,
COUNT(*) OVER (PARTITION BY canonical_customer_id) AS previous_cluster_size
FROM `identity_resolution.resolved_customers`, UNNEST(customer_records) AS node_id
),
current_run_clusters AS (
SELECT
new_cluster_id,
node_id,
COUNT(*) OVER (PARTITION BY new_cluster_id) AS current_cluster_size
FROM (
SELECT
canonical_customer_id AS new_cluster_id,
node_id
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
FROM `identity_resolution.incremental_resolved_customers`
)
),
cluster_intersections AS (
SELECT
c.new_cluster_id,
p.previous_persistent_id,
c.current_cluster_size,
p.previous_cluster_size,
COUNT(c.node_id) AS shared_node_count,
COUNT(c.node_id) / GREATEST(p.previous_cluster_size, 1) 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, p.previous_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. Поздравляем!
Поздравляем! Вы успешно создали комплексную систему идентификации клиентов в Google Cloud BigQuery, используя BigQuery Property Graph, запросы ISO GQL, гибридное сопоставление по сходству, инкрементальное дельта-сопоставление и гарантии стабильности кластера.
Что вы узнали
- Как развернуть облачную функцию второго поколения и предоставить к ней доступ в качестве удаленной функции BigQuery.
- Как выполнить предварительную обработку демографических данных клиентов с использованием фонетической кодировки
SOUNDEXи проверки адресов. - Как выполнить блокировку кандидатов и вычислить гибридные показатели сходства, используя расстояние Левенштейна (
EDIT_DISTANCE) и сходство токенов Жаккара. - Как построить граф свойств BigQuery (
CREATE PROPERTY GRAPH) на основе таблиц узлов и ребер. - Как выполнить запрос к путям графа с использованием ISO GQL (
GRAPH_TABLE) с квантификаторами k-hop{1, 2}. - Как определить канонические кластеры клиентов и оценить производительность модели по отношению к эталонным метрикам.
- Как выполнить инкрементальное дельта-сопоставление для ежедневных пакетных данных без полной переработки набора данных.
- Как применить гарантию порогового значения перекрытия (1 - ε) для поддержания стабильности кластера на протяжении всего выполнения конвейера.
Следующие шаги
- Изучите документацию BigQuery Property Graph .
- Попробуйте использовать векторный поиск BigQuery и текстовые встраивания Vertex AI для генерации семантических кандидатов.
- Ознакомьтесь с информацией о BigQuery Remote Functions для масштабируемой интеграции с внешними API.