Snowflake MigrationClickHouse Workshops

dbt บน ClickHouse

การตั้งค่า dbt-clickhouse: incremental strategy แบบ delete_insert, โมเดล ReplacingMergeTree และ refreshable materialized view

คู่มือนี้ครอบคลุมรูปแบบเฉพาะของ dbt-clickhouse ที่คุณจะใช้ใน Part 3 อ่านหลังจากทำเวิร์กชีต 1–4 เสร็จ และก่อนเวิร์กชีต 5 (การออกแบบโมเดล dbt)

ถ้าคุณมาจาก dbt-snowflake แนวคิดของ dbt ส่วนใหญ่เหมือนกันทุกอย่าง — sources, refs, tests, macros, รูปแบบชั้น staging/intermediate/analytics สิ่งที่เปลี่ยนคือชั้นการตั้งค่าที่เจาะจง ClickHouse: engine, order_by, incremental strategy และความหมายของ FINAL


1. ประเภทของ Materialization

dbt-clickhouse รองรับ materialization ห้าแบบ เลือกตามรูปแบบการอัปเดต ไม่ใช่ตามความชอบ

Materializationออบเจ็กต์จริงใช้เมื่อไร
viewClickHouse viewโมเดล staging: ทำความสะอาดและแปลงชนิดข้อมูลต้นทาง ไม่มีต้นทุนพื้นที่จัดเก็บ สร้างใหม่ทุกครั้งที่ query
ephemeralไม่มีออบเจ็กต์ (ฝังเป็น CTE)โมเดล intermediate ที่รวมโมเดล staging หลายตัวด้วย JOIN เลี่ยงการสร้างตารางจริงที่ซ้ำซ้อน
tableสร้างตัวแทนใหม่ทั้งชุดใน staging relation แล้วสลับเข้าที่แบบ atomic ด้วย EXCHANGE TABLES (หรือ rename เป็นคู่ในเวอร์ชันเก่ากว่า) ตารางเดิมถูก drop หลังการสลับตาราง dimension ขนาดเล็กที่ถูกแทนที่ทั้งชุดในทุกรอบ dbt run ไม่ต้องอัปเดตบางส่วน หมายเหตุ: การสร้างใหม่ทั้งชุดทำไม่ได้จริงกับตารางขนาดใหญ่ — ใช้ incremental กับตารางใดก็ตามที่เกินไม่กี่พันแถว
incrementalCREATE TABLE ในรอบแรก ตามด้วยรูปแบบ UPDATE แบบเลือกเฉพาะในรอบถัดไปตาราง fact และตาราง pre-aggregation ที่ควรประมวลผลเฉพาะแถวใหม่/แถวที่เปลี่ยนในแต่ละรอบ
materialized_viewClickHouse Materialized Viewค่า aggregate ที่รีเฟรชอัตโนมัติ ไม่เหมือนกับ incremental ของ dbt MV มาตรฐาน (แบบ trigger) จะทำงานหนึ่งครั้งต่อหนึ่ง INSERT และเห็นเฉพาะ batch นั้น — มันไม่สามารถคำนวณค่า aggregate ตลอดอายุข้อมูลได้ ส่วน REFRESHABLE MV จะรัน query ทั้งชุดใหม่ตามตารางเวลา จึงทำได้

ความต่างสำคัญจาก Snowflake: dbt-snowflake จัดการรายละเอียดพื้นที่จัดเก็บภายในตัวเอง ใน dbt-clickhouse โมเดล table และ incremental ต้องมีการตั้งค่า +engine อย่างชัดเจน — dbt ใช้ค่านี้สร้าง DDL CREATE TABLE ... ENGINE = ...

View ไม่มี engine ถ้าคุณเผลอเพิ่ม +engine ให้ materialization แบบ view dbt-clickhouse จะไม่สนใจมัน มีเพียง materialization แบบ table และ incremental เท่านั้นที่สร้างพื้นที่จัดเก็บถาวรซึ่งต้องมี engine

