BigQuery 그래프를 사용한 고객 ID 해결

1. 소개

이 Codelab에서는 Google Cloud BigQuery 내에서 직접 모듈식 엔드 투 엔드 고객 ID 확인 (엔티티 매칭) 엔진을 빌드합니다. 인프라 배포를 위한 Google Cloud Shell과 데이터 정리, 후보 점수 매기기, 속성 그래프 구성, ISO GQL (그래프 쿼리 언어) 경로 순회를 위한 BigQuery Studio SQL 편집기를 결합합니다.

ID 해결은 엔터프라이즈 고객 360, 사기 감지, 다중 시스템 데이터 통합을 위한 기본 기능입니다. 데이터 성숙도와 비즈니스 요구사항에 따라 ID 해결에 유효한 접근 방식이 많기 때문에 이 Codelab의 모든 단계는 모듈식이며 선택사항입니다. 이 파이프라인은 원격 UDF 주소 정규화, Soundex 음성 차단, 의미론적 벡터 검색 (AI.EMBED), 하이브리드 기능 점수 매기기, GQL 속성 그래프 클러스터링 등 다양한 일반적인 프로덕션 등급 업계 기술을 보여주도록 설계되었으므로 아키텍처에 맞는 패턴을 선택적으로 채택할 수 있습니다.

일치 방법과 점수 기준점은 조직의 결정적 일치와 확률적 일치에 대한 요구에 따라 조정해야 하며, 이는 타겟 사용 사례에 따라 결정됩니다. 예를 들어 엄격한 규정 준수, 청구 또는 재무 운영에서는 잘못된 연결을 방지하기 위해 고정밀 결정론적 규칙 (예: 정확한 사회보장번호 또는 세금 ID 일치)을 선호하는 반면, 마케팅 맞춤설정, 분석, 추천 엔진에서는 재현율을 극대화하고 미묘한 연결을 발견하기 위해 확률적 퍼지 매칭 및 의미론적 벡터 유사성을 사용하는 경우가 많습니다.

BigQuery 고객 ID 해결 엔진 아키텍처

실습할 내용

  • FEBRL3 벤치마크 데이터 세트 수집: 합성 고객 레코드와 그라운드 트루스 일치 쌍을 BigQuery에 로드합니다.
  • 주소 유효성 검사 원격 UDF 배포: Python Cloud 함수를 배포하고 BigQuery 원격 함수를 등록하여 도로 주소를 정규화합니다.
  • 프로필 데이터 및 음성 인코딩 사전 처리: SQL 데이터 정리 실행, 주소 UDF 호출, SOUNDEX 음성 키 및 레벤슈타인 편집 거리 계산:
    • Soundex 음성 인코딩: 영어로 발음되는 소리를 기준으로 이름을 색인화하는 음성 알고리즘입니다. 이름을 자음 소리 그룹을 나타내는 4자리 코드 (이니셜과 세 자리 숫자)로 변환하여 (예: "John""Jon"은 모두 J500에 매핑되고 "Smith""Smyth"S530에 매핑됨) 기능 점수 매기기 및 실시간 증분 델타 차단을 위한 음성 일치 신호를 제공합니다.
    • Levenshtein 거리 (EDIT_DISTANCE): 한 문자열을 다른 문자열로 변경하는 데 필요한 최소 단일 문자 수정 (삽입, 삭제, 대체)을 측정하는 문자열 측정항목으로, 정확한 퍼지 이름 및 주소 일치를 지원합니다.
  • 시맨틱 프로필 임베딩 및 벡터 검색 생성: 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 프로젝트 만들기

  1. Google Cloud 콘솔의 프로젝트 선택기 페이지에서 Google Cloud 프로젝트를 선택하거나 만듭니다.
  2. Cloud 프로젝트에 결제가 사용 설정되어 있는지 확인합니다. 프로젝트에 결제가 사용 설정되어 있는지 확인하는 방법을 알아보세요.

Cloud Shell 시작

