Snowflake MigrationClickHouse Workshops

dbt trên ClickHouse

Cấu hình dbt-clickhouse: chiến lược incremental delete_insert, model ReplacingMergeTree, và refreshable materialized view.

Hướng dẫn này bao quát các mẫu đặc thù của dbt-clickhouse mà bạn sẽ dùng trong Phần 3. Hãy đọc nó sau khi hoàn thành Worksheet 1–4 và trước Worksheet 5 (dbt Model Design).

Nếu bạn đến từ dbt-snowflake, phần lớn các khái niệm dbt là giống hệt nhau — sources, refs, tests, macros, và mẫu phân tầng staging/intermediate/analytics. Điều thay đổi là tầng cấu hình đặc thù cho ClickHouse: engine, order_by, chiến lược incremental, và ngữ nghĩa FINAL.


1. Các Kiểu Materialization

dbt-clickhouse hỗ trợ năm materialization. Hãy chọn dựa trên mẫu cập nhật, không phải theo sở thích.

MaterializationĐối tượng vật lýKhi nào dùng
viewClickHouse viewStaging model: làm sạch và ép kiểu dữ liệu nguồn; không tốn chi phí lưu trữ; dựng lại ở mỗi truy vấn
ephemeralKhông có đối tượng (nhúng inline như CTE)Intermediate model kết hợp nhiều staging model qua JOIN; tránh tạo một bảng vật lý dư thừa
tableDựng bản thay thế đầy đủ trong một staging relation, rồi hoán đổi nguyên tử vào đúng vị trí bằng EXCHANGE TABLES (hoặc cặp rename trên các phiên bản cũ hơn); bảng cũ bị drop sau khi hoán đổiBảng dimension nhỏ được thay thế toàn bộ ở mỗi lần dbt run; không cần cập nhật cục bộ. Lưu ý: dựng lại toàn bộ là không khả thi với bảng lớn — hãy dùng incremental cho bất kỳ bảng nào vượt vài nghìn dòng.
incrementalCREATE TABLE ở lần chạy đầu; mẫu UPDATE có chọn lọc ở các lần chạy sauBảng fact và bảng tổng hợp trước, nơi chỉ các dòng mới/đã thay đổi cần được xử lý ở mỗi lần chạy
materialized_viewClickHouse Materialized ViewTổng hợp tự làm mới; không giống incremental của dbt. Một MV chuẩn (dựa trên trigger) chỉ chạy một lần cho mỗi INSERT và chỉ thấy được lô đó — nó không thể tính một tổng hợp trên toàn bộ vòng đời dữ liệu. Ngược lại, một REFRESHABLE MV chạy lại toàn bộ truy vấn của nó theo lịch, nên nó làm được.

Khác biệt then chốt so với Snowflake: dbt-snowflake xử lý các chi tiết lưu trữ trong nội bộ. Trong dbt-clickhouse, các model table và incremental đòi hỏi cấu hình +engine tường minh — dbt dùng nó để sinh DDL CREATE TABLE ... ENGINE = ....

View không có engine. Nếu bạn vô tình thêm +engine vào một materialization view, dbt-clickhouse sẽ bỏ qua nó. Chỉ các materialization table và incremental mới tạo ra lưu trữ bền vững cần đến engine.

Refreshable materialized view. Materialization materialized_view của dbt-clickhouse nhận một khối config refreshable — một interval (và tùy chọn randomize) — khối này phát ra trực tiếp mệnh đề REFRESH trong câu lệnh CREATE MATERIALIZED VIEW mà nó sinh ra. Model mv_live_trip_feed của lab này không đặt refreshable, đó là lý do MV mà nó dựng không có lịch làm mới.


2. Biểu Đạt Config ClickHouse Trong dbt

Các thiết lập đặc thù cho ClickHouse được biểu đạt dưới dạng dbt model config, hoặc trong dbt_project.yml (cho mặc định toàn project), hoặc trong khối config() của một model (cho override riêng từng model).

Trong 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)"

Trong khối config() của một 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)'
) }}

Cả hai cách đều tương đương. dbt_project.yml được ưu tiên cho các mẫu áp dụng toàn project; khối config() được ưu tiên cho các override riêng từng model hoặc khi bạn muốn cấu hình nằm cùng chỗ với SQL.

