1. 소개
이 Codelab에서는 Google Cloud BigQuery 내에서 직접 모듈식 엔드 투 엔드 고객 ID 확인 (엔티티 매칭) 엔진을 빌드합니다. 인프라 배포를 위한 Google Cloud Shell과 데이터 정리, 후보 점수 매기기, 속성 그래프 구성, ISO GQL (그래프 쿼리 언어) 경로 순회를 위한 BigQuery Studio SQL 편집기를 결합합니다.
ID 해결은 엔터프라이즈 고객 360, 사기 감지, 다중 시스템 데이터 통합을 위한 기본 기능입니다. 데이터 성숙도와 비즈니스 요구사항에 따라 ID 해결에 유효한 접근 방식이 많기 때문에 이 Codelab의 모든 단계는 모듈식이며 선택사항입니다. 이 파이프라인은 원격 UDF 주소 정규화, Soundex 음성 차단, 의미론적 벡터 검색 (AI.EMBED), 하이브리드 기능 점수 매기기, GQL 속성 그래프 클러스터링 등 다양한 일반적인 프로덕션 등급 업계 기술을 보여주도록 설계되었으므로 아키텍처에 맞는 패턴을 선택적으로 채택할 수 있습니다.
일치 방법과 점수 기준점은 조직의 결정적 일치와 확률적 일치에 대한 요구에 따라 조정해야 하며, 이는 타겟 사용 사례에 따라 결정됩니다. 예를 들어 엄격한 규정 준수, 청구 또는 재무 운영에서는 잘못된 연결을 방지하기 위해 고정밀 결정론적 규칙 (예: 정확한 사회보장번호 또는 세금 ID 일치)을 선호하는 반면, 마케팅 맞춤설정, 분석, 추천 엔진에서는 재현율을 극대화하고 미묘한 연결을 발견하기 위해 확률적 퍼지 매칭 및 의미론적 벡터 유사성을 사용하는 경우가 많습니다.