Refreshable materialized view materialization materialized_view ของ dbt-clickhouse รับบล็อกการตั้งค่า refreshable — ค่า interval (และอาจมี randomize) — ซึ่งจะใส่ประโยค REFRESH ลงในคำสั่ง CREATE MATERIALIZED VIEW ที่มันสร้างโดยตรง โมเดล mv_live_trip_feed ของ lab นี้ไม่ได้ตั้ง refreshable ซึ่งเป็นเหตุผลที่ MV ที่มันสร้างไม่มีตารางเวลารีเฟรช


2. การแสดงค่าตั้ง ClickHouse ใน dbt

การตั้งค่าที่เจาะจง ClickHouse แสดงในรูปของ dbt model config ไม่ว่าจะใน dbt_project.yml (สำหรับค่าเริ่มต้นทั่วโปรเจกต์) หรือในบล็อก config() ของโมเดล (สำหรับการเขียนทับเฉพาะโมเดล)

ใน 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)"

ในบล็อก config() ของโมเดล

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

ทั้งสองวิธีเทียบเท่ากัน dbt_project.yml เหมาะกับรูปแบบที่ใช้ทั่วโปรเจกต์ ส่วนบล็อก config() เหมาะกับการเขียนทับเฉพาะโมเดล หรือเมื่อคุณต้องการให้การตั้งค่าอยู่ที่เดียวกับ SQL

พารามิเตอร์ config ที่สำคัญ

พารามิเตอร์ควบคุมอะไรการแมปไปที่ ClickHouse
+engineengine จัดเก็บของตารางENGINE = ... ใน CREATE TABLE
+order_byprimary key / ลำดับการเรียงORDER BY ... ใน CREATE TABLE ค่าเริ่มต้นเป็น tuple() ถ้าไม่ระบุ
+unique_keyคีย์สำหรับการกำจัดข้อมูลซ้ำของ delete_insertกำหนดว่าจะลบแถวใดก่อน insert
+incremental_strategyวิธีที่รอบ incremental อัปเดตข้อมูลตั้งเป็น delete_insert สำหรับ ClickHouse

กฎขอบเขต: การตั้งค่าใน dbt_project.yml ไหลจากระดับแม่ลงสู่ระดับลูก บล็อก config() ระดับโมเดลชนะค่าตั้งระดับโปรเจกต์เสมอ ให้ตั้ง engine ที่ใช้บ่อยที่สุดเป็นค่าเริ่มต้นของโปรเจกต์ แล้วเขียนทับเฉพาะโมเดลที่ต่างออกไป


3. กลไกของ delete_insert

delete_insert คือ incremental strategy มาตรฐานของชุมชน dbt-clickhouse มันเทียบเท่า MERGE INTO ของ Snowflake ได้ใกล้เคียงที่สุด — แต่กลไกต่างกัน

ข้อกำหนดเวอร์ชัน: delete_insert ใช้ lightweight delete ของ ClickHouse ซึ่งเริ่มมีใน 22.8 (ทดลอง) และพร้อมใช้งานจริงใน 23.3+ ClickHouse Cloud ผ่านข้อกำหนดนี้ หากต้องการเปิดใช้ ให้เพิ่ม use_lw_deletes: true ที่ target ของ ClickHouse ใน ~/.dbt/profiles.yml ของคุณ หรือตั้ง allow_experimental_lightweight_delete=1 ใน query_settings

มันทำอะไร

ในแต่ละรอบ incremental:

  1. DELETE แถวออกจากตารางเป้าหมายที่ unique_key ตรงกับแถวใดก็ได้ใน batch ที่เข้ามา
  2. INSERT ทุกแถวจาก batch ที่เข้ามา
-- 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;

ต่างจาก MERGE INTO ของ Snowflake อย่างไร

กลยุทธ์ merge ของ Snowflake สร้าง WHEN MATCHED THEN UPDATE / WHEN NOT MATCHED THEN INSERT แบบแถวต่อแถว ClickHouse ไม่มีคำสั่ง MERGE INTO delete_insert ให้ผลลัพธ์สุดท้ายเหมือนกัน — หนึ่งแถวต่อหนึ่ง unique key — โดยใช้การลบเป็น batch แล้วตามด้วยการ insert ทั้งชุด

มันทำงานร่วมกับ ReplacingMergeTree อย่างไร

