Snowflake MigrationClickHouse Workshops

dbt di ClickHouse

Mengonfigurasi dbt-clickhouse: strategi inkremental delete_insert, model ReplacingMergeTree, dan refreshable materialized view.

Panduan ini membahas pola khusus dbt-clickhouse yang akan Anda gunakan di Bagian 3. Bacalah setelah menyelesaikan Worksheet 1–4 dan sebelum Worksheet 5 (dbt Model Design).

Jika Anda datang dari dbt-snowflake, sebagian besar konsep dbt identik — source, ref, test, macro, pola lapisan staging/intermediate/analytics. Yang berubah adalah lapisan konfigurasi khusus ClickHouse: engine, order_by, strategi inkremental, dan semantik FINAL.


1. Tipe Materialization

dbt-clickhouse mendukung lima materialization. Pilih berdasarkan pola pembaruan, bukan selera.

MaterializationObjek fisikKapan digunakan
viewView ClickHouseModel staging: membersihkan dan meng-cast tipe data sumber; tanpa biaya penyimpanan; dibangun ulang pada setiap query
ephemeralTanpa objek (disisipkan sebagai CTE)Model intermediate yang menggabungkan beberapa model staging via JOIN; menghindari pembuatan tabel fisik yang redundan
tableMembangun pengganti penuh dalam sebuah relasi staging, lalu menukarnya ke tempatnya secara atomik via EXCHANGE TABLES (atau pasangan rename pada versi lama); tabel lama di-drop setelah pertukaranTabel dimensi kecil yang digantikan sepenuhnya pada setiap dbt run; tidak perlu pembaruan sebagian. Catatan: rebuild penuh tidak layak untuk tabel besar — gunakan incremental untuk tabel apa pun yang berisi lebih dari beberapa ribu baris.
incrementalCREATE TABLE pada run pertama; pola UPDATE selektif pada run berikutnyaTabel fact dan tabel pra-agregasi di mana hanya baris baru/berubah yang perlu diproses setiap run
materialized_viewMaterialized View ClickHouseAgregat yang menyegarkan diri otomatis; tidak sama dengan incremental milik dbt. MV standar (berbasis trigger) dipicu sekali per INSERT dan hanya pernah melihat batch tersebut — MV ini tidak dapat menghitung agregat sepanjang masa. Sebaliknya, MV REFRESHABLE menjalankan ulang seluruh query-nya sesuai jadwal, sehingga bisa.

Perbedaan utama dari Snowflake: dbt-snowflake menangani detail penyimpanan secara internal. Di dbt-clickhouse, model table dan incremental memerlukan konfigurasi +engine yang eksplisit — dbt memakainya untuk menghasilkan DDL CREATE TABLE ... ENGINE = ....

View tidak memiliki engine. Jika Anda tanpa sengaja menambahkan +engine pada materialization view, dbt-clickhouse akan mengabaikannya. Hanya materialization table dan incremental yang membuat penyimpanan persisten yang membutuhkan engine.

Refreshable materialized view. Materialization materialized_view pada dbt-clickhouse menerima blok config refreshable — sebuah interval (dan opsional randomize) — yang memancarkan klausa REFRESH langsung di dalam pernyataan CREATE MATERIALIZED VIEW yang dihasilkannya. Model mv_live_trip_feed di lab ini tidak menyetel refreshable, itulah sebabnya MV yang dibangunnya tidak memiliki jadwal refresh.


2. Menyatakan Config ClickHouse di dbt

Setelan khusus ClickHouse dinyatakan sebagai config model dbt, baik di dbt_project.yml (untuk default seluruh proyek) atau di blok config() sebuah model (untuk override khusus model).

Di dbt_project.yml

models:
  your_project:
    analytics:
      +schema: analytics
      +materialized: table
      +engine: "MergeTree()"          # default for all analytics tables

      fact_trips:
        +materialized: incremental
        +engine: "ReplacingMergeTree(updated_at)"   # overrides the default
        +incremental_strategy: delete_insert
        +unique_key: trip_id
        +order_by: "(toStartOfMonth(pickup_at), pickup_at, trip_id)"

Di blok config() sebuah model

{{ config(
    materialized         = 'incremental',
    engine               = 'ReplacingMergeTree(updated_at)',
    incremental_strategy = 'delete_insert',
    unique_key           = 'trip_id',
    order_by             = '(toStartOfMonth(pickup_at), pickup_at, trip_id)'
) }}

Kedua pendekatan setara. dbt_project.yml lebih disukai untuk pola seluruh proyek; blok config() lebih disukai untuk override khusus model atau ketika Anda ingin konfigurasi berada bersebelahan dengan SQL-nya.

Parameter config utama