Cloud Shell은 Google Cloud에서 실행되는 명령줄 환경으로, 필요한 도구가 미리 로드되어 제공됩니다.

  1. Google Cloud 콘솔 상단에서 Cloud Shell 활성화를 클릭합니다.
  2. 인증을 확인합니다.
gcloud auth list
  1. Cloud Shell에서 환경 변수를 구성합니다.
export GCP_PROJECT=$(gcloud config get-value project)
export REGION="us-central1"
export DATASET_ID="identity_resolution"

필수 API 사용 설정

사용자 계정을 사용하여 Cloud Shell에서 다음 명령어를 실행하여 필요한 모든 Google Cloud 서비스를 사용 설정합니다.

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

원활한 API 실행과 애플리케이션 기본 사용자 인증 정보 (ADC) 액세스를 보장하려면 전용 실습 서비스 계정을 만들고 gcloud 가장을 사용 설정하세요.

# 1. Create a Service Account for the lab (if it does not already exist)
gcloud iam service-accounts create identity-res-sa \
  --display-name="Identity Resolution Service Account" 2>/dev/null || true

# Wait 5 seconds for IAM propagation
sleep 5

# Extract Project Number for default build and compute service accounts
export PROJECT_NUMBER=$(gcloud projects describe ${GCP_PROJECT} --format="value(projectNumber)")

# 2. Grant specific required least-privilege roles to the lab Service Account
for role in roles/bigquery.admin \
            roles/bigquery.resourceAdmin \
            roles/run.admin \
            roles/cloudfunctions.admin \
            roles/resourcemanager.projectIamAdmin \
            roles/cloudbuild.builds.editor \
            roles/cloudbuild.builds.builder \
            roles/artifactregistry.repoAdmin \
            roles/artifactregistry.writer \
            roles/storage.admin \
            roles/logging.logWriter \
            roles/iam.serviceAccountUser \
            roles/aiplatform.user; do
  gcloud projects add-iam-policy-binding ${GCP_PROJECT} \
    --member="serviceAccount:identity-res-sa@${GCP_PROJECT}.iam.gserviceaccount.com" \
    --role="${role}" --quiet
done

# 3. Grant required build & storage permissions to default Compute Engine & Cloud Build service accounts (required for 2nd-gen Cloud Functions container builds)
for role in roles/cloudbuild.builds.builder \
            roles/logging.logWriter \
            roles/artifactregistry.writer \
            roles/storage.objectAdmin; do
  gcloud projects add-iam-policy-binding ${GCP_PROJECT} \
    --member="serviceAccount:${PROJECT_NUMBER}-compute@developer.gserviceaccount.com" \
    --role="${role}" --quiet || true
  gcloud projects add-iam-policy-binding ${GCP_PROJECT} \
    --member="serviceAccount:${PROJECT_NUMBER}@cloudbuild.gserviceaccount.com" \
    --role="${role}" --quiet || true
done

# 4. Grant Service Account Token Creator role to your user account
export SA_EMAIL="identity-res-sa@${GCP_PROJECT}.iam.gserviceaccount.com"

gcloud iam service-accounts add-iam-policy-binding \
  "${SA_EMAIL}" \
  --member="user:$(gcloud config get-value account)" \
  --role="roles/iam.serviceAccountTokenCreator" --quiet

# 5. Enable Service Account impersonation for gcloud
gcloud config set auth/impersonate_service_account "${SA_EMAIL}"

# 6. Wait for IAM role assignments and impersonation caches to propagate
echo "Waiting 90 seconds for IAM policies and impersonation caches to propagate..."
sleep 90

BigQuery 데이터 세트 만들기

고객 노드, 에지, 그래프 모델, 평가 뷰를 저장할 BigQuery 데이터 세트를 만듭니다.

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

다음과 비슷한 출력이 표시됩니다.

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

