03 프로비저닝과 마이그레이션
Terraform으로 ClickHouse Cloud를 프로비저닝하고, 계획에 따라 타깃 테이블을 만들고, 재개 가능한 Python 마이그레이션 스크립트로 5천만 행을 옮긴다.
시작 지점
모듈 02 완료 상태: migration-plan.md가 채워져 있고 Completion Checklist의 모든
체크박스가 체크되어 있으며, Snowflake 프로듀서가 여전히 실행 중이다. 이 모듈의
setup.sh는 그 파일을 확인하고 없거나 미완성이면 경고하지만, 절대 막지는 않는다 —
여기서 그것 없이 계속하는 것을 막는 것은 아무것도 없고, 다음 두 모듈에 대한 본인의
이해만이 막는다. 총 약 60분을 예상하라. 그중 약 40-50분은 백그라운드에 두고 지켜보지
않아도 되는 데이터 전송이다. 여기서부터 ClickHouse Cloud 트라이얼 지출이 시작된다.
서비스를 프로비저닝하고 이 모듈을 진행하면 트라이얼 크레딧 약 $1-2가 든다(랩 전체는
총 약 $2-4).
이유
이 모듈은 계획이 실물이 되는 곳이다. 모듈 02에서 migration-plan.md에 써넣은 모든
결정 — 테이블별 MergeTree 엔진, 실제 쿼리 워크로드에서 도출한 ORDER BY 키, Snowflake
전용 구문의 변환 방식 — 이 여기서 다시 도출되는 것이 아니라 테이블 DDL에 그대로 입력된다.
ClickHouse에는 나중에 덧붙일 수 있는 인덱스가 없다. 5천만 행이 테이블에 들어앉은 뒤에
ORDER BY 키가 틀렸다는 것이 드러나면 해결책은 빠른 ALTER가 아니라 전체 재적재다.
그래서 소프트 게이트가 사람을 막을 수 없더라도 중요하다. 완성된 계획 없이 이 모듈을
실행하면 기계적으로는 성공한다 — dbt run은 여전히 fact_trips를 ReplacingMergeTree로
만들고, 마이그레이션 스크립트는 여전히 5천만 행을 옮긴다 — 그러나 왜 평범한 MergeTree가
아니라 그 엔진인지, 왜 정렬 키가 그런 형태인지, 모듈 04가 나중에 보여 주는 약 6-9배의
벤치마크 속도 향상을 어떻게 정당화할지는 알지 못한다. 아래의 결정 정렬 표는 이 모듈이
구현하는 모든 선택을 그것이 답하는 워크시트 문제로 되짚어 매핑하므로, 무엇이든 프로비저닝
하기 전에 자신의 계획과 대조해 볼 수 있다.
개념 — 내부 동작
타깃 아키텍처. 일회성 Python 스크립트가 기존 5천만 행을 ClickHouse로 백필하는 동안에도 Snowflake는 트립 프로듀서를 통해 새 트립을 계속 쓴다 — 두 시스템은 컷오버가 아니라 마이그레이션 기간 내내 나란히 동작한다.
ClickHouse 쪽에서 trips_raw는 스크립트가 쓰는 랜딩 테이블이다. 그다음 dbt가 그 위에
staging 뷰와 나머지 analytics 레이어를 만든다 — 이 모듈이 만들지만 trips_raw 외에는
아직 채우지 않는 스키마다.
다이어그램 색상 범례:
- 초록 — 데이터 소스(컷오버 전후의 트립 프로듀서)
- 파랑 — Snowflake 테이블
- 주황 — dbt 모델과 파이프라인
- 빨강 — ClickHouse 테이블과 materialized view
- 시안 — Apache Superset 대시보드
- 점선 화살표 — 컷오버 이후의 흐름
네이티브 커넥터가 아니라 Python 스크립트를 쓰는 이유. Snowflake에서 ClickHouse로 데이터를 옮기는 방법은 여러 가지가 있다. 이 랩은 Python 배치 스크립트를 사용한다 — 대안들과 비교한 이유는 다음과 같다.
| 방법 | 동작 방식 | 여기서 쓰지 않는 이유 |
|---|---|---|
| ClickPipes (Snowflake 소스) | 네이티브 ClickHouse Cloud 커넥터 — zero-ETL, 관리형 UI | Snowflake는 지원되는 ClickPipes 소스가 아니다. ClickPipes는 Kafka, S3, Kinesis, PostgreSQL CDC, MySQL CDC, 객체 스토리지를 지원한다. |
| S3 익스포트 → ClickPipes S3 | COPY INTO @stage로 Parquet/CSV를 S3에 내보내고, ClickPipes S3 커넥터가 이를 ClickHouse로 적재한다 | S3 버킷, IAM 역할, Snowflake 스테이지, AWS 계정이 필요하다. 데이터가 조금이라도 움직이기 전에 설정 단계가 약 3개 늘어난다. 프로덕션에서는 쓸 만하지만 랩에는 인프라가 너무 많다. |
S3 익스포트 → clickhouse-client | 같은 S3 익스포트지만 INSERT INTO ... SELECT FROM s3(...)로 적재한다 | 같은 S3 사전 요건이 필요하다. 게다가 파트너가 파일 청킹과 재개 가능성을 직접 관리해야 한다. |
| Snowflake → Kafka → ClickHouse | Snowflake CDC 스트림이 Kafka 토픽에 공급하고, ClickPipes Kafka 커넥터가 인제스트한다 | 완전한 스트리밍 파이프라인 — 프로덕션에서 1분 미만 지연 요구가 있을 때 적합하다. Kafka 클러스터는 랩 환경에 지나치게 무겁다. |
| Python 스크립트(이 랩) | snowflake-connector-python이 10만 행 커서 배치로 읽고, clickhouse-connect가 직접 삽입한다 | 랩이 이미 필요로 하는 패키지 외에 추가 인프라가 전혀 없다. --resume(max(pickup_at) 워터마크)으로 재개 가능하다. 실시간 진행 출력이 있다. 초당 약 2만 행으로 5천만 행에 약 40-50분 — 일회성 마이그레이션 실습에는 충분하다. |
이 랩에서 Python 스크립트가 옳은 선택인 이유:
- AWS 계정이 필요 없다. S3 기반 접근법은 버킷 생성, IAM 정책, Snowflake 외부 스테이지를 요구한다 — ClickHouse와 아무 상관 없는 세 개의 설정 단계다.
- 자체 완결적이다. 두 패키지(
snowflake-connector-python,clickhouse-connect)는 dbt와 같은 venv에 설치된다. 새 서비스도, 새 자격 증명도 없다. - 재개 가능하다.
--resume덕분에 스크립트를 중단하고 재시작해도 안전하다.ReplacingMergeTree(_synced_at)가 재시도 시의 중복 삽입을 자동으로 중복 제거해 준다. - 투명하다. 파트너가 스크립트를 읽고, 컬럼 매핑을 이해하고, 자기 스키마에 맞게 고칠 수 있다 — UI 위저드를 클릭해 넘어가는 것보다 교육적이다.
마이그레이션 공백 처리. 마이그레이션 스크립트가 도는 동안(약 40-50분) Snowflake 프로듀서는 계속 실행된다. 그 구간에 Snowflake에 쓰인 트립은 ClickHouse에 없다. 이 랩은 컷오버 시점의 2패스 방식으로 그 공백을 닫으며, 모듈 05가 이를 직접 안내한다.
- Snowflake 프로듀서를 중지해 데이터셋을 고정한다.
python scripts/02_migrate_trips.py --resume를 실행한다 — 델타 행만 전송된다(분이 아니라 초 단위).- ClickHouse 프로듀서를 시작한다.
마이그레이션 재시도를 처리하는 것과 동일한 ReplacingMergeTree(_synced_at) 중복 제거가
이것도 처리한다. 이 모듈의 실행과 나중의 --resume 패스 사이에 겹치는 행이 있다면 나중의
_synced_at이 이긴다.
프로덕션에서 S3을 선택할 상황. 데이터셋이 5억 행을 넘거나, 전체 테이블 스캔의 Snowflake 웨어하우스 쿼리 비용이 상당하다면 S3 익스포트 경로가 더 낫다. Snowflake가 압축된 Parquet을 병렬로 내보내고(단일 커서보다 훨씬 빠르다), ClickHouse도 S3에서 병렬로 적재할 수 있다. 여기서 쓰는 Python 스크립트 방식은 랩 규모에서 잘 동작한다.
결정 정렬. 아래 표는 migration-plan.md의 워크시트 1(엔진 선택), 2(정렬 키),
3(스키마 변환)에 있는 것과 동일한 결정 목록을, 이 랩이 실제로 구축하는 것과 교차
확인한 것이다 — 무엇이든 프로비저닝하기 전에 자신의 계획과 대조하라.
| 결정 | 이 랩의 구현 | 이유 |
|---|---|---|
trips_raw 엔진 | ReplacingMergeTree(_synced_at) | Python 마이그레이션 스크립트는 중단 시 재시도될 수 있는 배치 INSERT를 사용한다. _synced_at DateTime DEFAULT now()가 모든 INSERT에 설정되므로 재시도된 행은 더 늦게 도착해 _synced_at 값이 더 크다 — RMT 중복 제거 시 나중 행이 이기므로 재시도가 멱등해진다. 컷오버 이후 프로듀서의 재시도도 같은 이유로 안전하다. stg_trips는 트립당 한 행을 보장하기 위해 FINAL로 쿼리한다. |
fact_trips 엔진 | ReplacingMergeTree(updated_at) | 트립은 정정될 수 있다(요금 조정). updated_at이 버전 컬럼이다 |
agg_hourly_zone_trips 엔진 | ReplacingMergeTree(updated_at) | 롤링 재계산 = upsert 패턴 |
dim_* 테이블 엔진 | MergeTree() | dbt 실행마다 전체 재적재. upsert 없음 |
fact_trips ORDER BY | (toStartOfMonth(pickup_at), pickup_at, trip_id) | Q1-Q7 전부 pickup_at으로 필터한다. trip_id는 리프 수준의 유일성을 보장한다 |
agg_hourly_zone_trips ORDER BY | (hour_bucket, zone_id) | 두 컬럼 모두 모든 집계 쿼리에 등장한다 |
| VARIANT → | String + JSONExtract* | 원본 JSON을 보존한다. 추출은 쿼리 시점에 일어난다 |
| QUALIFY → | ROW_NUMBER()를 감싼 서브쿼리 | ClickHouse는 v24.5부터 네이티브 QUALIFY 절을 갖고 있지만, QUALIFY가 없거나 그 이전인 ClickHouse 버전과 SQL 엔진으로 이식할 수 있어서 서브쿼리 형태를 가르친다 |
| MERGE INTO → | dbt의 delete_insert 증분 | dbt-clickhouse의 관용적 upsert 전략. 전체 테이블 재작성을 피한다 |
1단계 — ClickHouse 클러스터 프로비저닝
setup.sh는 한 가지만 한다. terraform apply를 실행하고 연결 정보를
.clickhouse_state에 쓴다. 또한 무엇이든 프로비저닝하기 전에 migration-plan.md를 다시
확인한다 — 위의 "이유"를 참고하라 — 그러나 그 확인은 경고만 하고 절대 막지 않는다.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
# Configure credentials
cp .env.example .env
vim .env
# Fill in: CLICKHOUSE_ORG_ID, CLICKHOUSE_TOKEN_KEY, CLICKHOUSE_TOKEN_SECRET, CLICKHOUSE_PASSWORD
# Provision
source .env && ./setup.sh.env는 gitignore 대상이다 — 절대 커밋하지 마라.
예상 출력: Terraform이 약 2-3분 안에 2개의 리소스(서비스 + IP 접근 목록)를 만든다.
Apply complete! Resources: 2 added, 0 changed, 0 destroyed.
Outputs:
clickhouse_host = "abc123xyz.us-east-1.aws.clickhouse.cloud"
clickhouse_port = 8443호스트와 포트는 .clickhouse_state에 저장된다. 어느 터미널에서든 이 파일을 source 하면
연결 정보를 가져올 수 있다.
source .clickhouse_state검증:
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 12단계 — 빈 테이블 생성
먼저 올바른 엔진으로 trips_raw를 직접 만든다. 마이그레이션 스크립트가 3단계에서 이
테이블로 데이터를 적재하므로, 행이 하나라도 도착하기 전에 버전 컬럼이 자리 잡도록
ReplacingMergeTree로 이미 존재해야 한다.
-- Run in the ClickHouse SQL console (cloud.clickhouse.com -> SQL console)
CREATE TABLE IF NOT EXISTS default.trips_raw (
trip_id String,
vendor_id UInt8,
pickup_at DateTime64(3, 'UTC'),
dropoff_at DateTime64(3, 'UTC'),
passenger_count UInt8,
trip_distance_miles Float32,
pickup_location_id UInt16,
dropoff_location_id UInt16,
payment_type_id UInt8,
rate_code_id UInt8,
store_fwd_flag String,
fare_amount_usd Float32,
extra_amount_usd Float32,
mta_tax_usd Float32,
tip_amount_usd Float32,
tolls_amount_usd Float32,
total_amount_usd Float32,
ingested_at DateTime64(3, 'UTC'),
trip_metadata String,
_synced_at DateTime DEFAULT now()
)
ENGINE = ReplacingMergeTree(_synced_at)
ORDER BY (pickup_at, trip_id);_synced_at은 모든 INSERT에서 자동으로 설정된다. 마이그레이션 스크립트가 중단되고
--resume로 다시 실행되면 같은 trip_id에 대한 중복 행이 잠시 존재할 수 있다 — RMT가
나중 행(더 큰 _synced_at)을 유지한다. stg_trips는 하위 모델이 데이터를 보기 전에
중복 제거를 강제하기 위해 trips_raw FINAL로 쿼리한다.
다음으로 존 참조 데이터를 시딩한다. 이는 dbt의 stg_taxi_zones가 소스로 읽는 정적
데이터(265개의 NYC TLC 존)다.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
clickhouse-client --host "${CLICKHOUSE_HOST}" --port 9440 --secure \
--user default --password "${CLICKHOUSE_PASSWORD}" \
--multiquery < scripts/00_seed_zones.sql(또는 scripts/00_seed_zones.sql의 내용을 ClickHouse SQL 콘솔에 직접 붙여 넣어도 된다.)
dbt 프로파일 구성. 이 프로젝트의 dbt_project.yml은 profile: 'nyc_taxi_ch'를
선언한다. ~/.dbt/profiles.yml에 대응하는 프로파일이 없으면 dbt run은
Could not find profile named 'nyc_taxi_ch'로 즉시 실패한다. 템플릿은
workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch/profiles.yml.example에
있다.
모듈 01이 이미 Snowflake용 nyc_taxi: 프로파일로 ~/.dbt/profiles.yml을 썼고, 그
모듈의 4단계 갱신 루프가 Snowflake 프로듀서가 도는 동안 계속 그 프로파일로 쿼리한다.
그 파일을 ClickHouse 템플릿으로 교체하지 마라 — profiles.yml.example로 덮어쓰면
nyc_taxi: 프로파일이 삭제되고 모듈 01의 갱신 루프가 깨진다. 대신 템플릿을 열어
nyc_taxi_ch: 블록을 기존 ~/.dbt/profiles.yml에 nyc_taxi:와 나란히 두 번째 최상위
프로파일로 병합하라.
nyc_taxi: # from module 01 — leave this one alone
target: dev
outputs:
dev:
type: snowflake
# ...
nyc_taxi_ch: # add this block
target: dev
outputs:
dev:
type: clickhouse
schema: nyc_taxi_ch
host: "{{ env_var('CLICKHOUSE_HOST') }}"
port: 8443
user: "{{ env_var('CLICKHOUSE_USER', 'default') }}"
password: "{{ env_var('CLICKHOUSE_PASSWORD') }}"
secure: truenyc_taxi_ch:는 env_var()를 통해 환경에서 CLICKHOUSE_HOST, CLICKHOUSE_USER,
CLICKHOUSE_PASSWORD를 읽으므로, 이 모듈의 어떤 dbt 명령보다 먼저 .env와
.clickhouse_state를 source 해야 한다 — 아래의 dbt run은 이미 그렇게 한다. Snowflake
프로파일과 마찬가지로 ~/.dbt/profiles.yml은 자격 증명을 담고 있고 gitignore 대상이다.
절대 커밋하지 마라. 이 병합은 이미 한 세트를 담고 있던 파일에 두 번째 자격 증명 세트를
추가하는 것이다.
검증:
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.clickhouse_state"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.env"
dbt debug
# Expected: "All checks passed!" — confirms dbt found the nyc_taxi_ch profile and
# connected to ClickHouse그다음 dbt run을 실행해 analytics 테이블과 staging 뷰를 만든다.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.clickhouse_state"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.env"
dbt deps # install packages (first run only)
dbt run # creates analytics tables and staging views; all empty at this point예상: 약 8개 모델이 2분 안에 생성된다(모든 테이블은 비어 있다).
검증:
# Check analytics tables were created
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SHOW+TABLES+IN+analytics" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: agg_hourly_zone_trips, dim_date, dim_payment_type, dim_vendor, dim_taxi_zones, fact_trips
# Check staging views were created
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SHOW+TABLES+IN+staging" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: stg_trips, stg_taxi_zones
# Check trips_raw exists with the correct engine
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+engine+FROM+system.tables+WHERE+database%3D%27default%27+AND+name%3D%27trips_raw%27" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: ReplacingMergeTree여섯 개의 analytics 테이블과 두 개의 staging 뷰가 이제 모두 존재하지만, 하나같이 아직
비어 있다 — dbt run은 스키마만 만들었다. 이 단계 이후 데이터가 있는 테이블은
trips_raw 하나뿐이고, 그마저도 아직 데이터가 없다. 그것이 다음 순서다.
3단계 — 데이터 마이그레이션
Python 배치 마이그레이션 스크립트로 Snowflake NYC_TAXI_DB.RAW.TRIPS_RAW의 모든 행을
ClickHouse default.trips_raw로 적재한다.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
source .venv/bin/activate
python scripts/02_migrate_trips.py예상 출력(5천만 행에 약 40-50분):
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
NYC Taxi Migration: Snowflake -> ClickHouse
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Rows to migrate: 50,000,000
Batch size: 100,000
Rows inserted Elapsed ETA Rate
-------------------- ------------ ---------------------- ---------------
100,000 0m 07s 56m 14s remaining 13,945 rows/s
200,000 0m 14s 55m 28s remaining 14,021 rows/s
...이 스크립트가 도는 내내 Snowflake 프로듀서는 TRIPS_RAW에 계속 쓰므로, ClickHouse는 이
전송 시간만큼 뒤처진다 — 그 공백은 예상된 것이고 여기가 아니라 모듈 05에서 다룬다.
스크립트가 중단되면 --resume으로 다시 실행해 마지막 체크포인트부터 이어서 진행하라.
python scripts/02_migrate_trips.py --resume--resume은 ClickHouse에서 max(pickup_at)을 읽어 이미 적재된 행을 건너뛰므로, 이
스크립트는 언제 중단하고 재시작해도 항상 안전하다 — 부분적이고 복구 불가능한 적재로
끝나는 일은 없다.
완료 확인 방법
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
# Row count in trips_raw
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+count()+FROM+default.trips_raw" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: approximately 50000000
# .clickhouse_state was written by setup.sh
ls -la .clickhouse_state
# Service is reachable
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1종료 상태
ClickHouse Cloud 서비스가 살아 있고 접근 가능하며, setup.sh가 CLICKHOUSE_HOST와
CLICKHOUSE_PORT를 담은 .clickhouse_state를 디스크에 썼다. 모든 타깃 테이블과 staging
뷰가 존재한다 — default.trips_raw, 두 개의 staging 뷰(stg_trips, stg_taxi_zones),
여섯 개의 analytics 테이블(fact_trips, agg_hourly_zone_trips, dim_taxi_zones,
dim_payment_type, dim_vendor, dim_date), 그리고 갱신 가능한 materialized view
analytics.mv_live_trip_feed — 즉 이 모듈의 dbt run이 만든 analytics 스키마의 총
일곱 개 객체다. mv_live_trip_feed를 제외한 모든 analytics 테이블은 아직 비어 있다.
mv_live_trip_feed는 dbt run이 뷰를 만들 때 생성한 스냅샷 한 행을 이미 담고 있다.
그 외에 데이터가 있는 것은 default.trips_raw뿐이다. 3단계에서 Python 마이그레이션
스크립트가 옮긴 약 5천만 행이다.
Snowflake 프로듀서는 여전히 실행 중이다. 이 모듈에서 한 번도 중지되지 않았고 여기서도
중지되지 않는다. 마이그레이션 스크립트의 마지막 배치 이후 Snowflake의 TRIPS_RAW에 쓰인
모든 트립은 ClickHouse에 없는 행이므로, ClickHouse는 이제 마이그레이션 구간 길이(약
40-50분에 이 모듈의 설정 시간을 더한 만큼)만큼 Snowflake보다 뒤처져 있다. 그 공백은 실재
하며 프로듀서가 도는 동안 계속 커진다. 이 모듈에서 그것을 닫지 마라. 모듈 05의
컷오버가 공백의 크기를 측정한 다음 제거하는 통제된 2패스 단계로 의도적으로 닫는다 —
지금 프로듀서를 멈추거나 마이그레이션 스크립트를 다시 실행하면 모듈 05가 보여 주도록
만들어진 바로 그것이 사라진다.