Các tham số config then chốt

Tham sốNó kiểm soát điều gìÁnh xạ sang ClickHouse
+engineStorage engine của bảngENGINE = ... trong CREATE TABLE
+order_byPrimary key / thứ tự sắp xếpORDER BY ... trong CREATE TABLE; mặc định là tuple() nếu bỏ trống
+unique_keyKhóa dùng để dedup trong delete_insertQuyết định những dòng nào bị xóa trước khi insert
+incremental_strategyCách các lần chạy incremental cập nhật dữ liệuĐặt thành delete_insert cho ClickHouse

Quy tắc phạm vi: Các thiết lập trong dbt_project.yml lan truyền từ cha xuống con. Khối config() ở cấp model luôn thắng project config. Hãy đặt engine phổ biến nhất làm mặc định của project, rồi override cho những model khác biệt.


3. Cơ Chế delete_insert

delete_insert là chiến lược incremental chuẩn của cộng đồng dbt-clickhouse. Nó là tương đương gần nhất với MERGE INTO của Snowflake — nhưng cơ chế thì khác.

Yêu cầu phiên bản: delete_insert dùng lightweight delete của ClickHouse, được giới thiệu ở 22.8 (thực nghiệm) và sẵn sàng cho production từ 23.3+. ClickHouse Cloud đáp ứng yêu cầu này. Để bật nó, thêm use_lw_deletes: true vào target ClickHouse trong ~/.dbt/profiles.yml của bạn, hoặc đặt allow_experimental_lightweight_delete=1 trong query_settings.

Nó làm gì

Ở mỗi lần chạy incremental:

  1. DELETE các dòng khỏi bảng đích nơi unique_key trùng với bất kỳ dòng nào trong lô dữ liệu đến
  2. INSERT toàn bộ các dòng từ lô dữ liệu đến
-- 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;

Nó khác MERGE INTO của Snowflake ra sao

Chiến lược merge của Snowflake sinh ra WHEN MATCHED THEN UPDATE / WHEN NOT MATCHED THEN INSERT theo từng dòng. ClickHouse không có câu lệnh MERGE INTO. delete_insert đạt được cùng kết quả cuối cùng — một dòng cho mỗi unique key — thông qua một lệnh xóa theo lô rồi insert toàn bộ.

Nó tương tác với ReplacingMergeTree ra sao

delete_insert là đường bảo đảm tính đúng đắn chính. ReplacingMergeTree là lưới an toàn.

Nếu một lần chạy delete_insert hoàn tất bình thường: bảng sạch (một dòng cho mỗi trip_id), không có bản trùng.

Nếu một lần chạy delete_insert bị ngắt giữa đường (crash sau DELETE, trước INSERT): dữ liệu rất có thể đang ở trạng thái không hợp lệ — những dòng đã bị xóa có thể chưa được insert lại. Lần chạy thành công tiếp theo sẽ phục hồi trạng thái đúng, nhưng đừng truy vấn bảng trong khoảng giữa một DELETE thất bại và lần chạy lại của nó.

Nếu vì lý do nào đó một lần chạy tạo ra bản trùng: merge nền của ReplacingMergeTree cuối cùng sẽ loại trùng chúng, giữ lại dòng có giá trị cột version cao nhất.

Đừng bao giờ chỉ dựa vào RMT mà không có delete_insert — merge nền là bất đồng bộ và có thể mất từ vài phút đến vài giờ trên các bảng lớn.

Khi nào dùng append

append insert các dòng mới mà không chạm tới các dòng đã có. Đó là chiến lược đúng cho các bảng chỉ insert thuần túy, nơi các dòng không bao giờ bị cập nhật — ví dụ, một event log bất biến hoặc một bảng ingest thô có ID chắc chắn duy nhất và không có chỉnh sửa. append không có yêu cầu phiên bản và không có rủi ro mutation.

Với fact_trips, append là sai: một chuyến đi có thể được chỉnh sửa sau đó (điều chỉnh giá vé, thay đổi trạng thái), nên cùng một trip_id lại đến với các giá trị mới. Với append, cả hai phiên bản tích tụ vĩnh viễn, và các phép tổng hợp (SUM giá vé, COUNT chuyến đi) sẽ đếm vượt cho tới lần merge nền RMT tiếp theo. Hãy dùng delete_insert bất cứ khi nào các dòng có thể bị cập nhật.

