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 | ออบเจ็กต์จริง | ใช้เมื่อไร |
|---|---|---|
view | ClickHouse view | โมเดล staging: ทำความสะอาดและแปลงชนิดข้อมูลต้นทาง ไม่มีต้นทุนพื้นที่จัดเก็บ สร้างใหม่ทุกครั้งที่ query |
ephemeral | ไม่มีออบเจ็กต์ (ฝังเป็น CTE) | โมเดล intermediate ที่รวมโมเดล staging หลายตัวด้วย JOIN เลี่ยงการสร้างตารางจริงที่ซ้ำซ้อน |
table | สร้างตัวแทนใหม่ทั้งชุดใน staging relation แล้วสลับเข้าที่แบบ atomic ด้วย EXCHANGE TABLES (หรือ rename เป็นคู่ในเวอร์ชันเก่ากว่า) ตารางเดิมถูก drop หลังการสลับ | ตาราง dimension ขนาดเล็กที่ถูกแทนที่ทั้งชุดในทุกรอบ dbt run ไม่ต้องอัปเดตบางส่วน หมายเหตุ: การสร้างใหม่ทั้งชุดทำไม่ได้จริงกับตารางขนาดใหญ่ — ใช้ incremental กับตารางใดก็ตามที่เกินไม่กี่พันแถว |
incremental | CREATE TABLE ในรอบแรก ตามด้วยรูปแบบ UPDATE แบบเลือกเฉพาะในรอบถัดไป | ตาราง fact และตาราง pre-aggregation ที่ควรประมวลผลเฉพาะแถวใหม่/แถวที่เปลี่ยนในแต่ละรอบ |
materialized_view | ClickHouse 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 |
|---|---|---|
+engine | engine จัดเก็บของตาราง | ENGINE = ... ใน CREATE TABLE |
+order_by | primary 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:
- DELETE แถวออกจากตารางเป้าหมายที่
unique_keyตรงกับแถวใดก็ได้ใน batch ที่เข้ามา - 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 %}