Contoh lengkap: rencana yang sudah selesai
Rencana migrasi yang sudah terisi untuk workload NYC taxi, untuk dibandingkan dengan rencana Anda sendiri setelah Anda menulisnya.
Ini adalah jawaban yang dikerjakan penuh untuk kelima worksheet
(1,
2,
3,
4,
5) yang diterapkan pada
workload NYC Taxi. Gunakan untuk:
- Memeriksa jawaban worksheet Anda setelah menyelesaikan setiap bagian
- Memahami alasan di balik keputusan yang diimplementasikan Bagian 3
- Membandingkan dengan tabel Decision Alignment di Bagian 3 jika Anda memilih berbeda
Ini adalah kunci jawaban — jangan mengisinya sebagai rencana Anda. Isilah migration-plan.md sebagai gantinya.
| Metrik | Nilai |
|---|
| Total tabel | 7 (TRIPS_RAW, FACT_TRIPS, AGG_HOURLY_ZONE_TRIPS, DIM_TAXI_ZONES, DIM_PAYMENT_TYPE, DIM_VENDOR, DIM_DATE) |
| Total view | 2 (STG_TRIPS, STG_TAXI_ZONES) |
| Stream | 1 (TRIPS_CDC_STREAM pada TRIPS_RAW) |
| Task | 2 (CDC_CONSUME_TASK, HOURLY_AGG_TASK) |
| Total baris di TRIPS_RAW | ~50,000,000 |
| Rentang tanggal | Jendela bergulir 4 tahun yang berakhir pada waktu setup |
| Kolom VARIANT | 1 (TRIPS_RAW.TRIP_METADATA) |
| Penggunaan QUALIFY yang terdeteksi | 1 (query Q3) |
| Penggunaan MERGE INTO yang terdeteksi | 2 (HOURLY_AGG_TASK, CDC_CONSUME_TASK) |
| Objek | Tipe | Schema | Baris | Tingkat Kompleksitas | Catatan |
|---|
trips_raw | Tabel | raw | ~50M | B | RMT dengan kolom versi _synced_at; tumpang tindih bulk + CDC memerlukan dedup; stg_trips harus memakai FINAL |
stg_trips | dbt View | staging | — | B | JSONExtract untuk TRIP_METADATA; memerlukan pengujian path JSON |
stg_taxi_zones | dbt View | staging | — | A | Passthrough; trivial |
int_trips_enriched | dbt Ephemeral | staging | — | A | CTE; perbedaan SQL ditangani di model induk |
fact_trips | dbt Incremental | analytics | ~50M | C | Engine RMT; delete_insert; penulisan ulang QUALIFY; FINAL wajib |
agg_hourly_zone_trips | dbt Incremental | analytics | ~140K | B | RMT; jendela rekalkulasi bergulir 2 jam; uji batas partisi dengan cermat |
dim_taxi_zones | dbt Table | analytics | 265 | A | Referensi statis; muat ulang penuh; trivial |
dim_payment_type | dbt Table | analytics | 6 | A | Referensi statis; trivial |
dim_vendor | dbt Table | analytics | 3 | A | Referensi statis; trivial |
taxi_zones_dict | Dictionary | analytics | 265 | B | Sintaks khusus ClickHouse; dictGet() pada saat query |
mv_hourly_revenue | Refreshable MV | analytics | — | B | Sintaks REFRESH EVERY; verifikasi penggantian atomik |
TRIPS_CDC_STREAM / CDC_CONSUME_TASK | Snowflake Stream + Task | — | — | D | Tidak ada padanan di ClickHouse; digantikan oleh cutover producer langsung di Bagian 3 |
| Tabel | Engine | Kolom Versi | Alasan |
|---|
trips_raw | ReplacingMergeTree(_synced_at) | _synced_at | Skrip migrasi Python (scripts/02_migrate_trips.py) mungkin mencoba ulang sebuah batch dan meng-insert kembali trip_id yang sama. Setelah cutover, producer live juga bisa mencoba ulang saat terjadi kegagalan sementara. _synced_at DateTime DEFAULT now() disetel otomatis pada INSERT — percobaan ulang yang lebih akhir punya timestamp lebih tinggi, sehingga RMT mempertahankan tulisan terbaru. stg_trips melakukan query dengan FINAL untuk menegakkan dedup sebelum model downstream mana pun berjalan. |
fact_trips | ReplacingMergeTree(updated_at) | updated_at | Trip dapat dikoreksi (penyesuaian tarif, perubahan status). trip_id yang sama di-insert ulang dengan nilai yang diperbarui. updated_at naik secara monoton pada setiap koreksi — nilai yang lebih tinggi menang selama dedup RMT. Selalu query dengan FINAL. |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | dbt menghitung ulang 2 jam terakhir dan meng-insert ulang. Tanpa RMT, agregat lama dan baru menumpuk dan terhitung dua kali. updated_at yang disetel ke now() pada setiap dbt run memastikan nilai terbaru yang menang. |
dim_taxi_zones | MergeTree() | — | Muat ulang penuh oleh dbt (pertukaran tabel atomik (rebuild penuh)). Tidak ada duplikat yang bisa menumpuk. Tidak perlu dedup. |
dim_payment_type | MergeTree() | — | Sama — muat ulang penuh. |
dim_vendor | MergeTree() | — | Sama — muat ulang penuh. |
mv_hourly_revenue | MergeTree() | — | MV REFRESHABLE menggantikan seluruh result set-nya secara atomik pada setiap REFRESH. Tidak ada upsert. |
| Tabel | ORDER BY | Alasan |
|---|
trips_raw | (pickup_at, trip_id) | Scan rentang waktu memfilter pickup_at terlebih dahulu. trip_id adalah key dedup RMT — key ini harus ada di ORDER BY agar RMT dapat mengidentifikasi baris mana yang duplikat. pickup_at didahulukan karena scan rentang analitis yang dominan; trip_id ditempatkan terakhir karena kardinalitasnya tinggi dan hanya berperan sebagai pembeda keunikan. |
fact_trips | (toStartOfMonth(pickup_at), pickup_at, trip_id) | Ketujuh query analitis memfilter pickup_at. Prefiks bulan mengelompokkan data bulan kalender ke dalam blok yang berdekatan — memungkinkan block skipping kasar untuk agregasi bulanan tanpa menambahkan PARTITION BY. trip_id ditempatkan terakhir untuk keunikan RMT tanpa mengganggu block skipping. |
agg_hourly_zone_trips | (hour_bucket, zone_id) | Q6 (dan semua query agregasi) memfilter hour_bucket dan zone_id. hour_bucket memiliki ~35K nilai berbeda; zone_id memiliki 265. hour_bucket didahulukan karena scan rentang waktu adalah pola akses utama. zone_id di posisi kedua untuk pemfilteran sekunder. |
dim_taxi_zones | (location_id) | 265 baris = satu granule. ORDER BY tidak relevan untuk performa. location_id sebagai join key adalah konvensi dan membantu keterbacaan. |
| Kolom | Tipe Snowflake | Tipe ClickHouse | Alasan Keputusan |
|---|
TRIP_METADATA | VARIANT | String | Mempertahankan JSON mentah secara persis. JSONExtract* menangani path arbitrer pada saat query. Map(String,String) kehilangan struktur bersarang; Tuple memerlukan schema tetap. String adalah pilihan aman untuk JSON arbitrer. |
PICKUP_DATETIME / PICKUP_AT | TIMESTAMP_NTZ(9) | DateTime64(3, 'UTC') | Presisi milidetik cukup untuk timestamp trip. Nanodetik (9) berlebihan. 'UTC' membuat timezone eksplisit dan menghindari kejutan terkait DST dalam agregasi rentang waktu. |
PICKUP_LOCATION_ID | INTEGER | UInt16 | Nilai 1–265. Maksimum UInt8 adalah 255 (terlalu kecil). Maksimum UInt16 adalah 65535 (benar). 2 byte dibanding 4 byte untuk Int32 — menghemat ~95MB tanpa kompresi pada 50M baris per kolom. |
VENDOR_ID | INTEGER | UInt8 | Nilai 1–3. Maksimum UInt8 adalah 255 — benar. 1 byte per baris. |
DRIVER_RATING | FLOAT | Nullable(Float32) | Sering NULL (tidak semua trip punya rating). Nullable mempertahankan semantik null yang benar. Float32 cukup untuk rentang 1.0–5.0. Float64 akan memboroskan penyimpanan tanpa menambah presisi yang berarti. |
UPDATED_AT | TIMESTAMP_NTZ(9) | DateTime64(3, 'UTC') | Kolom versi untuk ReplacingMergeTree. Harus memakai DateTime64, bukan DateTime — dua koreksi dalam detik yang sama akan non-deterministik dengan presisi detik. Presisi milidetik memastikan urutan dedup yang benar. |
| Ekspresi Snowflake | Padanan 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 inkremental — DELETE baris dengan key yang cocok, lalu INSERT semua baris baru |
| Wave | Objek | Dependensi | Catatan |
|---|
| Wave 0 | trips_raw (schema), dim_taxi_zones, dim_payment_type, dim_vendor | Tidak ada | dbt membuat tabel kosong. Tabel dim langsung diisi dari data referensi statis (tanpa dependensi pada trips). Jalankan: dbt run --select trips_raw dim_* |
| Wave 1 | Bulk load Python (scripts/02_migrate_trips.py) | Wave 0 (schema trips_raw harus sudah ada) | 50M baris dari TRIPS_RAW Snowflake. Dapat dilanjutkan dengan --resume. Verifikasi jumlah baris dengan scripts/01_verify_migration.sh. ~40-50 menit. |
| Wave 2 | stg_trips, stg_taxi_zones, int_trips_enriched, fact_trips, agg_hourly_zone_trips | Wave 1 selesai (trips_raw terisi) + Wave 0 (tabel dim ada) | dbt run penuh. stg_trips membaca trips_raw; int_trips_enriched melakukan join dengan dim; fact_trips dan agg_hourly_zone_trips dibangun di atasnya. |
| Wave 3 | taxi_zones_dict, mv_live_trip_feed | Wave 2 (dim_taxi_zones terisi untuk dict; fact_trips terisi untuk MV) | Dictionary dibuat via scripts/04_create_dictionary.sql. Refreshable MV dibuat via model dbt; mengaktifkan interval refresh-nya sesudahnya adalah langkah manual ALTER TABLE ... MODIFY REFRESH, bukan sesuatu yang dijalankan dbt secara otomatis. |
| Wave 4 | Cutover producer (scripts/03_cutover.sh) | Wave 1 selesai (bulk load terverifikasi) + Wave 2 selesai (lapisan analytics terbangun) | Hentikan producer Snowflake; jalankan producer ClickHouse yang menulis langsung ke ClickHouse Cloud; jalankan dbt run untuk mengisi agg_hourly_zone_trips dengan data live. |
| Objek | Risiko | Metode Verifikasi |
|---|
fact_trips | Query tanpa FINAL menghitung berlebih selama merge lag. Rentang partisi delete_insert harus dibatasi pada prefiks ORDER BY agar tidak menghapus partisi non-target. | SELECT COUNT(*) FINAL cocok dengan Snowflake ± CDC lag. Jalankan dbt test. Bandingkan hasil Q3 antar sistem. Periksa trip_id duplikat: SELECT trip_id, count() FROM fact_trips GROUP BY trip_id HAVING count() > 1 LIMIT 10. |
agg_hourly_zone_trips | Jendela rekalkulasi bergulir 2 jam harus membatasi rentang delete dengan benar. Jika terlalu luas, agregat lama terhapus; jika terlalu sempit, agregat basi tetap bertahan. | Periksa acak tuple (hour_bucket, zone_id) tertentu terhadap Snowflake. Verifikasi total trip_count di seluruh zone cocok dengan AGG_HOURLY_ZONE_TRIPS Snowflake untuk periode yang sama. |
| Cutover producer | Skrip migrasi yang terputus di tengah run meninggalkan celah jumlah baris; jalankan ulang dengan --resume untuk mengisinya. Percobaan ulang producer setelah cutover dapat meng-insert kembali trip yang sudah ada di ClickHouse. | scripts/01_verify_migration.sh — memeriksa paritas jumlah baris antara Snowflake dan ClickHouse. ReplacingMergeTree(_synced_at) menangani insert duplikat secara idempoten. |
Pemindahan data: skrip migrasi Python (scripts/02_migrate_trips.py)
Mengapa skrip Python dan bukan relay object storage atau ClickPipes?
remoteSecure() ditujukan untuk transfer data ClickHouse-ke-ClickHouse — tidak berlaku di sini.
- Relay object storage (Snowflake → S3 → fungsi tabel S3 ClickHouse) akan berhasil tetapi menambah kompleksitas: memerlukan provisioning bucket S3, role IAM, dan COPY INTO Snowflake — overhead yang tidak perlu untuk sebuah lab.
- ClickPipes tidak mendukung Snowflake sebagai sumber. Sumber yang didukungnya adalah Kafka, S3, Kinesis, PostgreSQL CDC, dan MySQL CDC.
- Skrip Python memakai
snowflake-connector-python dan clickhouse-connect — package yang sudah terinstal untuk lab ini. Skrip menampilkan progres secara real time, mendukung --resume untuk gangguan, dan kodenya sepenuhnya dapat diperiksa.
Strategi inkremental (dbt): delete_insert
Mengapa delete_insert dan bukan strategi append atau merge?
append meng-insert baris baru tanpa menyentuh baris yang sudah ada. Untuk fact_trips yang barisnya dapat diperbarui, ini menciptakan duplikat. Salah.
merge (jika tersedia) akan paling dekat dengan MERGE INTO milik Snowflake, tetapi strategi merge dbt-clickhouse punya keterbatasan dengan ReplacingMergeTree dan bukan pendekatan yang direkomendasikan.
delete_insert menghapus baris dalam rentang key batch yang masuk, lalu meng-insert semua baris baru. Ini idempoten (menjalankan ulang menghasilkan hasil yang sama), menangani insert maupun update, dan bekerja dengan benar bersama ReplacingMergeTree. Ini rekomendasi standar komunitas dbt-clickhouse untuk pola upsert.
| Kriteria | Ambang | Diukur Dengan |
|---|
| Paritas jumlah baris | cocok ≥ 99.9% (CH ≥ SF pasca-cutover memang diharapkan) | scripts/01_verify_migration.sh |
| Paritas checksum | MD5 cocok pada sampel 10K baris | scripts/02_validate_parity.sql |
| Tingkat kelulusan dbt test | 100% | dbt test di dbt/nyc_taxi_dbt_ch |
| Paritas hasil query | Ketujuh query mengembalikan hasil yang sama (dalam toleransi floating-point) | Perbandingan manual pada output scripts/run_benchmark.sh |
| Model | Materialization | Alasan |
|---|
stg_trips | view | Membaca dan membersihkan trips_raw; tidak ada pembaruan pada model ini; biaya penyimpanan nol; selalu mencerminkan keadaan sumber saat ini |
stg_taxi_zones | view | Sama — pembersihan passthrough atas tabel sumber |
int_trips_enriched | ephemeral | Logika join murni yang hanya dipakai oleh fact_trips; disisipkan sebagai CTE sehingga menghindari tabel fisik yang redundan; tidak ada model yang mengqueri-nya secara langsung |
fact_trips | incremental | Trip dapat dikoreksi setelah kejadian; hanya baris baru dan baris yang diperbarui yang perlu diproses per run |
agg_hourly_zone_trips | incremental | Rekalkulasi bergulir 2 jam adalah pola inkremental — proses baris terbaru, bukan seluruh 50M |
dim_taxi_zones | table | 265 zone statis; rebuild penuh pada setiap dbt run via pertukaran tabel atomik (rebuild penuh); tanpa pembaruan sebagian |
dim_payment_type | table | 6 tipe statis; alasan sama seperti dim_taxi_zones |
dim_vendor | table | 3 vendor; alasan sama |
| Model | ENGINE | Kolom Versi | Alasan |
|---|
fact_trips | ReplacingMergeTree(updated_at) | updated_at | Trip dapat dikoreksi; updated_at yang disetel ke now() pada setiap insert berarti versi terbaru menang selama dedup background RMT; delete_insert adalah jalur kebenaran utama, RMT adalah jaring pengaman |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | Rekalkulasi bergulir meng-insert ulang agregat untuk pasangan (hour_bucket, zone_id) yang sama; RMT memastikan agregat basi dihapus pada background merge |
dim_taxi_zones | MergeTree() | — | Muat ulang penuh oleh dbt berarti pertukaran tabel atomik (rebuild penuh) pada setiap run; duplikat tidak dapat menumpuk; tidak perlu dedup |
dim_payment_type | MergeTree() | — | Sama seperti dim_taxi_zones |
dim_vendor | MergeTree() | — | Sama seperti dim_taxi_zones |
| Model | unique_key | incremental_strategy | Filter inkremental | Mengapa filter ini? |
|---|
fact_trips | trip_id | delete_insert | WHERE updated_at > (SELECT max(updated_at) FROM {{ this }}) | High-watermark pada updated_at menangkap baik trip baru maupun trip yang dikoreksi (penyesuaian tarif meng-insert ulang trip_id yang sama dengan pickup_at yang sama tetapi updated_at yang lebih baru); watermark pickup_at akan melewatkan koreksi secara diam-diam |
agg_hourly_zone_trips | [hour_bucket, zone_id] | delete_insert | WHERE pickup_at >= now() - INTERVAL 2 HOUR | Jendela bergulir 2 jam memaksa agregasi ulang jam-jam batas sehingga hitungan jam yang belum lengkap selalu dikoreksi; high-watermark max(pickup_at) akan selamanya menghitung kurang pada jam batas |
| Model | FINAL di klausa FROM? | Alasan |
|---|
stg_trips | Ya — FROM trips_raw FINAL | trips_raw adalah ReplacingMergeTree; tabel ini bisa memiliki baris trip_id duplikat dari percobaan ulang skrip migrasi atau percobaan ulang producer pasca-cutover. stg_trips adalah satu-satunya titik penegakan: deduplikasi di sini agar setiap model downstream (int_trips_enriched, fact_trips, agg_hourly_zone_trips) menerima data yang bersih |
int_trips_enriched | Tidak | Membaca dari stg_trips (sebuah view), bukan tabel RMT; FINAL tidak relevan untuk view |
fact_trips | Tidak (di dalam body model) | delete_insert menjaga fact_trips tetap bersih setelah setiap run yang selesai; menambahkan FINAL di dalam model akan menerapkannya secara sia-sia pada subquery is_incremental() yang membaca max(updated_at) dari {{ this }}. Dashboard dan test dbt memakai FINAL secara eksternal ketika mengqueri fact_trips secara langsung |
Ini adalah contoh yang sudah lengkap. migration-plan.md Anda seharusnya cocok dengan keputusan-keputusan utama di sini — atau mendokumentasikan secara eksplisit mengapa Anda memilih berbeda.