Snowflake MigrationClickHouse Workshops

04 dbt 파이프라인 재구축

dbt-clickhouse로 ClickHouse에서 Medallion 파이프라인을 재구축한다 — delete_insert 증분 모델, ReplacingMergeTree, 갱신 가능한 materialized view — 그리고 존 딕셔너리를 만든다.

시작 지점

모듈 03 완료 상태: 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 — 모듈 03의 dbt run이 이미 일곱 개 모두를 만들었다. 그 빌드에서 나온 스냅샷 한 행을 이미 담고 있는 mv_live_trip_feed를 제외하면 모든 analytics 테이블은 아직 비어 있다. 그 외에 데이터가 있는 것은 default.trips_raw뿐이며, 약 5천만 행이다. Snowflake 프로듀서는 여전히 실행 중이므로 ClickHouse는 마이그레이션 구간 길이만큼 Snowflake보다 뒤처져 있다. 약 30분을 예상하라.

이유

모듈 03은 ClickHouse가 5천만 행을 담을 수 있음을 증명했다. 파이프라인이 ClickHouse에서 돌아갈 수 있음은 증명하지 않았다 — staging 뷰, 증분 fact 테이블, 차원 재적재, 파트너가 보기 전에 깨진 모델을 잡아내는 테스트 말이다. 이 모듈이 재구축하는 것이 바로 그것이다. 모듈 01과 동일한 Medallion 모델을 dbt-snowflake 대신 dbt-clickhouse로 표현해, 모듈 03이 이미 만든 테이블에 대해 실행한다.

모델 로직 자체는 아무것도 바뀌지 않는다 — stg_trips는 여전히 타입을 캐스팅하고 JSON을 추출하고, int_trips_enriched는 여전히 차원을 조인하고, fact_trips는 여전히 트립당 한 행으로 끝난다. 바뀌는 것은 그 아래의 머티리얼라이제이션 레이어다. MERGE INTO도 없고, Snowflake Task도 없고, cluster_by도 없다. 이 모듈이 보여 주는 것은 마이그레이션이 일회성 데이터 덤프가 아니라는 사실이다 — 파트너 팀이 매일 같은 스케줄로, 같은 dbt test 실행으로 게이트되어 돌리는 파이프라인이, 그 아래의 웨어하우스가 바뀌어도 계속 동작한다.

개념 — 내부 동작

이 모델 집합은 네 가지 면에서 Snowflake 파이프라인과 다르다. 전체 설정 레퍼런스는 dbt on ClickHouse이며, 여기 있는 것은 1단계에서 dbt run을 실행하기 전에 알아야 할 요약본이다. 이 모델들이 대체하는 소스 파이프라인은 dbt on Snowflake를 참고하라.

1. delete_insert가 MERGE를 대체한다. ClickHouse에는 MERGE INTO 문이 없다. Snowflake 파이프라인이 fact_trips와 agg_hourly_zone_trips를 upsert하기 위해 incremental_strategy: merge를 사용했던 곳에서, ClickHouse 모델은 incremental_strategy: delete_insert를 사용한다. dbt가 들어오는 배치의 unique_key에 해당하는 행을 삭제하고, 그다음 배치를 삽입한다. fact_trips의 unique_key는 trip_id이며, 증분 필터는 pickup_at이 아니라 updated_at에 워터마크를 잡는다 — 요금 정정은 같은 pickup_at을 가진 같은 trip_id를 더 최신의 updated_at으로 재삽입하므로, pickup_at에 워터마크를 잡으면 조용히 놓치게 된다.

2. ReplacingMergeTree는 delete_insert를 대체하는 것이 아니라 그 아래에 깔린 안전망이다. 두 증분 모델 모두 ReplacingMergeTree(updated_at)로 선언된다. delete_insert 실행이 정상적으로 완료되면 테이블에는 이미 키당 한 행만 있고 엔진이 정리할 것이 없다. 실행이 중간에 중단되면 — 삭제 후 삽입 전에 크래시 — 백그라운드 병합이 결국 남은 행들을 중복 제거해 updated_at이 가장 큰 행을 남긴다. delete_insert가 해야 할 중복 제거 작업을 ReplacingMergeTree 하나에만 의존하지 마라. 백그라운드 병합은 비동기이며 이 정도 크기의 테이블에서는 몇 분에서 몇 시간까지 지연될 수 있다.

