dbt di Snowflake
Cara pipeline Medallion di sisi sumber dibangun: source, view staging, model MERGE inkremental, snapshot, dan test.
Dokumen ini menjelaskan bagaimana dbt (data build tool) digunakan dalam NYC Taxi Snowflake Migration Lab — apa yang dilakukannya, mengapa setiap bagian ada, dan bagaimana cara menalarinya.
Yang Dilakukan dbt (dan Yang Tidak)
dbt mentransformasi data yang sudah ada di database Anda. dbt tidak memuat data dari luar, tidak memindahkan file, dan tidak mengelola infrastruktur. Tugasnya adalah mengambil tabel raw dan mengubahnya menjadi tabel yang bersih, teruji, dan siap dianalisis — dengan menjalankan SQL yang Anda tulis.
Anggap dbt sebagai build system untuk SQL. Setiap file .sql di models/ adalah sebuah model yang menjadi tabel atau view di Snowflake. dbt menangani boilerplate CREATE OR REPLACE, menyelesaikan dependensi antar model, dan menjalankan test Anda.
Struktur Proyek
dbt/nyc_taxi_dbt/
├── dbt_project.yml # Project config: name, folder layout, materialization defaults
├── profiles.yml.example # Connection config template (copy to ~/.dbt/profiles.yml)
├── packages.yml # Third-party dbt packages
│
├── macros/
│ ├── generate_schema_name.sql # Overrides dbt's default schema naming logic
│ └── generate_surrogate_key.sql # Wrapper for consistent surrogate key generation
│
└── models/
├── sources.yml # Declares RAW.TRIPS_RAW as an external source
│
├── staging/ # Layer 1: clean and rename raw columns
│ ├── schema.yml # Column-level tests for staging models
│ ├── stg_trips.sql
│ └── stg_taxi_zones.sql
│
├── intermediate/ # Layer 2: joins and enrichment (no physical table)
│ └── int_trips_enriched.sql
│
└── analytics/ # Layer 3: final tables consumed by dashboards
├── schema.yml
├── fact_trips.sql
├── agg_hourly_zone_trips.sql
├── dim_date.sql
├── dim_payment_type.sql
├── dim_taxi_zones.sql
└── dim_vendor.sqlTiga Lapisan (Arsitektur Medallion)
Lapisan 1 — Staging (models/staging/)
Tujuan: Mengambil data raw persis seperti saat tiba dan membuatnya dapat digunakan.
Model-model ini mendarat di schema STAGING sebagai view (tanpa biaya penyimpanan — dieksekusi saat query dijalankan). Setiap model staging mengerjakan satu tugas:
| Model | Sumber | Fungsinya |
|---|---|---|
stg_trips | RAW.TRIPS_RAW | Mengubah nama kolom menjadi snake_case, menambahkan duration_minutes, meratakan kolom VARIANT TRIP_METADATA menjadi kolom bertipe |
stg_taxi_zones | ANALYTICS.DIM_TAXI_ZONES | Pembersihan ringan, menambahkan penjaga COALESCE, menyediakan node lineage dbt |
Pekerjaan terpenting di sini adalah meratakan kolom VARIANT TRIP_METADATA. Sintaks colon-path Snowflake mengekstrak field JSON bersarang:
-- Snowflake: colon-path notation
TRIP_METADATA:driver.rating::FLOAT AS driver_rating,
TRIP_METADATA:app.surge_multiplier::FLOAT AS surge_multiplierIni salah satu tantangan migrasi — ClickHouse menggunakan JSONExtractFloat(TRIP_METADATA, 'driver', 'rating') sebagai gantinya.
Lapisan 2 — Intermediate (models/intermediate/)
Tujuan: Melakukan semua join di satu tempat sehingga tidak perlu diulang-ulang.
int_trips_enriched menggabungkan stg_trips dengan setiap dimensi (zone, payment type, vendor, tanggal) dan menghasilkan satu baris lebar yang sepenuhnya terdenormalisasi per trip. Model ini dideklarasikan ephemeral, artinya dbt menyisipkan (inline) SQL-nya ke dalam model mana pun yang mereferensikannya — tidak ada tabel atau view fisik yang dibuat di Snowflake.
-- dbt_project.yml
intermediate:
+materialized: ephemeral # compiled inline, no CREATE TABLEGunakan ephemeral ketika hasil intermediate hanya dibutuhkan oleh satu model downstream dan Anda tidak ingin menanggung biaya penyimpanan atau overhead kompilasi query.
Lapisan 3 — Analytics (models/analytics/)
Tujuan: Tabel final yang siap dipakai dashboard.
Model-model ini mendarat di schema ANALYTICS. Ada dua jenis:
Tabel dimensi statis — kecil, dimuat ulang sepenuhnya pada setiap dbt run:
| Model | Baris | Catatan |
|---|---|---|
dim_date | ~7,670 | Date spine 2009–2029, dengan kuartal fiskal dan hari libur federal AS |
dim_payment_type | 6 | Passthrough dari data seed |
dim_vendor | 3 | Passthrough dari data seed |
dim_taxi_zones | 265 | Passthrough melalui stg_taxi_zones |
Tabel fact/agregat inkremental — besar, diperbarui dengan MERGE pada setiap run:
| Model | Baris | Catatan |
|---|---|---|
fact_trips | 50M | Satu baris per trip, sepenuhnya terdenormalisasi |
agg_hourly_zone_trips | ~9M | Jumlah per jam yang telah diagregasi sebelumnya untuk setiap zone |
Materialization
Sebuah materialization menentukan objek apa yang dibuat dbt di Snowflake untuk sebuah model.
| Materialization | Objek Snowflake | Kapan digunakan |
|---|---|---|
view | CREATE VIEW | Murah; selalu mencerminkan data terbaru; dipakai untuk staging |
table | CREATE TABLE AS SELECT | Rebuild penuh setiap run; dipakai untuk dimensi kecil |
incremental | MERGE INTO ke tabel yang sudah ada | Tabel besar; hanya memproses baris baru |
ephemeral | (tanpa objek — disisipkan sebagai CTE) | Logika intermediate yang dipakai oleh satu model downstream |
Dua model inkremental berikut menunjukkan strategi inkremental yang berbeda:
fact_trips — memproses trip baru sejak run terakhir:
{% if is_incremental() %}
WHERE pickup_at > (SELECT MAX(pickup_at) FROM {{ this }})
{% endif %}agg_hourly_zone_trips — mengagregasi ulang jendela bergulir 2 jam untuk menangkap data yang datang terlambat:
{% if is_incremental() %}
WHERE pickup_at >= DATEADD('hour', -2, CURRENT_TIMESTAMP())
{% endif %}Pada run pertama (tabel masih kosong), is_incremental() mengembalikan false dan seluruh dataset diproses. Pada run berikutnya, hanya data baru yang diproses. Jika schema berubah dan Anda perlu membangun ulang dari awal, jalankan:
dbt run --full-refreshStrategi MERGE (Tantangan Migrasi Utama)
Ketika incremental_strategy = 'merge', dbt menghasilkan pernyataan MERGE INTO Snowflake:
MERGE INTO ANALYTICS.FACT_TRIPS AS target
USING (SELECT ...) AS source
ON target.trip_id = source.trip_id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT ...;Ini salah satu tantangan migrasi terpenting yang didokumentasikan dalam lab ini. ClickHouse tidak memiliki pernyataan MERGE. Padanannya di ClickHouse adalah menggunakan table engine ReplacingMergeTree dan menambahkan FINAL pada query, atau menggunakan CollapsingMergeTree untuk semantik insert/delete yang eksplisit.
Penamaan Schema: Macro generate_schema_name
Perilaku default dbt menggabungkan target schema dari profiles.yml dengan custom schema di dbt_project.yml:
target schema = STAGING + custom schema = ANALYTICS → STAGING_ANALYTICS (wrong)Proyek ini menimpa perilaku tersebut dengan macro kustom di macros/generate_schema_name.sql:
{% macro generate_schema_name(custom_schema_name, node) -%}
{%- if custom_schema_name is none -%}
{{ target.schema | upper }} -- no custom schema → use target schema
{%- else -%}
{{ custom_schema_name | upper }} -- custom schema → use it directly
{%- endif -%}
{%- endmacro %}Hasilnya: model dengan +schema: ANALYTICS mendarat di ANALYTICS, bukan STAGING_ANALYTICS.
Macro ini diperlukan setiap kali Anda memiliki beberapa schema dalam satu proyek dbt dan tidak ingin nama target schema ditambahkan sebagai awalan.
Koneksi dan Kredensial (profiles.yml)
dbt terhubung ke Snowflake menggunakan sebuah profile yang didefinisikan di ~/.dbt/profiles.yml (jangan pernah di-commit ke git). Nama profile di dbt_project.yml harus sama:
# dbt_project.yml
profile: 'nyc_taxi'
# ~/.dbt/profiles.yml
nyc_taxi:
target: dev
outputs:
dev:
type: snowflake
account: "{{ env_var('SNOWFLAKE_ORG') }}-{{ env_var('SNOWFLAKE_ACCOUNT') }}"
role: DBT_ROLE
database: NYC_TAXI_DB
warehouse: TRANSFORM_WH
schema: STAGING # ← this is the "target schema" / default schema
threads: 4Poin penting:
schema: STAGINGadalah schema default. Model tanpa override+schema:mendarat di sini.role: DBT_ROLEadalah role least-privilege yang dibuat oleh Terraform, hanya dengan izin yang dibutuhkan dbt.threads: 4mengatur berapa banyak model yang dibangun dbt secara paralel.- Kredensial berasal dari environment variable, dimuat dari
.envsebelum menjalankan setup.
Pengujian
Test dbt tersedia dalam dua bentuk:
Schema test (dideklarasikan di schema.yml)
- name: trip_id
tests:
- not_null
- unique
- name: total_amount_usd
tests:
- dbt_expectations.expect_column_values_to_be_between:
min_value: 0
max_value: 1000not_null dan unique sudah tersedia secara bawaan. Test dbt_expectations berasal dari package calogica/dbt_expectations yang dideklarasikan di packages.yml.
Test SQL kustom (tests/)
-- tests/assert_revenue_positive.sql
-- A passing test returns 0 rows
SELECT trip_id, total_amount_usd
FROM {{ ref('fact_trips') }}
WHERE total_amount_usd < 0Test kustom hanyalah query SQL. dbt menjalankannya dan gagal jika ada baris yang dikembalikan.
Jalankan semua test dengan:
dbt testPackage Pihak Ketiga (packages.yml)
packages:
- package: dbt-labs/dbt_utils
version: [">=1.0.0", "<2.0.0"]
- package: calogica/dbt_expectations
version: [">=0.10.0", "<1.0.0"]Instal terlebih dahulu sebelum penggunaan pertama:
dbt depsdbt_utils menyediakan generator date_spine yang dipakai di dim_date.sql. dbt_expectations menyediakan test rentang/distribusi yang melampaui not_null/unique bawaan.
Graf Dependensi
dbt membangun model dalam urutan yang benar secara otomatis dengan mengikuti panggilan {{ ref() }}:
RAW.TRIPS_RAW (source — not managed by dbt)
└── stg_trips (view)
└── int_trips_enriched (ephemeral)
├── fact_trips (incremental table)
└── agg_hourly_zone_trips (incremental table)
ANALYTICS.DIM_TAXI_ZONES (seeded by SQL script)
└── stg_taxi_zones (view)
├── int_trips_enriched
└── dim_taxi_zones (table)
dbt_utils.date_spine
└── dim_date (table){{ ref('stg_trips') }} adalah cara sebuah model mendeklarasikan dependensi pada model lain. {{ source('raw', 'TRIPS_RAW') }} mendeklarasikan dependensi pada tabel eksternal (didefinisikan di sources.yml).
Perintah Umum
| Perintah | Fungsinya |
|---|---|
dbt deps | Instal package dari packages.yml |
dbt run | Bangun semua model (inkremental jika memungkinkan) |
dbt run --full-refresh | Bangun ulang semua model inkremental dari awal |
dbt run -s fact_trips | Bangun hanya fact_trips dan dependensinya |
dbt test | Jalankan semua schema test dan test kustom |
dbt build | dbt run + dbt test sekaligus |
dbt compile | Hasilkan SQL tanpa mengeksekusinya (berguna untuk debugging) |
dbt docs generate && dbt docs serve | Bangun dan jelajahi graf lineage di browser |
Dalam proyek ini, dbt run --full-refresh dipicu otomatis oleh setup.sh jika fact_trips kosong (run pertama atau setelah tear-down).
Posisi dbt dalam Keseluruhan Setup
terraform apply → creates warehouses, database, schemas, roles
scripts/01_create_tables.sql → creates raw tables, seeds dimension data
scripts/02_seed_data.sql → loads 50M synthetic trip rows
dbt deps && dbt build → transforms raw data into analytics-ready tables
scripts/03_create_streams_tasks.sql → creates CDC stream and scheduled taskdbt berada di tengah pipeline. dbt tidak dapat berjalan sampai tabel raw ada dan berisi data. Skrip setup.sh menangani urutan ini.
Contoh lengkap: rencana yang sudah selesai
Rencana migrasi yang sudah terisi untuk workload NYC taxi, untuk dibandingkan dengan rencana Anda sendiri setelah Anda menulisnya.
dbt di ClickHouse
Mengonfigurasi dbt-clickhouse: strategi inkremental delete_insert, model ReplacingMergeTree, dan refreshable materialized view.