워크드 예제: 완성된 계획
NYC taxi 워크로드에 대한 작성 완료된 마이그레이션 계획으로, 여러분이 직접 작성한 계획과 비교해 보세요.
다섯 개 워크시트
(1,
2,
3,
4,
5) 전체를 NYC
Taxi 워크로드에 적용해 완전히 풀어낸 답안입니다. 다음 용도로 사용하세요:
- 각 섹션을 완료한 뒤 워크시트 답을 확인
- Part 3이 구현하는 결정들의 근거 이해
- 다르게 선택했다면 Part 3의 Decision Alignment 표와 비교
이것은 정답지입니다 — 여러분의 계획으로 여기에 작성하지 마세요. 대신 migration-plan.md에 작성하세요.
| 지표 | 값 |
|---|
| 전체 테이블 수 | 7 (TRIPS_RAW, FACT_TRIPS, AGG_HOURLY_ZONE_TRIPS, DIM_TAXI_ZONES, DIM_PAYMENT_TYPE, DIM_VENDOR, DIM_DATE) |
| 전체 뷰 수 | 2 (STG_TRIPS, STG_TAXI_ZONES) |
| 스트림 | 1 (TRIPS_RAW의 TRIPS_CDC_STREAM) |
| 태스크 | 2 (CDC_CONSUME_TASK, HOURLY_AGG_TASK) |
| TRIPS_RAW의 전체 행 수 | ~50,000,000 |
| 날짜 범위 | 설정 시점에서 끝나는 4년 롤링 윈도우 |
| VARIANT 컬럼 | 1 (TRIPS_RAW.TRIP_METADATA) |
| 탐지된 QUALIFY 사용 | 1 (Q3 쿼리) |
| 탐지된 MERGE INTO 사용 | 2 (HOURLY_AGG_TASK, CDC_CONSUME_TASK) |
| 객체 | 유형 | 스키마 | 행 수 | 복잡도 등급 | 비고 |
|---|
trips_raw | 테이블 | raw | ~50M | B | _synced_at 버전 컬럼을 가진 RMT; 벌크 + CDC 중첩으로 중복 제거 필요; stg_trips는 FINAL을 사용해야 함 |
stg_trips | dbt 뷰 | staging | — | B | TRIP_METADATA용 JSONExtract; JSON 경로 테스트 필요 |
stg_taxi_zones | dbt 뷰 | staging | — | A | 단순 전달; 간단함 |
int_trips_enriched | dbt Ephemeral | staging | — | A | CTE; SQL 차이는 상위 모델에서 처리 |
fact_trips | dbt Incremental | analytics | ~50M | C | RMT 엔진; delete_insert; QUALIFY 재작성; FINAL 필요 |
agg_hourly_zone_trips | dbt Incremental | analytics | ~140K | B | RMT; 롤링 2시간 재계산 윈도우; 파티션 경계를 신중히 테스트 |
dim_taxi_zones | dbt 테이블 | analytics | 265 | A | 정적 참조 데이터; 전체 재적재; 간단함 |
dim_payment_type | dbt 테이블 | analytics | 6 | A | 정적 참조 데이터; 간단함 |
dim_vendor | dbt 테이블 | analytics | 3 | A | 정적 참조 데이터; 간단함 |
taxi_zones_dict | 딕셔너리 | analytics | 265 | B | ClickHouse 고유 문법; 쿼리 시점의 dictGet() |
mv_hourly_revenue | Refreshable MV | analytics | — | B | REFRESH EVERY 문법; 원자적 교체 확인 |
TRIPS_CDC_STREAM / CDC_CONSUME_TASK | Snowflake Stream + Task | — | — | D | ClickHouse에 대응물 없음; Part 3에서 프로듀서 직접 전환으로 대체 |
| 테이블 | 엔진 | 버전 컬럼 | 근거 |
|---|
trips_raw | ReplacingMergeTree(_synced_at) | _synced_at | Python 마이그레이션 스크립트(scripts/02_migrate_trips.py)가 배치를 재시도해 같은 trip_id를 다시 삽입할 수 있습니다. 전환 후에는 실시간 프로듀서도 일시적 실패 시 재시도할 수 있습니다. _synced_at DateTime DEFAULT now()는 INSERT 시 자동으로 설정되므로 — 나중의 재시도는 더 큰 타임스탬프를 가지며, RMT는 가장 최근 쓰기를 유지합니다. stg_trips는 FINAL로 쿼리해 다운스트림 모델이 실행되기 전에 중복 제거를 강제합니다. |
fact_trips | ReplacingMergeTree(updated_at) | updated_at | 트립은 정정될 수 있습니다(요금 조정, 상태 변경). 같은 trip_id가 갱신된 값과 함께 다시 삽입됩니다. updated_at은 정정마다 단조 증가하므로 — RMT 중복 제거 시 더 큰 값이 이깁니다. 항상 FINAL로 쿼리하세요. |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | dbt가 최근 2시간을 재계산해 다시 삽입합니다. RMT가 없으면 기존 집계와 새 집계가 누적되어 이중 계산됩니다. dbt 실행마다 updated_at을 now()로 설정하면 최신 값이 이깁니다. |
dim_taxi_zones | MergeTree() | — | dbt에 의한 전체 재적재(원자적 테이블 교체(전체 재구축)). 중복이 누적될 수 없습니다. 중복 제거 불필요. |
dim_payment_type | MergeTree() | — | 동일 — 전체 재적재. |
dim_vendor | MergeTree() | — | 동일 — 전체 재적재. |
mv_hourly_revenue | MergeTree() | — | REFRESHABLE MV는 REFRESH마다 결과 집합 전체를 원자적으로 교체합니다. upsert 없음. |
| 테이블 | ORDER BY | 근거 |
|---|
trips_raw | (pickup_at, trip_id) | 시간 범위 스캔은 pickup_at을 먼저 필터링합니다. trip_id는 RMT 중복 제거 키이므로 — RMT가 어떤 행이 중복인지 식별할 수 있도록 ORDER BY에 있어야 합니다. 분석 범위 스캔이 지배적이므로 pickup_at을 앞에, trip_id는 카디널리티가 높고 유일성 판별자 역할만 하므로 뒤에 둡니다. |
fact_trips | (toStartOfMonth(pickup_at), pickup_at, trip_id) | 7개의 분석 쿼리 모두 pickup_at을 필터링합니다. 월 접두사는 달력상 월 데이터를 인접 블록으로 묶어 — PARTITION BY를 추가하지 않고도 월별 집계에서 대략적인 블록 스킵을 가능하게 합니다. trip_id는 블록 스킵을 방해하지 않으면서 RMT 유일성을 위해 마지막에 둡니다. |
agg_hourly_zone_trips | (hour_bucket, zone_id) | Q6(그리고 모든 집계 쿼리)은 hour_bucket과 zone_id를 필터링합니다. hour_bucket은 약 35K개의 서로 다른 값을, zone_id는 265개를 가집니다. 시간 범위 스캔이 주된 접근 패턴이므로 hour_bucket을 먼저, 2차 필터링을 위해 zone_id를 두 번째에 둡니다. |
dim_taxi_zones | (location_id) | 265행 = 하나의 그래뉼. ORDER BY는 성능에 무관합니다. 조인 키로서 location_id는 관례적이고 가독성에 도움이 됩니다. |
| 컬럼 | Snowflake 타입 | ClickHouse 타입 | 결정 근거 |
|---|
TRIP_METADATA | VARIANT | String | 원시 JSON을 정확히 보존합니다. JSONExtract*가 쿼리 시점에 임의 경로를 처리합니다. Map(String,String)은 중첩 구조를 잃고, Tuple은 고정 스키마를 요구합니다. 임의 JSON에는 String이 안전한 선택입니다. |
PICKUP_DATETIME / PICKUP_AT | TIMESTAMP_NTZ(9) | DateTime64(3, 'UTC') | 트립 타임스탬프에는 밀리초 정밀도가 충분합니다. 나노초(9)는 과도합니다. 'UTC'는 타임존을 명시해 시간 범위 집계에서 DST 관련 문제를 피합니다. |
PICKUP_LOCATION_ID | INTEGER | UInt16 | 값은 1–265. UInt8 최대는 255(너무 작음). UInt16 최대는 65535(적절). Int32의 4바이트 대비 2바이트 — 5천만 행에서 컬럼당 비압축 약 95MB를 절약합니다. |
VENDOR_ID | INTEGER | UInt8 | 값은 1–3. UInt8 최대는 255 — 적절합니다. 행당 1바이트. |
DRIVER_RATING | FLOAT | Nullable(Float32) | NULL이 자주 발생합니다(모든 트립에 평점이 있는 것은 아님). Nullable이 올바른 null 의미론을 보존합니다. Float32는 1.0–5.0 범위에 충분합니다. Float64는 의미 있는 정밀도 없이 스토리지만 낭비합니다. |
UPDATED_AT | TIMESTAMP_NTZ(9) | DateTime64(3, 'UTC') | ReplacingMergeTree의 버전 컬럼. DateTime이 아니라 DateTime64를 사용해야 합니다 — 초 정밀도라면 같은 초 안의 두 정정이 비결정적이 됩니다. 밀리초 정밀도가 올바른 중복 제거 순서를 보장합니다. |
| Snowflake 표현식 | ClickHouse 등가물 |
|---|
DATE_TRUNC('hour', pickup_at) | toStartOfHour(pickup_at) |
DATEADD('day', -7, CURRENT_DATE) | today() - 7 |
DATEDIFF('minute', pickup_at, dropoff_at) | dateDiff('minute', pickup_at, dropoff_at) |
TRIP_METADATA:driver.rating::FLOAT | JSONExtractFloat(trip_metadata, 'driver', 'rating') |
TRIP_METADATA:surge_multiplier::FLOAT | JSONExtractFloat(trip_metadata, 'surge_multiplier') |
QUALIFY ROW_NUMBER() OVER (PARTITION BY pickup_location_id ORDER BY fare_amount DESC) <= 10 | SELECT ... FROM (SELECT ..., ROW_NUMBER() OVER (...) AS rn FROM ...) WHERE rn <= 10 |
MERGE INTO fact_trips ... WHEN MATCHED THEN UPDATE | dbt delete_insert 증분 — 키가 일치하는 행을 DELETE한 뒤 모든 새 행을 INSERT |
| 웨이브 | 객체 | 의존성 | 비고 |
|---|
| Wave 0 | trips_raw(스키마), dim_taxi_zones, dim_payment_type, dim_vendor | 없음 | dbt가 빈 테이블을 생성합니다. 차원 테이블은 정적 참조 데이터로 즉시 채워집니다(트립에 의존하지 않음). 실행: dbt run --select trips_raw dim_* |
| Wave 1 | Python 벌크 로드 (scripts/02_migrate_trips.py) | Wave 0 (trips_raw 스키마가 존재해야 함) | Snowflake TRIPS_RAW에서 5천만 행. --resume으로 재개 가능. scripts/01_verify_migration.sh로 행 수를 검증하세요. 약 40-50분. |
| Wave 2 | stg_trips, stg_taxi_zones, int_trips_enriched, fact_trips, agg_hourly_zone_trips | Wave 1 완료(trips_raw 채워짐) + Wave 0(차원 테이블 존재) | 전체 dbt run. stg_trips는 trips_raw를 읽고, int_trips_enriched는 차원과 조인하며, fact_trips와 agg_hourly_zone_trips가 그 위에 만들어집니다. |
| Wave 3 | taxi_zones_dict, mv_live_trip_feed | Wave 2(딕셔너리를 위해 dim_taxi_zones 채워짐; MV를 위해 fact_trips 채워짐) | 딕셔너리는 scripts/04_create_dictionary.sql로 생성됩니다. Refreshable MV는 dbt 모델로 생성되며, 이후 갱신 주기를 활성화하는 것은 수동 ALTER TABLE ... MODIFY REFRESH 단계이지 dbt가 자동으로 실행하는 것이 아닙니다. |
| Wave 4 | 프로듀서 전환 (scripts/03_cutover.sh) | Wave 1 완료(벌크 로드 검증) + Wave 2 완료(analytics 계층 구축) | Snowflake 프로듀서를 중지하고, ClickHouse Cloud에 직접 쓰는 ClickHouse 프로듀서를 시작한 뒤, dbt run으로 실시간 데이터로 agg_hourly_zone_trips를 채웁니다. |
| 객체 | 리스크 | 검증 방법 |
|---|
fact_trips | FINAL 없는 쿼리는 병합 지연 동안 과다 집계합니다. delete_insert 파티션 범위는 대상이 아닌 파티션을 삭제하지 않도록 ORDER BY 접두사에 맞춰 한정해야 합니다. | SELECT COUNT(*) FINAL이 Snowflake와 ± CDC 지연 범위에서 일치합니다. dbt test를 실행하세요. 두 시스템 간 Q3 결과를 비교하세요. 중복 trip_id를 확인하세요: SELECT trip_id, count() FROM fact_trips GROUP BY trip_id HAVING count() > 1 LIMIT 10. |
agg_hourly_zone_trips | 롤링 2시간 재계산 윈도우는 삭제 범위를 정확히 한정해야 합니다. 너무 넓으면 기존 집계가 삭제되고, 너무 좁으면 오래된 집계가 남습니다. | 특정 (hour_bucket, zone_id) 조합을 Snowflake와 표본 비교하세요. 같은 기간에 대해 모든 존의 총 trip_count가 Snowflake AGG_HOURLY_ZONE_TRIPS와 일치하는지 확인하세요. |
| 프로듀서 전환 | 마이그레이션 스크립트가 실행 중 중단되면 행 수 격차가 남습니다; --resume으로 재실행해 채우세요. 전환 후 프로듀서 재시도로 이미 ClickHouse에 있는 트립이 다시 삽입될 수 있습니다. | scripts/01_verify_migration.sh — Snowflake와 ClickHouse 간 행 수 일치를 확인합니다. ReplacingMergeTree(_synced_at)가 중복 삽입을 멱등하게 처리합니다. |
데이터 이동: Python 마이그레이션 스크립트 (scripts/02_migrate_trips.py)
객체 스토리지 릴레이나 ClickPipes 대신 Python 스크립트를 쓰는 이유는?
remoteSecure()는 ClickHouse 간 데이터 전송용입니다 — 여기에는 해당되지 않습니다.
- 객체 스토리지 릴레이(Snowflake → S3 → ClickHouse S3 테이블 함수)도 동작하지만 복잡도를 더합니다: S3 버킷 프로비저닝, IAM 역할, Snowflake COPY INTO가 필요합니다 — 랩에는 불필요한 오버헤드입니다.
- ClickPipes는 Snowflake를 소스로 지원하지 않습니다. 지원 소스는 Kafka, S3, Kinesis, PostgreSQL CDC, MySQL CDC입니다.
- Python 스크립트는
snowflake-connector-python과 clickhouse-connect를 사용하며 — 랩에 이미 설치된 패키지입니다. 진행 상황을 실시간으로 보여주고, 중단 시 --resume을 지원하며, 코드를 전부 들여다볼 수 있습니다.
증분 전략 (dbt): delete_insert
append나 merge 전략 대신 delete_insert를 쓰는 이유는?
append는 기존 행을 건드리지 않고 새 행을 삽입합니다. 행이 갱신될 수 있는 fact_trips에서는 중복을 만듭니다. 잘못된 선택입니다.
merge(가능한 경우)는 Snowflake의 MERGE INTO에 가장 가깝지만, dbt-clickhouse의 merge 전략은 ReplacingMergeTree와 함께 쓸 때 제약이 있고 권장되는 접근이 아닙니다.
delete_insert는 들어오는 배치의 키 범위에 해당하는 행을 삭제한 뒤 모든 새 행을 삽입합니다. 멱등적이고(재실행해도 같은 결과), 삽입과 갱신을 모두 처리하며, ReplacingMergeTree와 올바르게 동작합니다. upsert 패턴에 대한 dbt-clickhouse 커뮤니티의 표준 권고입니다.
| 기준 | 임계값 | 측정 방법 |
|---|
| 행 수 일치 | ≥ 99.9% 일치 (전환 후 CH ≥ SF는 예상된 결과) | scripts/01_verify_migration.sh |
| 체크섬 일치 | 1만 행 샘플에서 MD5 일치 | scripts/02_validate_parity.sql |
| dbt test 통과율 | 100% | dbt/nyc_taxi_dbt_ch에서 dbt test |
| 쿼리 결과 일치 | 7개 쿼리 모두 동일한 결과 반환(부동소수점 허용 오차 내) | scripts/run_benchmark.sh 출력을 수동 비교 |
| 모델 | Materialization | 이유 |
|---|
stg_trips | view | trips_raw를 읽고 정리; 이 모델에 갱신 없음; 스토리지 비용 0; 항상 현재 소스 상태를 반영 |
stg_taxi_zones | view | 동일 — 소스 테이블의 단순 정리 전달 |
int_trips_enriched | ephemeral | fact_trips만 사용하는 순수 조인 로직; CTE로 인라인되어 불필요한 물리 테이블을 피함; 이 모델을 직접 쿼리하는 모델 없음 |
fact_trips | incremental | 트립은 사후에 정정될 수 있음; 실행마다 신규 및 갱신된 행만 처리해야 함 |
agg_hourly_zone_trips | incremental | 롤링 2시간 재계산은 증분 패턴 — 5천만 행 전체가 아니라 최근 행만 처리 |
dim_taxi_zones | table | 265개 정적 존; 원자적 테이블 교체(전체 재구축)로 dbt 실행마다 전체 재구축; 부분 갱신 없음 |
dim_payment_type | table | 6개 정적 유형; dim_taxi_zones와 동일한 근거 |
dim_vendor | table | 3개 벤더; 동일한 근거 |
| 모델 | ENGINE | 버전 컬럼 | 이유 |
|---|
fact_trips | ReplacingMergeTree(updated_at) | updated_at | 트립은 정정될 수 있음; 삽입마다 updated_at을 now()로 설정하면 RMT 백그라운드 중복 제거 시 최신 버전이 이김; delete_insert가 주된 정확성 경로이고 RMT는 안전망 |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | 롤링 재계산이 동일한 (hour_bucket, zone_id) 쌍에 대해 집계를 다시 삽입함; RMT가 백그라운드 병합 시 오래된 집계를 제거해 줌 |
dim_taxi_zones | MergeTree() | — | dbt의 전체 재적재는 실행마다 원자적 테이블 교체(전체 재구축)를 의미함; 중복이 누적될 수 없음; 중복 제거 불필요 |
dim_payment_type | MergeTree() | — | dim_taxi_zones와 동일 |
dim_vendor | MergeTree() | — | dim_taxi_zones와 동일 |
| 모델 | unique_key | incremental_strategy | 증분 필터 | 이 필터를 쓰는 이유 |
|---|
fact_trips | trip_id | delete_insert | WHERE updated_at > (SELECT max(updated_at) FROM {{ this }}) | updated_at 기준 하이워터마크는 새 트립과 정정된 트립을 모두 포착합니다(요금 조정은 같은 pickup_at과 같은 trip_id에 더 새로운 updated_at으로 다시 삽입됨); pickup_at 워터마크는 정정을 조용히 놓칩니다 |
agg_hourly_zone_trips | [hour_bucket, zone_id] | delete_insert | WHERE pickup_at >= now() - INTERVAL 2 HOUR | 롤링 2시간 윈도우는 경계 시간대의 재집계를 강제해 부분 시간대 카운트가 항상 정정되게 합니다; max(pickup_at) 하이워터마크는 경계 시간대를 영구적으로 과소 집계합니다 |
| 모델 | FROM 절에 FINAL? | 이유 |
|---|
stg_trips | 예 — FROM trips_raw FINAL | trips_raw는 ReplacingMergeTree이며, 마이그레이션 스크립트 재시도나 전환 후 프로듀서 재시도로 중복 trip_id 행이 생길 수 있습니다. stg_trips가 단일 시행 지점입니다: 여기서 중복을 제거해 모든 다운스트림 모델(int_trips_enriched, fact_trips, agg_hourly_zone_trips)이 깨끗한 데이터를 받게 합니다 |
int_trips_enriched | 아니요 | RMT 테이블이 아닌 stg_trips(뷰)에서 읽음; 뷰에는 FINAL이 무관함 |
fact_trips | 아니요 (모델 본문에서) | delete_insert가 완료된 각 실행 후 fact_trips를 깨끗하게 유지합니다; 모델 안에 FINAL을 넣으면 {{ this }}에서 max(updated_at)을 읽는 is_incremental() 서브쿼리에까지 낭비되게 적용됩니다. 대시보드와 dbt 테스트는 fact_trips를 직접 쿼리할 때 외부에서 FINAL을 사용합니다 |
이것은 완성된 예제입니다. 여러분의 migration-plan.md는 여기의 핵심 결정과 일치해야 합니다 — 또는 다르게 선택한 이유를 명시적으로 문서화해야 합니다.