delete_insert คือ เส้นทางความถูกต้องหลัก ส่วน ReplacingMergeTree คือ ตาข่ายรองรับ

ถ้ารอบ delete_insert ทำงานจบตามปกติ: ตารางสะอาด (หนึ่งแถวต่อหนึ่ง trip_id) ไม่มีข้อมูลซ้ำ

ถ้ารอบ delete_insert ถูกขัดจังหวะกลางทาง (crash หลัง DELETE ก่อน INSERT): ข้อมูลมีแนวโน้มจะอยู่ในสภาพที่ไม่ถูกต้อง — แถวที่ถูกลบอาจยังไม่ถูก insert กลับ รอบถัดไปที่สำเร็จจะกู้สภาพให้ถูกต้อง แต่อย่า query ตารางในช่วงระหว่าง DELETE ที่ล้มเหลวกับการรันซ้ำ

ถ้ารอบใดสร้างข้อมูลซ้ำด้วยเหตุผลใดก็ตาม: การ merge เบื้องหลังของ ReplacingMergeTree จะกำจัดข้อมูลซ้ำในที่สุด โดยเก็บแถวที่มีค่าคอลัมน์ version สูงสุดไว้

อย่าพึ่ง RMT เพียงลำพังโดยไม่มี delete_insert — การ merge เบื้องหลังทำงานแบบ asynchronous และอาจใช้เวลาหลายนาทีถึงหลายชั่วโมงกับตารางขนาดใหญ่

เมื่อไรควรใช้ append

append insert แถวใหม่โดยไม่แตะแถวที่มีอยู่ มันเป็นกลยุทธ์ที่ถูกต้องสำหรับตารางที่ insert-only ล้วน ๆ ซึ่งแถวไม่เคยถูกอัปเดต — ตัวอย่างเช่น event log ที่ไม่เปลี่ยนแปลง หรือตาราง raw ingest ที่รับประกัน ID ไม่ซ้ำและไม่มีการแก้ไข append ไม่มีข้อกำหนดเวอร์ชันและไม่มีความเสี่ยงจาก mutation

สำหรับ fact_trips append ผิด: ทริปหนึ่งอาจถูกแก้ไขภายหลัง (ปรับค่าโดยสาร เปลี่ยนสถานะ) ดังนั้น trip_id เดิมจะมาถึงอีกครั้งพร้อมค่าใหม่ ด้วย append ทั้งสองเวอร์ชันจะสะสมอยู่อย่างถาวร และค่า aggregate (SUM ของค่าโดยสาร, COUNT ของทริป) จะนับเกินจนกว่าจะถึงการ merge เบื้องหลังของ RMT ครั้งถัดไป ใช้ delete_insert เมื่อใดก็ตามที่แถวสามารถถูกอัปเดตได้

ทำไมไม่ใช้กลยุทธ์ merge

กลยุทธ์ merge (ค่าเริ่มต้นเดิมก่อนจะมี delete_insert) สร้างตารางชั่วคราว เติมด้วยแถวเดิมที่ไม่เปลี่ยนบวกกับ batch ใหม่ แล้วแทนที่ตารางเดิมแบบ atomic ต่างจาก delete_insert มันไม่ใช้ lightweight delete — มันเขียนตารางทั้งตารางใหม่ในทุกรอบ incremental สำหรับตาราง fact_trips ขนาด 50M แถว นี่จะแพงมาก delete_insert ประมวลผลเฉพาะแถวใน batch ปัจจุบัน ส่วน merge แตะทุกแถวในตาราง ให้ใช้ delete_insert


4. กลยุทธ์การวาง FINAL

การกำจัดข้อมูลซ้ำของ ReplacingMergeTree เกิดขึ้นเบื้องหลัง — ClickHouse merge parts แบบ asynchronous ระหว่างการ merge แถวที่ซ้ำจะอยู่ร่วมกัน FINAL บังคับให้กำจัดข้อมูลซ้ำแบบ synchronous ตอนอ่าน

FINAL ควรอยู่ตรงไหนใน pipeline ของ dbt