ParameterYang dikendalikanPemetaan ClickHouse
+engineEngine penyimpanan tabelENGINE = ... di CREATE TABLE
+order_byPrimary key / urutan sortORDER BY ... di CREATE TABLE; default ke tuple() jika dihilangkan
+unique_keyKey untuk dedup delete_insertMenentukan baris mana yang dihapus sebelum insert
+incremental_strategyCara run inkremental memperbarui dataSetel ke delete_insert untuk ClickHouse

Aturan scoping: Setelan di dbt_project.yml mengalir (cascade) dari induk ke anak. Blok config() di tingkat model selalu menang atas config proyek. Setel engine yang paling umum sebagai default proyek, lalu override untuk model yang berbeda.


3. Mekanika delete_insert

delete_insert adalah strategi inkremental standar komunitas dbt-clickhouse. Ini padanan terdekat dari MERGE INTO milik Snowflake — tetapi mekanikanya berbeda.

Persyaratan versi: delete_insert menggunakan lightweight delete ClickHouse, diperkenalkan di 22.8 (eksperimental) dan siap produksi di 23.3+. ClickHouse Cloud memenuhi persyaratan ini. Untuk mengaktifkannya, tambahkan use_lw_deletes: true ke target ClickHouse di ~/.dbt/profiles.yml Anda, atau setel allow_experimental_lightweight_delete=1 di query_settings.

Apa yang dilakukannya

Pada setiap run inkremental:

  1. DELETE baris dari tabel target di mana unique_key cocok dengan baris mana pun dalam batch yang masuk
  2. INSERT semua baris dari batch yang masuk
-- Step 1: dbt generates this DELETE
ALTER TABLE analytics.fact_trips
DELETE WHERE trip_id IN (SELECT trip_id FROM incoming_batch);

-- Step 2: dbt generates this INSERT
INSERT INTO analytics.fact_trips
SELECT * FROM incoming_batch;

Perbedaannya dari MERGE INTO Snowflake

Strategi merge Snowflake menghasilkan WHEN MATCHED THEN UPDATE / WHEN NOT MATCHED THEN INSERT baris per baris. ClickHouse tidak memiliki pernyataan MERGE INTO. delete_insert mencapai hasil akhir yang sama — satu baris per unique key — melalui delete batch yang diikuti insert penuh.

Cara interaksinya dengan ReplacingMergeTree

delete_insert adalah jalur kebenaran utama. ReplacingMergeTree adalah jaring pengamannya.

Jika sebuah run delete_insert selesai normal: tabel bersih (satu baris per trip_id), tanpa duplikat.

Jika sebuah run delete_insert terputus di tengah jalan (crash setelah DELETE, sebelum INSERT): data kemungkinan besar berada dalam keadaan tidak valid — baris yang telah dihapus mungkin belum di-insert ulang. Run berikutnya yang berhasil akan memulihkan keadaan yang benar, tetapi jangan mengquery tabel tersebut di antara DELETE yang gagal dan run ulangnya.

Jika sebuah run menghasilkan duplikat karena alasan apa pun: background merge ReplacingMergeTree pada akhirnya akan mendeduplikasinya, mempertahankan baris dengan nilai kolom versi tertinggi.

Jangan pernah mengandalkan RMT saja tanpa delete_insert — background merge bersifat asinkron dan bisa memakan waktu beberapa menit hingga beberapa jam pada tabel besar.

Kapan menggunakan append

append meng-insert baris baru tanpa menyentuh baris yang sudah ada. Ini strategi yang benar untuk tabel yang murni insert-only di mana baris tidak pernah diperbarui — misalnya, event log yang immutable atau tabel ingest raw dengan ID yang dijamin unik dan tanpa koreksi. append tidak punya persyaratan versi dan tidak punya risiko mutasi.

Untuk fact_trips, append salah: sebuah trip dapat dikoreksi setelah fakta (penyesuaian tarif, perubahan status), sehingga trip_id yang sama datang lagi dengan nilai baru. Dengan append, kedua versi menumpuk secara permanen, dan agregat (SUM tarif, COUNT trip) menghitung berlebih sampai background merge RMT berikutnya. Gunakan delete_insert kapan pun baris dapat diperbarui.

Mengapa bukan strategi merge?

Strategi merge (default lama sebelum delete_insert) membuat tabel temporer, mengisinya dengan baris lama yang tidak berubah ditambah batch baru, lalu menggantikan tabel asli secara atomik. Tidak seperti delete_insert, strategi ini tidak memakai lightweight delete — ia menulis ulang seluruh tabel pada setiap run inkremental. Untuk tabel fact_trips dengan 50M baris, ini akan sangat mahal. delete_insert hanya memproses baris dalam batch saat ini; merge menyentuh setiap baris dalam tabel. Gunakan delete_insert.


4. Strategi Penempatan FINAL

Deduplikasi ReplacingMergeTree terjadi di background — ClickHouse me-merge part secara asinkron. Di antara merge, baris duplikat hidup bersamaan. FINAL memaksa deduplikasi sinkron pada saat baca.