3. 갱신 가능한 materialized view가 예약 태스크를 대체한다. Snowflake 파이프라인은 롤링 집계를 최신으로 유지하기 위해 저장 프로시저를 실행하는 예약 Task를 사용했다. ClickHouse의 dbt 프로젝트는 대신 mv_live_trip_feed를 materialized = 'materialized_view'와 engine = 'ReplacingMergeTree(refreshed_at)'로 선언하며 — dbt run이 갱신 가능한 materialized view로 만든다. 같은 효과를 내는 Snowflake CREATE TASK DDL 30여 줄과 비교되는 지점이다. 이 모듈은 갱신 주기를 켜지 않는다. 그렇게 하려면 이 랩이 스크립트로 만들지 않은 수동 ALTER TABLE analytics.mv_live_trip_feed MODIFY REFRESH EVERY ... 문이 필요하다(이유는 모듈 05를 참고하라).

4. mv_live_trip_feed에는 Snowflake 대응물이 아예 없다. 기존 모델의 변환이 아니라 마이그레이션이 추가하는 새로운 역량이다. 표준 ClickHouse materialized view는 INSERT마다 한 번 발동하고 그 배치에 속한 행만 보므로, 총 트립 수나 평균 요금 같은 전체 기간 집계를 올바르게 계산할 수 없다. REFRESHABLE materialized view는 대신 자신의 쿼리 전체 — 여기서는 SELECT ... FROM {{ ref('fact_trips') }} — 를 스케줄에 따라 다시 실행하므로 갱신할 때마다 전체 테이블을 본다. 이 워크샵의 Snowflake 쪽에는 애초에 이 선택지가 없었다.

dbt가 실제로 만드는 모델과 각 모델이 데이터를 갖게 되는 시점:

모델레이어머티리얼라이제이션데이터 채워지는 시점비고
stg_tripsstaging뷰1단계(매 실행)타입 캐스팅, trip_metadata에 대한 JSONExtract*
stg_taxi_zonesstaging뷰1단계(매 실행)존 차원 패스스루
int_trips_enrichedstagingEphemeral—(CTE로 인라인된다)모든 차원 조인. 물리 테이블 없음
fact_tripsanalytics증분1단계trip_id를 키로 한 delete_insert, updated_at에 워터마크. ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id)
agg_hourly_zone_tripsanalytics증분모듈 05, 컷오버 이후롤링 2시간 윈도우. 증분 필터가 라이브 프로듀서 행만 매칭한다 — 아래 1단계 참고
dim_taxi_zonesanalytics테이블1단계실행마다 전체 재적재. 2단계의 존 딕셔너리 소스
dim_payment_typeanalytics테이블1단계실행마다 전체 재적재
dim_vendoranalytics테이블1단계실행마다 전체 재적재
dim_dateanalytics테이블1단계2009-2029 정적 날짜 스파인. 실행마다 전체 재적재
mv_live_trip_feedanalyticsMaterialized view(갱신 가능)모듈 03의 dbt run. 갱신 주기는 켜지지 않는다Snowflake 대응물 없음 — 위 4번 참고

1단계 — analytics 레이어 채우기

모듈 00에서 구축한 dbt-clickhouse venv를 활성화한 다음 dbt run을 두 번째로 실행하라. 모듈 03이 스키마를 만들기 위해 빈 테이블에 대해 이미 한 번 실행했고, 이번 실행은 뒤에 실제 데이터가 있다 — trips_raw가 이제 5천만 행을 담고 있다.

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .venv/bin/activate
source .env && source .clickhouse_state

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt run

dbt는 dbt/nyc_taxi_dbt_ch/profiles.yml.example을 바탕으로 한 ~/.dbt/profiles.yml에서 ClickHouse 연결을 읽는다 — 모듈 03이 빈 스키마를 만들 때 쓴 것과 같은 프로파일이다.

예상: 약 8-12분(증분 모델이 5천만 행을 처리한다).

이 실행 후 agg_hourly_zone_trips는 비어 있을 것이다 — 실패가 아니라 예상된 결과다. 이 모델의 증분 필터는 WHERE pickup_at >= now() - INTERVAL 2 HOUR이며, 라이브 프로듀서가 쓴 행만 매칭한다. 방금 마이그레이션한 모든 행은 과거 데이터이므로, 지금 시점에서 측정한 2시간 윈도우 안에 들어오는 것이 없다. 이 테이블은 모듈 05의 컷오버가 ClickHouse 프로듀서를 시작할 때까지 비어 있다 — 이것을 깨진 파이프라인으로 보고 디버깅에 시간을 쓰지 마라.