실습할 내용
- FEBRL3 벤치마크 데이터 세트 수집: 합성 고객 레코드와 그라운드 트루스 일치 쌍을 BigQuery에 로드합니다.
- 주소 유효성 검사 원격 UDF 배포: Python Cloud 함수를 배포하고 BigQuery 원격 함수를 등록하여 도로 주소를 정규화합니다.
- 프로필 데이터 및 음성 인코딩 사전 처리: SQL 데이터 정리 실행, 주소 UDF 호출,
SOUNDEX음성 키 및 레벤슈타인 편집 거리 계산:- Soundex 음성 인코딩: 영어로 발음되는 소리를 기준으로 이름을 색인화하는 음성 알고리즘입니다. 이름을 자음 소리 그룹을 나타내는 4자리 코드 (이니셜과 세 자리 숫자)로 변환하여 (예:
"John"와"Jon"은 모두J500에 매핑되고"Smith"와"Smyth"은S530에 매핑됨) 기능 점수 매기기 및 실시간 증분 델타 차단을 위한 음성 일치 신호를 제공합니다. - Levenshtein 거리 (
EDIT_DISTANCE): 한 문자열을 다른 문자열로 변경하는 데 필요한 최소 단일 문자 수정 (삽입, 삭제, 대체)을 측정하는 문자열 측정항목으로, 정확한 퍼지 이름 및 주소 일치를 지원합니다.
- Soundex 음성 인코딩: 영어로 발음되는 소리를 기준으로 이름을 색인화하는 음성 알고리즘입니다. 이름을 자음 소리 그룹을 나타내는 4자리 코드 (이니셜과 세 자리 숫자)로 변환하여 (예:
- 시맨틱 프로필 임베딩 및 벡터 검색 생성:
AI.EMBED(text-embedding-005)를 사용하여 SQL에서 직접 텍스트 임베딩을 생성하고VECTOR_SEARCH를 사용하여 상위 K 최근접 이웃을 찾아 선형 이하 후보 생성 레이어로 제공합니다. - 후보 쌍 점수 지정 및 하이브리드 에지 기능 융합: 벡터 검색 후보 쌍을 활용하여 O(N²) 교차 조인 복잡성을 제거하고, 다중 기능 가중 유사성 점수 (SSN, Levenshtein 편집 거리, DOB, 주소 Jaccard)를 계산하고, 에지를 통합 후보 테이블로 융합합니다.
- 속성 그래프 구성 및 ISO GQL 경로 순회: BigQuery
PROPERTY GRAPH를 구성하고 ISO GQL{1, 2}경로 쿼리 (GRAPH_TABLE)를 실행하여 연결된 고객 클러스터를 해결하고, 개별 평가 측정항목을 계산하고, Adamic-Adar 그래프 가중치를 사용하여 소프트 가구 클러스터링을 실행합니다. - 증분 해결 및 지속적인 안정성: 증분 델타 매칭으로 일일 배치 인테이크를 처리합니다.
- 클러스터 통합 및 클러스터 안정성 (1-ε 중복): (1-ε) 중복 임곗값 보장을 사용하여 파이프라인 실행 전반에서 지속적인 클러스터 안정성을 적용합니다.
필요한 항목
- 웹브라우저(예: Chrome)
- 결제가 사용 설정된 Google Cloud 프로젝트.
이 Codelab은 초보자를 포함한 모든 수준의 데이터 엔지니어, 데이터베이스 개발자, AI/ML 실무자를 대상으로 합니다.
예상 소요 시간: 45분
예상 비용: 미화 2달러 미만 (사용한 만큼만 지불하는 Cloud Functions 및 BigQuery 쿼리 처리 사용)
2. 시작하기 전에
Google Cloud 프로젝트 만들기
- Google Cloud 콘솔의 프로젝트 선택기 페이지에서 Google Cloud 프로젝트를 선택하거나 만듭니다.
- 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 예약 및 할당 만들기 (선택사항 / 권장)
주문형 CPU 한도나 공유 할당량에 제약받지 않고 벡터 색인 검색, 그래프 집계, 원격 함수 실행을 위한 전용 컴퓨팅 용량을 확보하려면 Cloud Shell에서 자동 확장 기능이 있는 Enterprise Edition 예약을 만드세요.
# 1. Create a BigQuery Enterprise reservation with 0 baseline slots and 100 max autoscaling slots
bq mk --reservation \
--project_id=${GCP_PROJECT} \
--location=US \
--edition=ENTERPRISE \
--slots=0 \
--autoscale_max_slots=100 \
--ignore_idle_slots=true \
identity-res-reservation
# 2. Assign your Cloud project to the newly created reservation for query execution
bq mk --reservation_assignment \
--project_id=${GCP_PROJECT} \
--location=US \
--reservation_id=identity-res-reservation \
--job_type=QUERY \
--assignee_type=PROJECT \
--assignee_id=${GCP_PROJECT}
3. FEBRL3 고객 노드 데이터 세트 수집
주소 유효성 검사 원격 함수를 배포하고 ID 확인을 수행하기 전에 Python의 recordlinkage 라이브러리를 사용하여 합성 FEBRL3 엔티티 확인 벤치마크 데이터 세트 (고객당 최대 5개의 중복이 있는 다중 중복 클러스터가 포함된 고객 레코드 5,000개 포함)를 로드하고 BigQuery DataFrames (bigframes)를 사용하여 원시 고객 노드 (customer_nodes)와 그라운드 트루스 일치 링크 (ground_truth_links)를 BigQuery에 작성합니다.
Cloud Shell에서 다음 명령어를 실행하여 종속 항목을 설치하고 수집 스크립트를 실행합니다.
# 1. Install recordlinkage dataset library & bigframes (if outside Cloud Shell, activate your virtual environment first)
pip install recordlinkage bigframes --quiet
# 2. Write and execute the FEBRL3 dataset ingestion script
cat << 'EOF' > ingest_febrl.py
import os
import pandas as pd
import bigframes.pandas as bpd
from recordlinkage.datasets import load_febrl3
GCP_PROJECT = os.environ.get("GCP_PROJECT", "your-project-id")
DATASET_ID = "identity_resolution"
table_raw_id = f"{GCP_PROJECT}.{DATASET_ID}.customer_nodes"
table_gt_id = f"{GCP_PROJECT}.{DATASET_ID}.ground_truth_links"
print("Loading FEBRL3 benchmark dataset...")
df_nodes, true_links = load_febrl3(return_links=True)
df_nodes = df_nodes.reset_index()
df_nodes['dataset_source'] = 'febrl3'
for col in df_nodes.columns:
if df_nodes[col].dtype == 'object':
df_nodes[col] = df_nodes[col].fillna('')
print("Ingesting raw customer nodes into BigQuery via BigQuery DataFrames...")
bf_nodes = bpd.read_pandas(df_nodes)
bf_nodes.to_gbq(table_raw_id, if_exists="replace")
print("Ingesting ground truth links into BigQuery via BigQuery DataFrames...")
df_gt = pd.DataFrame(list(true_links), columns=["source_id", "target_id"])
bf_gt = bpd.read_pandas(df_gt)
bf_gt.to_gbq(table_gt_id, if_exists="replace")
print(f"Raw customer nodes ingested into `{table_raw_id}` ({len(df_nodes):,} rows).")
print(f"Ground truth links ingested into `{table_gt_id}` ({len(df_gt):,} pairs).")
EOF
python3 ingest_febrl.py
Google Cloud 콘솔에서 BigQuery Studio로 이동하고 새 SQL 쿼리 탭 (+)을 연 다음 아래 쿼리를 실행하여 수집된 고객 노드 테이블을 검사합니다.
SELECT rec_id, given_name, surname, street_number, address_1, address_2, suburb, postcode, state, date_of_birth, soc_sec_id, dataset_source
FROM `identity_resolution.customer_nodes`
LIMIT 5;
다음과 비슷한 출력이 표시됩니다.
rec_id | given_name | 성 | street_number | address_1 | address_2 | 교외 지역 | 우편번호 | 주 | date_of_birth | soc_sec_id | dataset_source |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| |
|
|
|
|
|
|
|
|
|
|
| |
|
|
|
|
|
|
|
|
|
| ||
|
|
|
|
|
|
|
|
|
|
|
|
벤치마크 데이터 세트가 중복 클러스터에 현실적인 더티 데이터를 도입하는 방법을 확인하세요.
- 발음 및 맞춤법 변형:
brentvs.brnt/bernt,woodvs.woode/wod,cliftonvs.cliffton - 주소 약어 및 오타 처리:
girdlestone circuit대girdelstone circut/girdlestone cir/girdlestone crt, 번지수11대 OCR 오류15 - 문자 자리바꿈 및 누락된 값: 생년월일 자리바꿈 (
19340706vs.19340760), 누락된 주 ( ), 누락된 사회보장번호 ( ).
다음 단계에서는 SOUNDEX 음성 인코딩, 주소 정규화 UDF, 레벤슈타인 편집 거리, AI.EMBED 벡터 검색을 사용하여 이러한 불일치를 해소하고 중복 프로필을 정확하게 연결합니다.
4. Address Validation 원격 함수 UDF 배포
주소 정규화는 일치 작업을 실행하기 전에 도로명, 교외 경계, 우편번호를 표준화합니다. Google Maps Address Validation API는 주소를 수락하고 주소 구성요소를 식별하여 유효성을 검사하는 서비스입니다. 이 단계에서는 주소 유효성 검사 및 정규화 UDF를 BigQuery에 노출하는 Python Cloud 함수를 Cloud Shell에 배포합니다.
Cloud 함수 소스 파일 작성
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 함수 배포 및 IAM 권한 구성
Cloud Shell에서 다음 명령어를 실행하여 2세대 Cloud 함수를 배포하고 BigQuery Cloud 리소스 연결을 구성합니다.
# 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 함수 배포가 완료되고 IAM 바인딩이 성공적으로 적용되었음을 나타내는 출력이 표시됩니다.
원격 주소 정규화 함수 등록
이제 BigQuery 테이블 행을 배포된 Cloud 함수 엔드포인트 (${FUNCTION_URL})에 연결하는 BigQuery 원격 함수 DDL (validate_address_udf)을 등록합니다.
Cloud Shell에서 다음 명령어를 실행하여 배포된 Cloud 함수 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. 프로필 데이터 및 음성 인코딩 전처리
이 단계에서는 수집된 customer_nodes 테이블에 대해 BigQuery SQL 전처리 쿼리를 실행합니다.
데이터 정리 및 음성 기능 쿼리 실행
BigQuery Studio SQL 편집기에서 아래 쿼리를 실행하여 customer_nodes_cleaned를 만듭니다. 이 쿼리는 다음과 같은 작업을 합니다.
validate_address_udf를 호출하여 정규화된 주소와 검증 결과를 가져옵니다.- 맞춤법 변형을 처리하기 위해
given_name및surname의SOUNDEX음성 인코딩을 생성합니다. - 구조화된
profile_text필드를 생성합니다.
CREATE OR REPLACE TABLE `identity_resolution.customer_nodes_cleaned` AS
WITH raw_data AS (
SELECT
rec_id, dataset_source,
TRIM(LOWER(given_name)) AS given_name_clean,
TRIM(LOWER(surname)) AS surname_clean,
`identity_resolution.validate_address_udf`(street_number, address_1, address_2, suburb, state, postcode) AS addr_json,
TRIM(suburb) AS suburb, TRIM(state) AS state, TRIM(postcode) AS postcode,
TRIM(date_of_birth) AS date_of_birth, TRIM(soc_sec_id) AS soc_sec_id
FROM `identity_resolution.customer_nodes`
)
SELECT
rec_id, dataset_source,
given_name_clean AS given_name,
SOUNDEX(given_name_clean) AS given_name_soundex,
surname_clean AS surname,
SOUNDEX(surname_clean) AS surname_soundex,
CONCAT(given_name_clean, ' ', surname_clean) AS full_name,
STRING(addr_json.formatted_address) AS formatted_address,
BOOL(addr_json.address_is_valid) AS address_is_valid,
STRING(addr_json.validation_granularity) AS validation_granularity,
STRING(addr_json.possible_next_action) AS possible_next_action,
suburb, state, postcode, date_of_birth, soc_sec_id,
CONCAT('Name: ', CONCAT(given_name_clean, ' ', surname_clean), '; Address: ', STRING(addr_json.formatted_address), '; DOB: ', date_of_birth, '; SSN: ', soc_sec_id) AS profile_text
FROM raw_data;
정리된 노드 테이블을 쿼리합니다.
SELECT rec_id, given_name, given_name_soundex, surname, surname_soundex, formatted_address
FROM `identity_resolution.customer_nodes_cleaned`
LIMIT 5;
다음과 비슷한 출력이 표시됩니다.
rec_id | given_name | given_name_soundex | 성 | surname_soundex | formatted_address |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
6. 시맨틱 프로필 임베딩 및 벡터 검색 생성
BigQuery는 주소 정규화, Soundex 음성 키, Levenshtein 편집 거리 외에도 AI.EMBED를 통해 기본 제공 생성형 AI 임베딩 함수를 지원합니다.
AI.EMBED를 사용하면 BigQuery에서 수동 벡터 색인 DDL 없이 기반 모델 (예: text-embedding-005)을 사용하여 SQL에서 직접 텍스트 삽입을 생성합니다.
프로필 임베딩 생성 및 Top-K 벡터 검색 실행
BigQuery Studio SQL 편집기에서 아래 쿼리를 실행합니다.
-- 1. Generate Customer Profile Embeddings using AI.EMBED (offloaded to Vertex AI)
CREATE OR REPLACE TABLE `identity_resolution.customer_embeddings` AS
SELECT
rec_id,
dataset_source,
profile_text,
AI.EMBED(profile_text, connection_id => 'us.address_val_conn', endpoint => 'text-embedding-005').result AS text_embedding
FROM `identity_resolution.customer_nodes_cleaned`;
-- 2. Execute VECTOR_SEARCH for Top-K Nearest Neighbors Candidate Generation
CREATE OR REPLACE TABLE `identity_resolution.vector_candidate_edges` AS
SELECT
query.rec_id AS source_id,
base.rec_id AS target_id,
distance AS vector_distance
FROM VECTOR_SEARCH(
TABLE `identity_resolution.customer_embeddings`,
'text_embedding',
TABLE `identity_resolution.customer_embeddings`,
top_k => 5,
distance_type => 'COSINE'
)
WHERE query.rec_id < base.rec_id AND distance <= 0.20;
7. 후보 쌍 점수 매기기 및 하이브리드 Edge 기능 융합
데이터 세트 규모가 커지면 가능한 모든 고객 레코드 쌍을 평가하는 것이 (O(N²) 2차 성장) 계산상 불가능해집니다. BigQuery와 같은 관계형 데이터베이스 엔진에서 여러 열에 걸쳐 복잡한 OR 조인 조건을 사용하여 규칙 기반 차단을 구현하려고 하면 쿼리 최적화 프로그램이 단일 equi-join 키에서 확장 가능한 해시 조인 또는 정렬-병합 조인을 사용하지 못합니다 (예: a.soc_sec_id = b.soc_sec_id OR a.given_name_soundex = b.given_name_soundex OR ...에서 조인). 대신 엔진은 O(N²) 교차 조인으로 대체하고 각 쌍을 필터링하므로 확장 시 실패합니다.
이 단계에서는 벡터 검색 테이블 (vector_candidate_edges)에서 생성된 후보 쌍을 가져와 빠르고 색인이 생성된 동등 조인 (ON c.source_id = a.rec_id 및 ON c.target_id = b.rec_id)을 통해 customer_nodes_cleaned에 조인합니다. 그런 다음 다음을 결합하여 가중치 일치 점수를 계산합니다.
- SSN 일치 점수 (가중치:
0.30) - Levenshtein 거리
EDIT_DISTANCE를 사용하는 성 수정 유사성 (가중치:0.20) - 이름 수정 유사성 (가중치:
0.20) - 생년월일 일치 점수 (가중치:
0.15) SPLIT(LOWER(formatted_address), ' ')에 대한 주소 토큰 Jaccard 유사성 (가중치:0.15)
후보 가장자리 및 가중 유사성 점수 계산
BigQuery Studio SQL 편집기에서 다음 쿼리를 실행하여 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 |
|
|
|
|
|
|
|
|
|
규칙 기반 및 벡터 검색 에지를 통합 테이블로 병합
규칙 기반 퍼지 일치와 시맨틱 벡터 검색의 후보 가장자리를 중복이 삭제된 단일 final_matched_edges 테이블로 결합합니다.
CREATE OR REPLACE TABLE `identity_resolution.final_matched_edges` AS
SELECT
source_id,
target_id,
MAX(edge_weight) AS edge_weight,
IF(COUNT(DISTINCT edge_type) > 1, 'HYBRID', MAX(edge_type)) AS edge_type
FROM (
SELECT source_id, target_id, match_score AS edge_weight, 'RULE_BASED' AS edge_type
FROM `identity_resolution.matched_edges`
UNION ALL
SELECT source_id, target_id, ROUND(1.0 - vector_distance, 4) AS edge_weight, 'VECTOR_SEARCH' AS edge_type
FROM `identity_resolution.vector_candidate_edges`
WHERE vector_distance <= 0.05
)
GROUP BY source_id, target_id;
8. 속성 그래프 구성 및 ISO GQL 경로 순회
BigQuery는 속성 그래프를 통해 ISO GQL (Graph Query Language)을 기본적으로 지원합니다. 속성 그래프는 데이터 중복 없이 관계형 BigQuery 테이블에 논리적 그래프 뷰를 만듭니다.
이 단계에서는 통합 후보 에지 테이블 (final_matched_edges)을 사용하여 속성 그래프 customer_identity_graph를 빌드하고 {1, 2} k-hop 경로 순회를 사용하여 고객 프로필 전반에서 그래프 연결을 쿼리합니다.

BigQuery 속성 그래프 DDL 만들기
BigQuery Studio SQL 편집기에서 다음 DDL 문을 실행합니다.
CREATE OR REPLACE PROPERTY GRAPH `identity_resolution.customer_identity_graph`
NODE TABLES (
`identity_resolution.customer_nodes_cleaned` AS `Customer`
KEY (rec_id)
)
EDGE TABLES (
`identity_resolution.final_matched_edges`
KEY (source_id, target_id)
SOURCE KEY (source_id) REFERENCES `Customer`(rec_id)
DESTINATION KEY (target_id) REFERENCES `Customer`(rec_id)
LABEL MATCHED_TO
);
K-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;

표준 고객 클러스터 해결
다음 쿼리를 실행하여 엔티티 클러스터를 resolved_customers로 변환합니다.
CREATE OR REPLACE TABLE `identity_resolution.resolved_customers` AS
WITH graph_paths AS (
SELECT
source_node_id, target_node_id
FROM GRAPH_TABLE(
`identity_resolution.customer_identity_graph`
MATCH (c1:Customer)-[e:MATCHED_TO]->{1, 2}(c2:Customer)
COLUMNS (c1.rec_id AS source_node_id, c2.rec_id AS target_node_id)
)
),
all_connections AS (
SELECT source_node_id AS node_id, target_node_id AS connected_id FROM graph_paths
UNION DISTINCT
SELECT target_node_id AS node_id, source_node_id AS connected_id FROM graph_paths
UNION DISTINCT
SELECT rec_id AS node_id, rec_id AS connected_id FROM `identity_resolution.customer_nodes_cleaned`
),
clusters AS (
SELECT
node_id,
MIN(connected_id) AS canonical_customer_id
FROM all_connections
GROUP BY node_id
)
SELECT
canonical_customer_id,
ARRAY_AGG(node_id) AS customer_records,
COUNT(node_id) AS record_count
FROM clusters
GROUP BY canonical_customer_id;
해결된 클러스터 테이블을 쿼리합니다.
SELECT canonical_customer_id, record_count, customer_records
FROM `identity_resolution.resolved_customers`
ORDER BY record_count DESC;
다음과 비슷한 출력이 표시됩니다.
canonical_customer_id | record_count | customer_records |
|
|
|
|
|
|
|
|
|
평가 측정항목 뷰 만들기
ground_truth_links 테이블에 대한 정밀도, 재현율, F1 점수를 계산하려면 다음을 실행하세요.
CREATE OR REPLACE VIEW `identity_resolution.evaluation_metrics` AS
WITH predictions AS (
SELECT
LEAST(source_id, target_id) AS source_id,
GREATEST(source_id, target_id) AS target_id
FROM `identity_resolution.final_matched_edges`
WHERE edge_weight >= 0.55
),
ground_truth AS (
SELECT
LEAST(source_id, target_id) AS source_id,
GREATEST(source_id, target_id) AS target_id
FROM `identity_resolution.ground_truth_links`
),
stats AS (
SELECT
COUNT(g.source_id) AS total_ground_truth,
COUNT(p.source_id) AS total_predictions,
COUNTIF(p.source_id IS NOT NULL AND g.source_id IS NOT NULL) AS true_positives,
COUNTIF(p.source_id IS NOT NULL AND g.source_id IS NULL) AS false_positives,
COUNTIF(p.source_id IS NULL AND g.source_id IS NOT NULL) AS false_negatives
FROM ground_truth g
FULL OUTER JOIN predictions p ON g.source_id = p.source_id AND g.target_id = p.target_id
)
SELECT
total_ground_truth, total_predictions, true_positives, false_positives, false_negatives,
ROUND(true_positives / NULLIF(true_positives + false_positives, 0), 4) AS precision,
ROUND(true_positives / NULLIF(true_positives + false_negatives, 0), 4) AS recall,
ROUND(2 * true_positives / NULLIF((2 * true_positives) + false_positives + false_negatives, 0), 4) AS f1_score
FROM stats;
평가 측정항목 뷰를 쿼리합니다.
SELECT * FROM `identity_resolution.evaluation_metrics`;
다음과 비슷한 출력이 표시됩니다.
total_ground_truth | total_predictions | true_positives | false_positives | false_negatives | precision | recall | f1_score |
|
|
|
|
|
|
|
|
Adamic-Adar 그래프 가중치를 통한 가구 클러스터 해결
개별 ID 해결은 동일한 사람에 속하는 레코드를 해결하지만, 엔터프라이즈 고객 360 아키텍처에서는 주소를 공유하는 동거인을 그룹화하는 상위 수준의 가구 엔티티가 필요한 경우가 많습니다.
합성 벤치마크 (예: FEBRL3)는 개인 수준 정답을 평가하므로 가구 확인은 다운스트림 단계로 실행됩니다. 타임스탬프가 지정된 이전 기록이 없으면 여러 주소에 연결된 개인이 과도한 병합이나 클러스터 조각화를 유발할 수 있습니다. 이 문제를 해결하기 위해 Adamic-Adar 그래프 가중치를 사용하여 소프트 가구 멤버십을 구성합니다.
BigQuery Studio SQL 편집기에서 아래 쿼리를 실행하여 Adamic-Adar 그래프 가중치를 사용하여 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,
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을 통해 엔드 투 엔드 ID 계층 구조 시각화
클러스터링되지 않은 원시 고객을 해결된 고객 엔티티에 연결하고 해결된 가구 엔티티로 연결하는 등 완전한 3계층 ID 계층 구조를 시각적으로 추적하려면 BigQuery Studio에서 다음 DDL 및 ISO GQL 쿼리를 실행하세요.
-- 1. Create Household Node Table
CREATE OR REPLACE TABLE `identity_resolution.household_nodes` AS
SELECT DISTINCT
canonical_household_id,
household_address,
total_residents
FROM `identity_resolution.household_clusters`;
-- 2. Create Unresolved Record to Resolved Entity Edge Table
CREATE OR REPLACE TABLE `identity_resolution.customer_entity_edges` AS
SELECT DISTINCT
rec_id,
canonical_customer_id
FROM `identity_resolution.resolved_customers`,
UNNEST(customer_records) AS rec_id;
-- 3. Create Primary Household Edge Table (Highest Weighted Household Rank = 1)
CREATE OR REPLACE TABLE `identity_resolution.primary_household_edges` AS
SELECT
canonical_customer_id,
canonical_household_id,
household_membership_weight,
household_rank
FROM `identity_resolution.household_clusters`
WHERE household_rank = 1;
-- 4. Update Unified Property Graph DDL
CREATE OR REPLACE PROPERTY GRAPH `identity_resolution.customer_identity_graph`
NODE TABLES (
`identity_resolution.customer_nodes_cleaned` AS `RawCustomer`
KEY (rec_id),
`identity_resolution.resolved_customers` AS `ResolvedCustomer`
KEY (canonical_customer_id),
`identity_resolution.household_nodes` AS `ResolvedHousehold`
KEY (canonical_household_id)
)
EDGE TABLES (
`identity_resolution.final_matched_edges`
KEY (source_id, target_id)
SOURCE KEY (source_id) REFERENCES `RawCustomer`(rec_id)
DESTINATION KEY (target_id) REFERENCES `RawCustomer`(rec_id)
LABEL MATCHED_TO,
`identity_resolution.customer_entity_edges`
KEY (rec_id, canonical_customer_id)
SOURCE KEY (rec_id) REFERENCES `RawCustomer`(rec_id)
DESTINATION KEY (canonical_customer_id) REFERENCES `ResolvedCustomer`(canonical_customer_id)
LABEL RESOLVED_TO,
`identity_resolution.primary_household_edges`
KEY (canonical_customer_id, canonical_household_id)
SOURCE KEY (canonical_customer_id) REFERENCES `ResolvedCustomer`(canonical_customer_id)
DESTINATION KEY (canonical_household_id) REFERENCES `ResolvedHousehold`(canonical_household_id)
LABEL BELONGS_TO_HOUSEHOLD
);
-- 5. Execute 3-Tier GQL Query for Multi-Resident Household Visualization
GRAPH `identity_resolution.customer_identity_graph`
MATCH p = (raw:RawCustomer)-[e1:RESOLVED_TO]->(c:ResolvedCustomer)-[e2:BELONGS_TO_HOUSEHOLD]->(h:ResolvedHousehold)
WHERE h.total_residents > 1
RETURN TO_JSON(p) AS multi_resident_household_hierarchy_path
LIMIT 15;
BigQuery Studio에서 이 GQL 쿼리를 실행하면 공유 다중 거주 가구 엔티티 (ResolvedHousehold)에 연결된 개별 표준 엔티티 (ResolvedCustomer)로 확인된 원시 고객 프로필 레코드 (RawCustomer)를 보여주는 3계층 대화형 그래프 시각화 캔버스가 렌더링됩니다.

9. 점진적 해결 및 지속적인 안정성
실제 엔터프라이즈 애플리케이션에서는 일일 또는 실시간 배치 인테이크를 통해 신규 고객 레코드가 지속적으로 도착합니다. 전체 이전 데이터 세트에 대해 전체 그래프 해결을 다시 실행하는 대신 증분 델타 매칭 엔진은 새로 수신된 레코드를 기존에 해결된 기준 클러스터 (resolved_customers)와 비교합니다.
엔진은 이를 효율적으로 달성하기 위해 벡터 검색 (
VECTOR_SEARCH
)을 동적 클러스터링의 한 형태로 사용합니다. VECTOR_SEARCH는 각 수신 레코드를 쿼리 포인트로 취급하여 기록 기준 임베딩 색인에서 상위 K개의 최근접 이웃 집합을 가져옵니다. 수신 레코드가 유사성 기준점 이상의 기존 고객 프로필과 일치하면 해당 클러스터로 동적으로 병합되고 기준 canonical_customer_id (MATCHED_TO_EXISTING_CLUSTER)를 상속합니다. 기준점 이상의 가장 가까운 기준 이웃이 없으면 새 항목 UUID가 생성됩니다 (NEW_CUSTOMER_ENTITY).
샘플 증분 일괄 수집 기록 수집
BigQuery Studio SQL 편집기에 다음 DDL을 붙여넣고 실행하여 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
)
]);
증분 델타 일치 쿼리 실행
BigQuery SQL 편집기에서 다음 쿼리를 실행하여 해결된 기준 데이터 세트에 대해 델타 매칭을 실행합니다.
CREATE OR REPLACE TABLE `identity_resolution.incremental_resolved_customers` AS
WITH historical_resolved_base AS (
SELECT c.rec_id, c.given_name, c.surname, c.given_name_soundex, c.surname_soundex, c.date_of_birth, c.soc_sec_id, c.formatted_address, c.postcode, r.canonical_customer_id
FROM `identity_resolution.customer_nodes_cleaned` c
JOIN (
SELECT canonical_customer_id, node_id
FROM `identity_resolution.resolved_customers`, UNNEST(customer_records) AS node_id
) r ON c.rec_id = r.node_id
),
incremental_intake AS (
SELECT
rec_id AS new_rec_id, given_name, surname, given_name_soundex, surname_soundex, date_of_birth, soc_sec_id, formatted_address, postcode,
CONCAT('Name: ', CONCAT(given_name, ' ', surname), '; Address: ', formatted_address, '; DOB: ', date_of_birth, '; SSN: ', soc_sec_id) AS profile_text
FROM `identity_resolution.incremental_daily_intake`
),
rule_delta_matches AS (
SELECT
i.new_rec_id,
h.canonical_customer_id AS matched_canonical_id,
h.rec_id AS matched_baseline_rec_id,
(
0.30 * (CASE WHEN i.soc_sec_id = h.soc_sec_id AND i.soc_sec_id != '' THEN 1.0 ELSE 0.0 END) +
0.20 * (1.0 - (EDIT_DISTANCE(i.surname, h.surname) / GREATEST(LENGTH(i.surname), LENGTH(h.surname), 1))) +
0.20 * (1.0 - (EDIT_DISTANCE(i.given_name, h.given_name) / GREATEST(LENGTH(i.given_name), LENGTH(h.given_name), 1))) +
0.15 * (CASE WHEN i.date_of_birth = h.date_of_birth THEN 1.0 ELSE 0.0 END) +
0.15 * (1.0 - (EDIT_DISTANCE(i.formatted_address, h.formatted_address) / GREATEST(LENGTH(i.formatted_address), LENGTH(h.formatted_address), 1)))
) AS match_score
FROM incremental_intake i
JOIN historical_resolved_base h
ON (i.soc_sec_id = h.soc_sec_id AND i.soc_sec_id != '')
OR (i.date_of_birth = h.date_of_birth AND i.given_name_soundex = h.given_name_soundex)
),
vector_intake_embeddings AS (
SELECT new_rec_id, AI.EMBED(profile_text, connection_id => 'us.address_val_conn', endpoint => 'text-embedding-005').result AS text_embedding
FROM incremental_intake
),
vector_delta_matches AS (
SELECT
v.query.new_rec_id,
h.canonical_customer_id AS matched_canonical_id,
h.rec_id AS matched_baseline_rec_id,
ROUND(1.0 - v.distance, 4) AS match_score
FROM VECTOR_SEARCH(
TABLE `identity_resolution.customer_embeddings`,
'text_embedding',
TABLE vector_intake_embeddings,
top_k => 3,
distance_type => 'COSINE'
) v
JOIN historical_resolved_base h ON v.base.rec_id = h.rec_id
WHERE v.distance <= 0.20
),
combined_delta AS (
SELECT
new_rec_id, matched_canonical_id, matched_baseline_rec_id, match_score,
'RULE_BASED' AS match_strategy
FROM rule_delta_matches WHERE match_score >= 0.55
UNION ALL
SELECT
new_rec_id, matched_canonical_id, matched_baseline_rec_id, match_score,
'VECTOR_SEARCH' AS match_strategy
FROM vector_delta_matches
),
aggregated_delta AS (
SELECT
new_rec_id, matched_canonical_id, matched_baseline_rec_id,
MAX(match_score) AS match_score,
CASE
WHEN COUNT(DISTINCT match_strategy) > 1 THEN 'BOTH'
ELSE MAX(match_strategy)
END AS match_strategy
FROM combined_delta
GROUP BY new_rec_id, matched_canonical_id, matched_baseline_rec_id
),
best_matches AS (
SELECT
new_rec_id, matched_canonical_id, matched_baseline_rec_id, match_score, match_strategy,
ROW_NUMBER() OVER(PARTITION BY new_rec_id ORDER BY match_score DESC) AS rank
FROM aggregated_delta
)
SELECT
i.new_rec_id AS record_id,
i.given_name, i.surname,
COALESCE(b.matched_canonical_id, GENERATE_UUID()) AS persistent_canonical_customer_id,
h.canonical_household_id AS assigned_household_id,
CASE WHEN b.matched_canonical_id IS NOT NULL THEN 'MATCHED_TO_EXISTING_CLUSTER' ELSE 'NEW_CUSTOMER_ENTITY' END AS assignment_type,
b.matched_baseline_rec_id,
b.match_score,
COALESCE(b.match_strategy, 'NONE') AS match_strategy
FROM incremental_intake i
LEFT JOIN best_matches b ON i.new_rec_id = b.new_rec_id AND b.rank = 1
LEFT JOIN `identity_resolution.primary_household_edges` h ON b.matched_canonical_id = h.canonical_customer_id;
증분 해결 결과를 쿼리합니다.
SELECT record_id, persistent_canonical_customer_id, assigned_household_id, assignment_type, match_score, match_strategy
FROM `identity_resolution.incremental_resolved_customers`;
다음과 비슷한 출력이 표시됩니다.
record_id | persistent_canonical_customer_id | assigned_household_id | assignment_type | match_score | match_strategy |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
10. 클러스터링 통합 및 클러스터 안정성 (1-ε 중복)
프로덕션 엔터프라이즈 시스템에서 팀은 일반적으로 새로운 에지와 데이터 소스를 통합하기 위해 반복적으로 (예: 주별 또는 월별) 전체 그래프를 재클러스터링합니다. 새로운 관계가 형성되면 전체 그래프 재클러스터링으로 인해 파이프라인 실행 전반에서 클러스터 식별자가 임의로 이동하거나 뒤집힐 수 있습니다.
다운스트림 CRM, CDP, 결제 시스템의 지속적인 고객 ID를 유지하기 위해 클러스터 안정성은 (1 - ε) 중복 임계값 (여기서 ε = 0.30, 최소 70% 노드 중복 필요)을 사용하여 현재 실행 (t) 클러스터와 이전 실행 (t-1) 클러스터 간의 노드 중복을 평가합니다.
Run (t)에서 새로 계산된 클러스터가 Run (t-1)의 클러스터와 회원 레코드의 70% 이상을 공유하는 경우 이전 영구 고객 ID (STABLE_EVOLUTION)를 상속합니다. 완전히 새로운 클러스터에는 새로 생성된 UUID (NEW_CLUSTER_CREATED)가 부여됩니다.
클러스터 중복 및 안정성 쿼리 실행
BigQuery Studio SQL 편집기에서 다음 쿼리를 실행하여 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 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
11. 삭제
Google Cloud 계정에 지속적으로 비용이 청구되지 않도록 배포된 리소스와 BigQuery 데이터 세트를 정리합니다.
Cloud Shell에서 다음을 실행합니다.
# 1. Delete BigQuery Dataset
bq rm -r -f -d ${GCP_PROJECT}:${DATASET_ID}
# 2. Delete Cloud Function (2nd-Gen)
gcloud functions delete validate_address_udf --region=${REGION} --gen2 --quiet
# 3. Delete BigQuery Cloud Connection
bq rm -f --connection US.address_val_conn
이 실습을 위해 전용 Google Cloud 프로젝트를 만든 경우 프로젝트를 삭제할 수 있습니다.
gcloud projects delete ${GCP_PROJECT}
12. 축하합니다
수고하셨습니다 BigQuery 속성 그래프, ISO GQL 쿼리, 하이브리드 유사성 매칭, 증분 델타 매칭, 지속적인 클러스터 안정성 보장을 사용하여 Google Cloud BigQuery 내에서 엔드 투 엔드 고객 ID 해결 엔진을 성공적으로 빌드했습니다.
학습한 내용
- 2세대 Cloud 함수를 배포하고 BigQuery 원격 함수로 노출하는 방법
SOUNDEX음성 인코딩 및 주소 유효성 검사를 사용하여 고객 인구통계를 사전 처리하는 방법- Levenshtein 거리 (
EDIT_DISTANCE) 및 토큰 Jaccard 유사성을 사용하여 후보 차단을 실행하고 하이브리드 유사성 점수를 계산하는 방법 - 노드 및 에지 테이블에 BigQuery 속성 그래프 (
CREATE PROPERTY GRAPH)를 구성하는 방법 {1, 2}k-hop 수량자를 사용하여 ISO GQL (GRAPH_TABLE)로 그래프 경로를 쿼리하는 방법- 표준 고객 클러스터를 해결하고 정답 측정항목을 기준으로 모델 성능을 평가하는 방법
- 전체 데이터 세트를 재처리하지 않고 일일 배치 인테이크에 대해 증분 델타 매칭을 실행하는 방법
- 파이프라인 실행 전반에서 지속적인 클러스터 안정성을 유지하기 위해 (1 - ε) 중복 임계값 보장을 적용하는 방법
다음 단계
- BigQuery 속성 그래프 문서를 살펴봅니다.
- BigQuery 벡터 검색과 Vertex AI 텍스트 임베딩을 사용하여 시맨틱 후보를 생성해 보세요.
- 확장 가능한 외부 API 통합을 위한 BigQuery 원격 함수에 대해 알아보세요.