04 Bangun ulang pipeline dbt
Bangun ulang pipeline Medallion di ClickHouse dengan dbt-clickhouse — model inkremental delete_insert, ReplacingMergeTree, materialized view yang refreshable — dan buat dictionary zona.
Titik awal
Modul 03 selesai: layanan ClickHouse Cloud hidup dan dapat dijangkau, dan
.clickhouse_state sudah ditulis ke disk oleh setup.sh dengan CLICKHOUSE_HOST dan
CLICKHOUSE_PORT. Semua tabel target dan view staging ada — default.trips_raw, kedua view
staging (stg_trips, stg_taxi_zones), keenam tabel analytics (fact_trips,
agg_hourly_zone_trips, dim_taxi_zones, dim_payment_type, dim_vendor, dim_date), dan
materialized view yang refreshable analytics.mv_live_trip_feed — dbt run modul 03 sudah
membangun ketujuhnya. Setiap tabel analytics masih kosong kecuali mv_live_trip_feed, yang
sudah menyimpan satu baris snapshot dari build tersebut. Selain itu hanya
default.trips_raw yang berisi data: kurang lebih 50 juta baris. Producer Snowflake masih
berjalan, jadi ClickHouse tertinggal dari Snowflake kurang lebih sepanjang jendela migrasi.
Siapkan sekitar 30 menit.
Mengapa
Modul 03 membuktikan ClickHouse bisa menampung 50 juta baris. Ia belum membuktikan pipeline-nya bisa berjalan di ClickHouse — view staging, tabel fact inkremental, reload dimensi, tes yang menangkap model rusak sebelum seorang partner melihatnya. Itulah yang dibangun ulang modul ini: model Medallion yang sama dari modul 01, diekspresikan dengan dbt-clickhouse alih-alih dbt-snowflake, berjalan terhadap tabel yang sudah dibuat modul 03.
Tidak ada logika model yang berubah — stg_trips tetap melakukan type-cast dan mengekstrak
JSON, int_trips_enriched tetap men-join dimensi, fact_trips tetap berakhir sebagai satu
baris per perjalanan. Yang berubah adalah lapisan materialisasi di bawahnya: tanpa
MERGE INTO, tanpa Snowflake Task, tanpa cluster_by. Inilah modul yang menunjukkan bahwa
migrasi bukan sekadar dump data satu kali — pipeline yang dijalankan tim partner setiap hari,
pada jadwal yang sama, digerbangi eksekusi dbt test yang sama, tetap bekerja setelah
warehouse di bawahnya berganti.
Konsep — di balik layar
Kumpulan model ini berbeda dari pipeline Snowflake dalam empat hal. Referensi konfigurasi
lengkapnya ada di
dbt di ClickHouse — ini adalah versi
singkat yang Anda perlukan sebelum menjalankan dbt run di Langkah 1. Untuk pipeline sumber
yang digantikan model-model ini, lihat
dbt di Snowflake.
1. delete_insert menggantikan MERGE. ClickHouse tidak punya pernyataan MERGE INTO.
Di tempat pipeline Snowflake memakai incremental_strategy: merge untuk meng-upsert
fact_trips dan agg_hourly_zone_trips, model ClickHouse memakai
incremental_strategy: delete_insert: dbt menghapus baris yang cocok dengan unique_key
untuk batch masuk, lalu menyisipkan batch itu. Untuk fact_trips, unique_key adalah
trip_id, dan filter inkrementalnya memakai watermark pada updated_at alih-alih
pickup_at — koreksi tarif menyisipkan ulang trip_id yang sama dengan pickup_at yang
sama tetapi updated_at yang lebih baru, jadi watermark pada pickup_at akan diam-diam
melewatkannya.
2. ReplacingMergeTree adalah jaring pengaman di bawah delete_insert, bukan
penggantinya. Kedua model inkremental dideklarasikan sebagai
ReplacingMergeTree(updated_at). Jika eksekusi delete_insert selesai normal, tabelnya sudah
punya satu baris per kunci dan engine tidak punya apa pun untuk dibersihkan. Jika eksekusi
terinterupsi di tengah jalan — crash setelah delete, sebelum insert — merge latar belakang
akhirnya mendeduplikasi baris sisa mana pun, menyimpan baris dengan updated_at tertinggi.
Jangan pernah mengandalkan ReplacingMergeTree sendirian untuk melakukan pekerjaan
deduplikasi yang mestinya dilakukan delete_insert: merge latar belakang bersifat asinkron
dan bisa tertinggal beberapa menit sampai beberapa jam pada tabel sebesar ini.
3. Materialized view yang refreshable menggantikan task terjadwal. Pipeline Snowflake
memakai Task terjadwal yang menjalankan stored procedure untuk menjaga agregat bergulir tetap
mutakhir. Proyek dbt ClickHouse justru mendeklarasikan mv_live_trip_feed dengan
materialized = 'materialized_view' dan engine = 'ReplacingMergeTree(refreshed_at)' —
dibangun oleh dbt run sebagai materialized view yang refreshable, versus lebih dari 30 baris
DDL CREATE TASK Snowflake untuk efek yang sama. Modul ini tidak menyalakan interval
refresh; melakukannya adalah pernyataan manual
ALTER TABLE analytics.mv_live_trip_feed MODIFY REFRESH EVERY ... yang tidak di-skrip oleh
lab (lihat modul 05 untuk alasannya).
4. mv_live_trip_feed sama sekali tidak punya padanan di Snowflake. Ia bukan terjemahan
dari model yang sudah ada — ia adalah kapabilitas baru yang ditambahkan migrasi. Materialized
view ClickHouse standar terpicu satu kali per INSERT dan hanya melihat baris dalam batch
tersebut, jadi ia tidak bisa menghitung agregat sepanjang masa seperti total perjalanan atau
rata-rata tarif dengan benar. Materialized view REFRESHABLE justru menjalankan ulang seluruh
query-nya — di sini, SELECT ... FROM {{ ref('fact_trips') }} — pada sebuah jadwal, jadi
setiap refresh melihat seluruh tabel. Sisi Snowflake dari workshop ini tidak pernah punya
opsi ini.
Model yang benar-benar dibangun dbt, dan kapan masing-masing mendapat data:
| Model | Lapisan | Materialisasi | Terisi | Catatan |
|---|---|---|---|---|
stg_trips | staging | View | Langkah 1 (setiap eksekusi) | Type-cast, JSONExtract* untuk trip_metadata |
stg_taxi_zones | staging | View | Langkah 1 (setiap eksekusi) | Passthrough dimensi zona |
int_trips_enriched | staging | Ephemeral | — (di-inline sebagai CTE) | Semua join dimensi; tanpa tabel fisik |
fact_trips | analytics | Inkremental | Langkah 1 | delete_insert berkunci trip_id, watermark pada updated_at; ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id) |
agg_hourly_zone_trips | analytics | Inkremental | Modul 05, pasca-cutover | Jendela bergulir 2 jam; filter inkrementalnya hanya cocok dengan baris dari producer live — lihat Langkah 1 di bawah |
dim_taxi_zones | analytics | Tabel | Langkah 1 | Reload penuh per eksekusi; sumber untuk dictionary zona di Langkah 2 |
dim_payment_type | analytics | Tabel | Langkah 1 | Reload penuh per eksekusi |
dim_vendor | analytics | Tabel | Langkah 1 | Reload penuh per eksekusi |
dim_date | analytics | Tabel | Langkah 1 | Tulang belakang tanggal statis, 2009-2029; reload penuh per eksekusi |
mv_live_trip_feed | analytics | Materialized view (refreshable) | dbt run modul 03; interval refresh tidak pernah dinyalakan | Tanpa padanan Snowflake — lihat poin 4 di atas |
Langkah 1 — Isi lapisan analytics
Aktifkan venv dbt-clickhouse yang Anda bangun di modul 00, lalu jalankan dbt run untuk
kedua kalinya. Modul 03 sudah menjalankannya sekali terhadap tabel kosong untuk membuat
skemanya; eksekusi kali ini punya data nyata di belakangnya — trips_raw kini menyimpan 50
juta baris.
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 rundbt membaca koneksi ClickHouse dari ~/.dbt/profiles.yml, berdasarkan
dbt/nyc_taxi_dbt_ch/profiles.yml.example — profil yang sama yang dipakai modul 03 untuk
membuat skema kosong.
Diharapkan: sekitar 8-12 menit (50 juta baris diproses oleh model inkremental).
agg_hourly_zone_trips akan kosong setelah eksekusi ini — itu memang diharapkan, bukan
kegagalan. Filter inkrementalnya adalah WHERE pickup_at >= now() - INTERVAL 2 HOUR, yang
hanya cocok dengan baris yang ditulis producer live. Setiap baris yang baru Anda migrasikan
bersifat historis, jadi tidak ada yang masuk ke dalam jendela 2 jam yang diukur dari sekarang.
Tabel ini tetap kosong sampai cutover di modul 05 menjalankan producer ClickHouse — jangan
membuang waktu men-debug ini sebagai pipeline yang rusak.
Lalu jalankan rangkaian tesnya:
dbt testDiharapkan: semua tes lulus.
Verifikasi:
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: 265Langkah 2 — Buat dictionary zona
analytics.dim_taxi_zones sekarang terisi dengan seluruh 265 zona NYC TLC. Bangun
analytics.taxi_zones_dict, sebuah dictionary in-memory yang ditopang tabel tersebut,
sehingga query hilir bisa mencari borough sebuah zona dengan dictGet() alih-alih JOIN.
Apa keuntungan dictionary dibanding join. Sebuah dictionary dimuat ke memori satu kali
dan tetap hangat; pencarian terhadapnya praktis gratis pada setiap query berikutnya. Sebuah
JOIN terhadap dim_taxi_zones membaca ulang dan mencocokkan ulang tabel dimensi setiap
kali dijalankan. Untuk tabel referensi kecil yang jarang berubah seperti ini — 265 baris,
di-reload penuh oleh setiap dbt run — pertukaran itu jelas sepihak. Eksekusi benchmark di
modul 05 melakukan query taxi_zones_dict langsung dengan dictGet, jadi langkah ini adalah
dependensi keras untuk modul itu, bukan tambahan opsional.
Source detail koneksinya, lalu muat DDL dictionary baik melalui clickhouse-client atau API
HTTP — pilih mana pun yang tersedia bagi Anda:
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.sqlVerifikasi:
-- 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';Cara memverifikasi bahwa Anda sudah selesai
SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 millionSELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0Nol adalah nilai yang benar di sini, bukan kegagalan. Filter inkremental
agg_hourly_zone_trips (WHERE pickup_at >= now() - INTERVAL 2 HOUR) hanya cocok dengan
baris yang ditulis producer live, dan setiap baris di ClickHouse saat ini adalah data
historis yang dipindahkan skrip migrasi modul 03 — tak satu pun berumur kurang dari 2 jam
relatif terhadap now(). Tabel ini baru terisi setelah cutover di modul 05 menjalankan
producer ClickHouse; sampai saat itu, setiap chart dashboard yang ditopang tabel ini akan
menampilkan data kosong, dan itu memang diharapkan, bukan sesuatu untuk di-debug.
SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt testDiharapkan: semua tes lulus.
-- Should return a borough name, e.g. 'Manhattan'
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));Kondisi akhir
Lapisan analytics sudah terisi dan teruji: fact_trips menyimpan kurang lebih 50 juta baris,
dim_taxi_zones, dim_payment_type, dim_vendor, dan dim_date terisi penuh, dbt test
lulus dari ujung ke ujung, dan analytics.taxi_zones_dict hidup serta mengembalikan borough
melalui dictGet(). agg_hourly_zone_trips masih kosong — sesuai desain, bukan karena
cacat — dan tetap demikian sampai cutover di modul 05.
Dashboard, benchmark ClickHouse-vs-Snowflake, dan cutover adalah modul 05, bukan modul ini.
Producer Snowflake masih berjalan, dan jeda antara Snowflake dan ClickHouse masih terbuka. Tidak ada di modul ini yang menyentuh producer atau skrip migrasi — modul 05 menutup jeda itu dengan sengaja, dalam langkah dua tahap terkendali yang sama yang dipratinjau modul 03. Jangan hentikan producer sekarang.
03 Provisioning dan migrasi
Provisioning ClickHouse Cloud dengan Terraform, buat tabel target dari rencana Anda, dan pindahkan 50 juta baris dengan skrip migrasi Python yang dapat dilanjutkan.
05 Benchmark dan cutover
Bangun ulang dashboard di ClickHouse, benchmark ketujuh query terhadap kedua engine, cutover producer, verifikasi paritas, dan bongkar semuanya.