Tại sao không dùng chiến lược merge?

Chiến lược merge (mặc định cũ trước delete_insert) tạo một bảng tạm, nạp vào đó các dòng hiện có không thay đổi cộng với lô mới, rồi thay thế nguyên tử bảng gốc. Khác với delete_insert, nó không dùng lightweight delete — nó ghi lại toàn bộ bảng ở mỗi lần chạy incremental. Với một bảng fact_trips 50M dòng, điều này sẽ cực kỳ đắt đỏ. delete_insert chỉ xử lý các dòng trong lô hiện tại; merge chạm tới mọi dòng trong bảng. Hãy dùng delete_insert.


4. Chiến Lược Đặt FINAL

Việc loại trùng của ReplacingMergeTree diễn ra ở chế độ nền — ClickHouse merge các part một cách bất đồng bộ. Giữa các lần merge, các dòng trùng cùng tồn tại. FINAL buộc loại trùng đồng bộ tại thời điểm đọc.

FINAL thuộc về đâu trong một pipeline dbt

Ở tầng đọc từ một nguồn ReplacingMergeTree và tạo ra dữ liệu phân tích sạch.

Với 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 là điểm thực thi duy nhất cho việc loại trùng trips_raw. Mọi model hạ nguồn đọc stg_trips đều tự động nhận được dữ liệu nguồn sạch, đã loại trùng. Bạn không cần FINAL trong int_trips_enriched hay fact_trips vì chúng đọc từ stg_trips (một view, không phải một bảng RMT).

Các truy vấn dashboard và dbt test đọc trực tiếp từ fact_trips thì dùng FINAL từ bên ngoài. Bản thân model không nhúng FINAL vì nó sẽ áp dụng cho mọi lần scan bên trong truy vấn của model — kể cả subquery is_incremental() đọc max(updated_at) từ {{ this }}.

Ảnh hưởng hiệu năng của FINAL

FINAL thêm độ trễ tỉ lệ với số dòng trùng. Trên một bảng RMT được bảo trì tốt (merge nền thường xuyên), FINAL thêm rất ít overhead vì có ít bản trùng cần giải quyết. Trên một bảng mới nạp với nhiều part chưa merge, FINAL có thể chậm hơn đáng kể.

Với dbt test và các truy vấn kiểm chứng, hãy luôn dùng FINAL trên bảng RMT. Với các truy vấn benchmark nơi mục đích là so sánh độ trễ với Snowflake, các truy vấn ClickHouse đã dùng FINAL — nên phép so sánh là công bằng.


5. Macro generate_schema_name

Theo mặc định, dbt thêm tiền tố cho schema của model bằng tên target schema lấy từ profile. Nếu dbt profile của bạn trỏ tới schema nyc_taxi_ch, một model có +schema: analytics sẽ nằm trong nyc_taxi_ch_analytics — không phải analytics.

Điều này vô hại trong Snowflake (schema là namespace bên trong một database) nhưng tạo ra những cái tên khó coi trong ClickHouse, nơi schema chính là database. nyc_taxi_ch_analytics là một tên database ClickHouse hợp lệ, nhưng nó xấu hơn analytics và không khớp với các tên database đích được dùng trong kiến trúc ClickHouse của Phần 3.

Cách sửa là 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 này:

  • Trả về custom_schema_name nguyên trạng (đã chuyển thành chữ thường) khi một model chỉ định +schema: analytics
  • Trả về target schema của profile (đã chuyển thành chữ thường) cho các model không có custom schema

Filter | lower cũng bảo đảm tên schema luôn ở dạng chữ thường nhất quán, khớp với quy tắc phân biệt chữ hoa/chữ thường trong định danh của ClickHouse (Snowflake ở Phần 1 dùng | upper).

Nó nằm ở đâu: macros/generate_schema_name.sql — trong thư mục macros/ ở cấp cao nhất; dbt_project.yml đặt macro-paths: ["macros"].


Ghép Lại Với Nhau: Tóm Tắt dbt Config Cho 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 %}

Trên trang này

VI