그다음 테스트 스위트를 실행하라.

dbt test

예상: 모든 테스트 통과.

검증:

SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 million

SELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0 (normal — populated after cutover in module 05)

SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265

2단계 — 존 딕셔너리 생성

analytics.dim_taxi_zones에 이제 265개의 NYC TLC 존이 모두 채워져 있다. 하위 쿼리가 JOIN 대신 dictGet()으로 존의 자치구를 조회할 수 있도록, 그 테이블을 기반으로 한 인메모리 딕셔너리 analytics.taxi_zones_dict를 만들어라.

딕셔너리가 조인보다 이득인 점. 딕셔너리는 메모리에 한 번 적재되어 계속 뜨거운 상태로 남는다. 이후 모든 쿼리에서 그 조회는 사실상 무료다. dim_taxi_zones에 대한 JOIN은 실행할 때마다 차원 테이블을 다시 읽고 다시 매칭한다. 이 테이블처럼 작고 거의 변하지 않는 참조 테이블 — 265행, 매 dbt run마다 전체 재적재 — 이라면 그 거래는 한쪽으로 기울어 있다. 모듈 05의 벤치마크 실행은 dictGet으로 taxi_zones_dict를 직접 쿼리하므로, 이 단계는 그 모듈의 선택적 부가물이 아니라 강한 의존성이다.

연결 정보를 source 한 다음, clickhouse-client 또는 HTTP API로 딕셔너리 DDL을 적재하라 — 가진 쪽을 골라 쓰면 된다.

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .clickhouse_state

# Via clickhouse-client
clickhouse-client \
  --host "${CLICKHOUSE_HOST}" \
  --port 9440 \
  --user default \
  --password "${CLICKHOUSE_PASSWORD}" \
  --secure \
  --multiquery \
  < scripts/04_create_dictionary.sql

# Or via HTTP API
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/" \
  --user "default:${CLICKHOUSE_PASSWORD}" \
  --data-binary @scripts/04_create_dictionary.sql

검증:

-- Should return 'Manhattan' for zone 42
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));

-- Should show status = LOADED, element_count = 265
SELECT name, status, element_count
FROM system.dictionaries
WHERE name = 'taxi_zones_dict';

완료 확인 방법

SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 million
SELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0

여기서 0은 실패가 아니라 정상이다. agg_hourly_zone_trips의 증분 필터 (WHERE pickup_at >= now() - INTERVAL 2 HOUR)는 라이브 프로듀서가 쓴 행만 매칭하며, 지금 ClickHouse에 있는 모든 행은 모듈 03의 마이그레이션 스크립트가 옮긴 과거 데이터다 — now() 기준으로 2시간 미만인 것이 하나도 없다. 이 테이블은 모듈 05의 컷오버가 ClickHouse 프로듀서를 시작한 뒤에야 채워진다. 그때까지 이 테이블을 사용하는 대시보드 차트에는 데이터가 표시되지 않으며, 그것은 디버깅할 대상이 아니라 예상된 동작이다.

SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt test

예상: 모든 테스트 통과.

-- Should return a borough name, e.g. 'Manhattan'
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));

종료 상태

analytics 레이어가 채워지고 테스트를 통과했다. fact_trips가 약 5천만 행을 담고 있고, dim_taxi_zones, dim_payment_type, dim_vendor, dim_date가 완전히 적재되었고, dbt test가 처음부터 끝까지 통과하고, analytics.taxi_zones_dict가 살아 있으며 dictGet()으로 자치구를 반환한다. agg_hourly_zone_trips는 여전히 비어 있다 — 결함이 아니라 설계상 그렇고 — 모듈 05의 컷오버까지 그 상태로 남는다.

대시보드, ClickHouse 대 Snowflake 벤치마크, 컷오버는 이 모듈이 아니라 모듈 05다.

Snowflake 프로듀서는 여전히 실행 중이고, Snowflake와 ClickHouse 사이의 공백도 여전히 열려 있다. 이 모듈의 어떤 것도 프로듀서나 마이그레이션 스크립트를 건드리지 않았다 — 모듈 05가 모듈 03에서 미리 보여 준 것과 동일한 통제된 2패스 단계로 그 공백을 의도적으로 닫는다. 지금 프로듀서를 멈추지 마라.

이 페이지의 내용

Track your progress?

Optional. We email a link to confirm your address; progress records once you open it.

Please use your work email address, not a personal one.

Progress tracking also requires accepting the current Terms of Service in Privacy settings.

KO