주문형 CPU 한도나 공유 할당량에 제약받지 않고 벡터 색인 검색, 그래프 집계, 원격 함수 실행을 위한 전용 컴퓨팅 용량을 확보하려면 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

rec-10-org

brent

wood

11

girdlestone circuit

kingston tower

clifton springs

4152

nsw

19340706

1075870

febrl3

rec-10-dup-0

brnt

woode

11

girdelstone circut

cliffton springs

4152

nsw

19340706

1075870

febrl3

rec-10-dup-1

bernt

wood

11

girdlestone cir

kingston twr

clifton spngs

4152

19340760

1075870

febrl3

rec-10-dup-2

brent

wod

15

girdlestone crt

clifton springs

4152

nsw

19340706

febrl3

rec-25-org

mccarthy

henry

13

beasley street

crystal brook farm

ingleburn

6164

vic

19770913

1347524

febrl3

벤치마크 데이터 세트가 중복 클러스터에 현실적인 더티 데이터를 도입하는 방법을 확인하세요.

  • 발음 및 맞춤법 변형: brent vs. brnt / bernt, wood vs. woode / wod, clifton vs. cliffton
  • 주소 약어 및 오타 처리: girdlestone circuitgirdelstone circut / girdlestone cir / girdlestone crt, 번지수 11 대 OCR 오류 15
  • 문자 자리바꿈 및 누락된 값: 생년월일 자리바꿈 (19340706 vs. 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.pyrequirements.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를 만듭니다. 이 쿼리는 다음과 같은 작업을 합니다.

  1. validate_address_udf를 호출하여 정규화된 주소와 검증 결과를 가져옵니다.
  2. 맞춤법 변형을 처리하기 위해 given_namesurnameSOUNDEX 음성 인코딩을 생성합니다.
  3. 구조화된 profile_text 필드를 생성합니다.
CREATE OR REPLACE TABLE `identity_resolution.customer_nodes_cleaned` AS
WITH raw_data AS (
  SELECT 
    rec_id, dataset_source,
    TRIM(LOWER(given_name)) AS given_name_clean,
    TRIM(LOWER(surname)) AS surname_clean,
    `identity_resolution.validate_address_udf`(street_number, address_1, address_2, suburb, state, postcode) AS addr_json,
    TRIM(suburb) AS suburb, TRIM(state) AS state, TRIM(postcode) AS postcode,
    TRIM(date_of_birth) AS date_of_birth, TRIM(soc_sec_id) AS soc_sec_id
  FROM `identity_resolution.customer_nodes`
)
SELECT
  rec_id, dataset_source,
  given_name_clean AS given_name,
  SOUNDEX(given_name_clean) AS given_name_soundex,
  surname_clean AS surname,
  SOUNDEX(surname_clean) AS surname_soundex,
  CONCAT(given_name_clean, ' ', surname_clean) AS full_name,
  STRING(addr_json.formatted_address) AS formatted_address,
  BOOL(addr_json.address_is_valid) AS address_is_valid,
  STRING(addr_json.validation_granularity) AS validation_granularity,
  STRING(addr_json.possible_next_action) AS possible_next_action,
  suburb, state, postcode, date_of_birth, soc_sec_id,
  CONCAT('Name: ', CONCAT(given_name_clean, ' ', surname_clean), '; Address: ', STRING(addr_json.formatted_address), '; DOB: ', date_of_birth, '; SSN: ', soc_sec_id) AS profile_text
FROM raw_data;

정리된 노드 테이블을 쿼리합니다.

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

다음과 비슷한 출력이 표시됩니다.

rec_id

given_name

given_name_soundex

surname_soundex

formatted_address

rec-001-A

John

J500

Smith

S530

12 high st richmond vic 3121

rec-001-B

Jon

J500

Smith

S530

12 high st richmond vic 3121

rec-002-A

Elizabeth

E421

Taylor

T460

45 park rd suite 4 south yarra vic 3141

6. 시맨틱 프로필 임베딩 및 벡터 검색 생성

BigQuery는 주소 정규화, Soundex 음성 키, Levenshtein 편집 거리 외에도 AI.EMBED를 통해 기본 제공 생성형 AI 임베딩 함수를 지원합니다.

AI.EMBED를 사용하면 BigQuery에서 수동 벡터 색인 DDL 없이 기반 모델 (예: text-embedding-005)을 사용하여 SQL에서 직접 텍스트 삽입을 생성합니다.

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_idON 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

rec-001-A

rec-001-B

0.9400

rec-002-A

rec-002-B

0.9100

rec-003-A

rec-003-B

0.7300

규칙 기반 및 벡터 검색 에지를 통합 테이블로 병합

규칙 기반 퍼지 일치와 시맨틱 벡터 검색의 후보 가장자리를 중복이 삭제된 단일 final_matched_edges 테이블로 결합합니다.

CREATE OR REPLACE TABLE `identity_resolution.final_matched_edges` AS
SELECT 
  source_id, 
  target_id, 
  MAX(edge_weight) AS edge_weight,
  IF(COUNT(DISTINCT edge_type) > 1, 'HYBRID', MAX(edge_type)) AS edge_type
FROM (
  SELECT source_id, target_id, match_score AS edge_weight, 'RULE_BASED' AS edge_type
  FROM `identity_resolution.matched_edges`
  UNION ALL
  SELECT source_id, target_id, ROUND(1.0 - vector_distance, 4) AS edge_weight, 'VECTOR_SEARCH' AS edge_type
  FROM `identity_resolution.vector_candidate_edges`
  WHERE vector_distance <= 0.05
)
GROUP BY source_id, target_id;

8. 속성 그래프 구성 및 ISO GQL 경로 순회

BigQuery는 속성 그래프를 통해 ISO GQL (Graph Query Language)을 기본적으로 지원합니다. 속성 그래프는 데이터 중복 없이 관계형 BigQuery 테이블에 논리적 그래프 뷰를 만듭니다.

이 단계에서는 통합 후보 에지 테이블 (final_matched_edges)을 사용하여 속성 그래프 customer_identity_graph를 빌드하고 {1, 2} k-hop 경로 순회를 사용하여 고객 프로필 전반에서 그래프 연결을 쿼리합니다.

전이적 연결 및 K-홉 순회

BigQuery 속성 그래프 DDL 만들기

BigQuery Studio SQL 편집기에서 다음 DDL 문을 실행합니다.

CREATE OR REPLACE PROPERTY GRAPH `identity_resolution.customer_identity_graph`
NODE TABLES (
  `identity_resolution.customer_nodes_cleaned` AS `Customer`
  KEY (rec_id)
)
EDGE TABLES (
  `identity_resolution.final_matched_edges`
  KEY (source_id, target_id)
  SOURCE KEY (source_id) REFERENCES `Customer`(rec_id)
  DESTINATION KEY (target_id) REFERENCES `Customer`(rec_id)
  LABEL MATCHED_TO
);

K-Hop 그래프 클러스터 시각화

아래 쿼리를 실행하여 1~2 관계 홉에서 일치하는 고객 클러스터를 시각화합니다.

GRAPH `identity_resolution.customer_identity_graph`
MATCH p = (c1:Customer)-[e:MATCHED_TO]->{1, 2}(c2:Customer)
RETURN TO_JSON(p) AS graph_cluster_path
LIMIT 10;

K-Hop 그래프 클러스터 시각화

표준 고객 클러스터 해결

다음 쿼리를 실행하여 엔티티 클러스터를 resolved_customers로 변환합니다.

CREATE OR REPLACE TABLE `identity_resolution.resolved_customers` AS
WITH graph_paths AS (
  SELECT 
    source_node_id, target_node_id
  FROM GRAPH_TABLE(
    `identity_resolution.customer_identity_graph`
    MATCH (c1:Customer)-[e:MATCHED_TO]->{1, 2}(c2:Customer)
    COLUMNS (c1.rec_id AS source_node_id, c2.rec_id AS target_node_id)
  )
),
all_connections AS (
  SELECT source_node_id AS node_id, target_node_id AS connected_id FROM graph_paths
  UNION DISTINCT
  SELECT target_node_id AS node_id, source_node_id AS connected_id FROM graph_paths
  UNION DISTINCT
  SELECT rec_id AS node_id, rec_id AS connected_id FROM `identity_resolution.customer_nodes_cleaned`
),
clusters AS (
  SELECT 
    node_id,
    MIN(connected_id) AS canonical_customer_id
  FROM all_connections
  GROUP BY node_id
)
SELECT 
  canonical_customer_id,
  ARRAY_AGG(node_id) AS customer_records,
  COUNT(node_id) AS record_count
FROM clusters
GROUP BY canonical_customer_id;

해결된 클러스터 테이블을 쿼리합니다.

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

다음과 비슷한 출력이 표시됩니다.

canonical_customer_id

record_count

customer_records

rec-001-A

2

['rec-001-A', 'rec-001-B']

rec-002-A

2

['rec-002-A', 'rec-002-B']

rec-003-A

2

['rec-003-A', 'rec-003-B']

평가 측정항목 뷰 만들기

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

6538

5620

5608

13

930

0.9977

0.8579

0.9225

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

rec-9999-new-1

rec-529-dup-0

hh-8745613133408211212

MATCHED_TO_EXISTING_CLUSTER

1.0000

BOTH

rec-9999-new-2

rec-875-dup-0

hh-4802217335228845918

MATCHED_TO_EXISTING_CLUSTER

1.0000

BOTH

rec-9999-new-3

rec-359-dup-0

hh-8891673998566910207

MATCHED_TO_EXISTING_CLUSTER

0.8739

VECTOR_SEARCH

rec-9999-new-4

197e02f9-7175-4484-9c32-b7edbc731c9e

NULL

NEW_CUSTOMER_ENTITY

NULL

NONE

10. 클러스터링 통합 및 클러스터 안정성 (1-ε 중복)

프로덕션 엔터프라이즈 시스템에서 팀은 일반적으로 새로운 에지와 데이터 소스를 통합하기 위해 반복적으로 (예: 주별 또는 월별) 전체 그래프를 재클러스터링합니다. 새로운 관계가 형성되면 전체 그래프 재클러스터링으로 인해 파이프라인 실행 전반에서 클러스터 식별자가 임의로 이동하거나 뒤집힐 수 있습니다.

다운스트림 CRM, CDP, 결제 시스템의 지속적인 고객 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

rec-529-dup-0

rec-9999-new-1

hh-8745613133408211212

STABLE_EVOLUTION

rec-875-dup-0

rec-9999-new-2

hh-4802217335228845918

STABLE_EVOLUTION

rec-359-dup-0

rec-9999-new-3

hh-8891673998566910207

STABLE_EVOLUTION

589f8102-1204-4530-8910-bc10294810a4

rec-9999-new-4

NULL

NEW_CLUSTER_CREATED

11. 삭제

Google Cloud 계정에 지속적으로 비용이 청구되지 않도록 배포된 리소스와 BigQuery 데이터 세트를 정리합니다.

Cloud Shell에서 다음을 실행합니다.

# 1. Delete BigQuery Dataset
bq rm -r -f -d ${GCP_PROJECT}:${DATASET_ID}

# 2. Delete Cloud Function (2nd-Gen)
gcloud functions delete validate_address_udf --region=${REGION} --gen2 --quiet

# 3. Delete BigQuery Cloud Connection
bq rm -f --connection US.address_val_conn

이 실습을 위해 전용 Google Cloud 프로젝트를 만든 경우 프로젝트를 삭제할 수 있습니다.

gcloud projects delete ${GCP_PROJECT}

12. 축하합니다

수고하셨습니다 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 - ε) 중복 임계값 보장을 적용하는 방법

다음 단계

참조 문서