ที่ชั้นซึ่งอ่านจาก source ที่เป็น ReplacingMergeTree และผลิตข้อมูลเชิงวิเคราะห์ที่สะอาด

สำหรับ 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 เป็นจุดบังคับใช้จุดเดียวสำหรับการกำจัดข้อมูลซ้ำของ trips_raw ทุกโมเดลปลายทางที่อ่าน stg_trips จะได้ข้อมูลต้นทางที่สะอาดและกำจัดข้อมูลซ้ำแล้วโดยอัตโนมัติ คุณไม่ต้องใช้ FINAL ใน int_trips_enriched หรือ fact_trips เพราะทั้งคู่อ่านจาก stg_trips (ซึ่งเป็น view ไม่ใช่ตาราง RMT)

Query ของ dashboard และการทดสอบ dbt ที่อ่านจาก fact_trips โดยตรงจะใช้ FINAL จากภายนอก ตัวโมเดลเองไม่ฝัง FINAL เพราะมันจะถูกนำไปใช้กับทุกการสแกนภายใน query ของโมเดล — รวมถึง subquery ของ is_incremental() ที่อ่าน max(updated_at) จาก {{ this }}

ผลกระทบด้านประสิทธิภาพของ FINAL

FINAL เพิ่ม latency ตามสัดส่วนของจำนวนแถวที่ซ้ำ บนตาราง RMT ที่ดูแลดี (merge เบื้องหลังบ่อย) FINAL เพิ่ม overhead น้อยมาก เพราะมีข้อมูลซ้ำให้แก้ไขไม่มาก บนตารางที่โหลดมาสด ๆ และมี parts ที่ยังไม่ merge จำนวนมาก FINAL อาจช้าลงอย่างมีนัยสำคัญ

สำหรับการทดสอบ dbt และ query ตรวจสอบ ให้ใช้ FINAL กับตาราง RMT เสมอ สำหรับ query benchmark ที่ประเด็นคือการเปรียบเทียบ latency กับ Snowflake นั้น query ของ ClickHouse ใช้ FINAL อยู่แล้ว — ดังนั้นการเปรียบเทียบจึงยุติธรรม


5. มาโคร generate_schema_name

โดยค่าเริ่มต้น dbt เติมชื่อ target schema จาก profile ไว้ข้างหน้า schema ของโมเดล ถ้า profile ของ dbt คุณชี้ไปที่ schema nyc_taxi_ch โมเดลที่มี +schema: analytics จะลงใน nyc_taxi_ch_analytics — ไม่ใช่ analytics

เรื่องนี้ไม่มีผลเสียใน Snowflake (schema เป็น namespace ภายในฐานข้อมูล) แต่สร้างชื่อที่เก้กังใน ClickHouse ที่ schema คือ ฐานข้อมูล nyc_taxi_ch_analytics เป็นชื่อฐานข้อมูล ClickHouse ที่ใช้ได้ แต่มันดูไม่สวยเท่า analytics และไม่ตรงกับชื่อฐานข้อมูลเป้าหมายที่ใช้ในสถาปัตยกรรม ClickHouse ของ Part 3

วิธีแก้คือเขียนทับมาโคร 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 %}

มาโครนี้:

  • คืนค่า custom_schema_name ตามที่เป็น (แปลงเป็นตัวพิมพ์เล็ก) เมื่อโมเดลระบุ +schema: analytics
  • คืนค่า target schema ของ profile (แปลงเป็นตัวพิมพ์เล็ก) สำหรับโมเดลที่ไม่มี custom schema

ตัวกรอง | lower ยังทำให้ชื่อ schema เป็นตัวพิมพ์เล็กอย่างสม่ำเสมอ ซึ่งตรงกับกฎการแยกตัวพิมพ์ของตัวระบุใน ClickHouse (Part 1 ที่เป็น Snowflake ใช้ | upper)

อยู่ที่ไหน: macros/generate_schema_name.sql — ในไดเรกทอรี macros/ ระดับบนสุด โดย dbt_project.yml ตั้ง macro-paths: ["macros"]


ประกอบร่างเข้าด้วยกัน: สรุปค่าตั้ง 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 %}

ในหน้านี้

TH