Snowflake MigrationClickHouse Workshops

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.sql

Tiga 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:

ModelSumberFungsinya
stg_tripsRAW.TRIPS_RAWMengubah nama kolom menjadi snake_case, menambahkan duration_minutes, meratakan kolom VARIANT TRIP_METADATA menjadi kolom bertipe
stg_taxi_zonesANALYTICS.DIM_TAXI_ZONESPembersihan 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_multiplier

Ini 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 TABLE

Gunakan 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:

ModelBarisCatatan
dim_date~7,670Date spine 2009–2029, dengan kuartal fiskal dan hari libur federal AS
dim_payment_type6Passthrough dari data seed
dim_vendor3Passthrough dari data seed
dim_taxi_zones265Passthrough melalui stg_taxi_zones

Tabel fact/agregat inkremental — besar, diperbarui dengan MERGE pada setiap run:

ModelBarisCatatan
fact_trips50MSatu baris per trip, sepenuhnya terdenormalisasi
agg_hourly_zone_trips~9MJumlah per jam yang telah diagregasi sebelumnya untuk setiap zone

Materialization

Sebuah materialization menentukan objek apa yang dibuat dbt di Snowflake untuk sebuah model.

MaterializationObjek SnowflakeKapan digunakan
viewCREATE VIEWMurah; selalu mencerminkan data terbaru; dipakai untuk staging
tableCREATE TABLE AS SELECTRebuild penuh setiap run; dipakai untuk dimensi kecil
incrementalMERGE INTO ke tabel yang sudah adaTabel 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-refresh

Strategi 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: 4

Poin penting:

  • schema: STAGING adalah schema default. Model tanpa override +schema: mendarat di sini.
  • role: DBT_ROLE adalah role least-privilege yang dibuat oleh Terraform, hanya dengan izin yang dibutuhkan dbt.
  • threads: 4 mengatur berapa banyak model yang dibangun dbt secara paralel.
  • Kredensial berasal dari environment variable, dimuat dari .env sebelum 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: 1000

not_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 < 0

Test kustom hanyalah query SQL. dbt menjalankannya dan gagal jika ada baris yang dikembalikan.

Jalankan semua test dengan:

dbt test

Package 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 deps

dbt_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

PerintahFungsinya
dbt depsInstal package dari packages.yml
dbt runBangun semua model (inkremental jika memungkinkan)
dbt run --full-refreshBangun ulang semua model inkremental dari awal
dbt run -s fact_tripsBangun hanya fact_trips dan dependensinya
dbt testJalankan semua schema test dan test kustom
dbt builddbt run + dbt test sekaligus
dbt compileHasilkan SQL tanpa mengeksekusinya (berguna untuk debugging)
dbt docs generate && dbt docs serveBangun 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 task

dbt berada di tengah pipeline. dbt tidak dapat berjalan sampai tabel raw ada dan berisi data. Skrip setup.sh menangani urutan ini.

Di halaman ini

ID