Di mana FINAL seharusnya berada dalam pipeline dbt

Di lapisan yang membaca dari sumber ReplacingMergeTree dan menghasilkan data analitis yang bersih.

Untuk workload NYC Taxi:

trips_raw (RMT)
    ↓
stg_trips (view): SELECT ... FROM trips_raw FINAL   ← FINAL goes here
    ↓
int_trips_enriched (ephemeral CTE)
    ↓
fact_trips (incremental, RMT)                        ← NO FINAL in model
    ↓
Dashboard queries: SELECT ... FROM fact_trips FINAL  ← FINAL goes here (externally)

stg_trips adalah satu-satunya titik penegakan untuk deduplikasi trips_raw. Setiap model downstream yang membaca stg_trips otomatis mendapat data sumber yang bersih dan terdeduplikasi. Anda tidak memerlukan FINAL di int_trips_enriched atau fact_trips karena keduanya membaca dari stg_trips (sebuah view, bukan tabel RMT).

Query dashboard dan test dbt yang membaca langsung dari fact_trips memakai FINAL secara eksternal. Model itu sendiri tidak menanamkan FINAL karena FINAL akan berlaku pada setiap scan di dalam query model — termasuk subquery is_incremental() yang membaca max(updated_at) dari {{ this }}.

Dampak FINAL terhadap performa

FINAL menambah latensi yang sebanding dengan jumlah baris duplikat. Pada tabel RMT yang terpelihara baik (background merge sering terjadi), FINAL menambah overhead minimal karena sedikit duplikat yang harus diselesaikan. Pada tabel yang baru dimuat dengan banyak part yang belum di-merge, FINAL bisa jauh lebih lambat.

Untuk test dbt dan query verifikasi, selalu gunakan FINAL pada tabel RMT. Untuk query benchmark yang tujuannya membandingkan latensi terhadap Snowflake, query ClickHouse sudah memakai FINAL — sehingga perbandingannya adil.


5. Macro generate_schema_name

Secara default, dbt memberi awalan nama target schema dari profile pada schema model. Jika profile dbt Anda menargetkan schema nyc_taxi_ch, model dengan +schema: analytics mendarat di nyc_taxi_ch_analytics — bukan analytics.

Ini tidak berbahaya di Snowflake (schema adalah namespace di dalam sebuah database) tetapi menciptakan nama yang aneh di ClickHouse, di mana schema adalah database. nyc_taxi_ch_analytics adalah nama database ClickHouse yang valid, tetapi lebih jelek daripada analytics dan tidak cocok dengan nama database target yang dipakai dalam arsitektur ClickHouse di Bagian 3.

Solusinya adalah override macro generate_schema_name:

-- macros/generate_schema_name.sql
{% macro generate_schema_name(custom_schema_name, node) -%}
  {%- if custom_schema_name is none -%}
    {{ target.schema | lower }}
  {%- else -%}
    {{ custom_schema_name | lower }}
  {%- endif -%}
{%- endmacro %}

Macro ini:

  • Mengembalikan custom_schema_name apa adanya (dalam huruf kecil) ketika sebuah model menyatakan +schema: analytics
  • Mengembalikan target schema dari profile (dalam huruf kecil) untuk model tanpa custom schema

Filter | lower juga memastikan nama schema konsisten dalam huruf kecil, sesuai aturan identifier ClickHouse yang case-sensitive (Snowflake di Bagian 1 memakai | upper).

Lokasinya: macros/generate_schema_name.sql — di direktori macros/ tingkat atas; dbt_project.yml menyetel macro-paths: ["macros"].


Menyatukan Semuanya: Ringkasan Config dbt NYC Taxi

# dbt_project.yml (abbreviated)
models:
  nyc_taxi_dbt_ch:
    staging:
      +schema: staging
      +materialized: view           # no engine — views need none

    intermediate:
      +schema: staging
      +materialized: ephemeral      # inlined as CTE

    analytics:
      +schema: analytics
      +materialized: table
      +engine: "MergeTree()"        # default for dim_* tables

      fact_trips:
        +materialized: incremental
        +engine: "ReplacingMergeTree(updated_at)"
        +incremental_strategy: delete_insert
        +unique_key: trip_id

      agg_hourly_zone_trips:
        +materialized: incremental
        +engine: "ReplacingMergeTree(updated_at)"
        +incremental_strategy: delete_insert
        +unique_key: [hour_bucket, zone_id]
-- stg_trips.sql (staging view — the FINAL enforcement point)
SELECT ... FROM {{ source('raw', 'trips_raw') }} FINAL

-- fact_trips.sql (incremental — no FINAL in model body)
SELECT ... FROM {{ ref('int_trips_enriched') }}
{% if is_incremental() %}
WHERE updated_at > (SELECT max(updated_at) FROM {{ this }})
{% endif %}

Di halaman ini

ID