04 Dựng lại pipeline dbt
Dựng lại pipeline Medallion trên ClickHouse với dbt-clickhouse — các model incremental delete_insert, ReplacingMergeTree, materialized view có thể refresh — và tạo dictionary vùng.
Điểm khởi đầu
Module 03 đã hoàn tất: service ClickHouse Cloud đang hoạt động và truy cập được, và
.clickhouse_state đã được setup.sh ghi ra đĩa với CLICKHOUSE_HOST và
CLICKHOUSE_PORT. Mọi bảng đích và view staging đều tồn tại — default.trips_raw, hai view
staging (stg_trips, stg_taxi_zones), sáu bảng analytics (fact_trips,
agg_hourly_zone_trips, dim_taxi_zones, dim_payment_type, dim_vendor, dim_date), và
materialized view có thể refresh analytics.mv_live_trip_feed — lệnh dbt run của module 03
đã dựng cả bảy. Mọi bảng analytics vẫn còn rỗng, trừ mv_live_trip_feed, vốn đã chứa một
dòng ảnh chụp từ lượt build đó. Ngoài ra chỉ có default.trips_raw là có dữ liệu: khoảng 50
triệu dòng. Producer Snowflake vẫn đang chạy, nên ClickHouse tụt lại sau Snowflake khoảng bằng
độ dài của cửa sổ di chuyển. Hãy dự trù khoảng 30 phút.
Vì sao
Module 03 đã chứng minh ClickHouse có thể chứa 50 triệu dòng. Nó chưa chứng minh rằng pipeline có thể chạy trên ClickHouse — các view staging, bảng fact incremental, các lượt nạp lại bảng chiều, các test bắt được một model hỏng trước khi một đối tác kịp nhìn thấy. Đó là những gì module này dựng lại: cùng các model Medallion từ module 01, được diễn đạt bằng dbt-clickhouse thay vì dbt-snowflake, chạy trên các bảng mà module 03 đã tạo.
Không có gì trong logic của model thay đổi — stg_trips vẫn ép kiểu và trích xuất JSON,
int_trips_enriched vẫn join các bảng chiều, fact_trips vẫn kết thúc với một dòng cho mỗi
chuyến đi. Điều thay đổi là tầng materialization bên dưới: không MERGE INTO, không Snowflake
Task, không cluster_by. Đây là module cho thấy cuộc di chuyển không phải là một lần dump dữ
liệu duy nhất — pipeline mà đội của một đối tác chạy mỗi ngày, theo cùng lịch trình, được gác
bởi cùng một lượt dbt test, vẫn tiếp tục hoạt động khi warehouse bên dưới nó thay đổi.
Khái niệm — bên dưới lớp vỏ
Tập model khác pipeline Snowflake ở bốn điểm. Tài liệu cấu hình đầy đủ là
dbt trên ClickHouse — đây là bản ngắn
bạn cần trước khi chạy dbt run ở Bước 1. Về pipeline nguồn mà các model này thay thế, xem
dbt trên Snowflake.
1. delete_insert thay cho MERGE. ClickHouse không có câu lệnh MERGE INTO. Ở nơi mà
pipeline Snowflake dùng incremental_strategy: merge để upsert fact_trips và
agg_hourly_zone_trips, các model ClickHouse dùng incremental_strategy: delete_insert: dbt
xóa các dòng khớp unique_key cho lô dữ liệu đang đến, rồi chèn lô đó vào. Với fact_trips,
unique_key là trip_id, và filter incremental đặt mốc nước trên updated_at thay vì
pickup_at — một lần sửa giá cước sẽ chèn lại cùng trip_id với cùng pickup_at nhưng
updated_at mới hơn, nên đặt mốc nước trên pickup_at sẽ âm thầm bỏ sót nó.
2. ReplacingMergeTree là lưới an toàn bên dưới delete_insert, không phải là vật thay thế
cho nó. Cả hai model incremental đều được khai báo là ReplacingMergeTree(updated_at). Nếu
một lượt delete_insert hoàn tất bình thường, bảng đã có sẵn một dòng cho mỗi khóa và engine
không còn gì phải dọn. Nếu một lượt bị ngắt giữa đường — sập sau khi xóa, trước khi chèn — thì
các lượt merge chạy nền cuối cùng sẽ loại trùng bất kỳ dòng sót lại, giữ lại dòng có
updated_at cao nhất. Đừng bao giờ chỉ dựa vào ReplacingMergeTree để làm công việc loại
trùng mà delete_insert phải làm: các lượt merge nền là bất đồng bộ và có thể trễ từ vài phút
đến vài giờ trên một bảng cỡ này.
3. Materialized view có thể refresh thay cho task theo lịch. Pipeline Snowflake dùng một
Task theo lịch chạy một stored procedure để giữ một bảng tổng hợp cuốn luôn cập nhật. Thay vào
đó, dự án dbt của ClickHouse khai báo mv_live_trip_feed với
materialized = 'materialized_view' và engine = 'ReplacingMergeTree(refreshed_at)' — được
dbt run dựng thành một materialized view có thể refresh, so với hơn 30 dòng DDL
CREATE TASK của Snowflake cho cùng tác dụng. Module này không bật một chu kỳ refresh nào;
làm vậy đòi hỏi một câu lệnh
ALTER TABLE analytics.mv_live_trip_feed MODIFY REFRESH EVERY ... thủ công mà lab không viết
script cho (xem module 05 để biết vì sao).
4. mv_live_trip_feed hoàn toàn không có đối ứng ở Snowflake. Nó không phải là bản dịch
của một model có sẵn — nó là năng lực mới mà cuộc di chuyển mang lại. Một materialized view
tiêu chuẩn của ClickHouse chỉ chạy một lần cho mỗi lệnh INSERT và chỉ nhìn thấy các dòng trong
lô đó, nên nó không thể tính đúng một giá trị tổng hợp trên toàn bộ vòng đời như tổng số
chuyến đi hay giá cước trung bình. Một materialized view REFRESHABLE thì chạy lại toàn bộ truy
vấn của nó — ở đây là SELECT ... FROM {{ ref('fact_trips') }} — theo lịch, nên mỗi lượt
refresh đều thấy toàn bảng. Phía Snowflake của workshop này chưa bao giờ có lựa chọn đó.
Các model mà dbt thực sự dựng, và thời điểm mỗi model có dữ liệu:
| Model | Tầng | Materialization | Có dữ liệu từ | Ghi chú |
|---|---|---|---|---|
stg_trips | staging | View | Bước 1 (mỗi lượt chạy) | Ép kiểu, JSONExtract* cho trip_metadata |
stg_taxi_zones | staging | View | Bước 1 (mỗi lượt chạy) | Chuyển tiếp bảng chiều vùng |
int_trips_enriched | staging | Ephemeral | — (được nội tuyến thành một CTE) | Toàn bộ các phép join bảng chiều; không có bảng vật lý |
fact_trips | analytics | Incremental | Bước 1 | delete_insert theo khóa trip_id, mốc nước trên updated_at; ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id) |
agg_hourly_zone_trips | analytics | Incremental | Module 05, sau cutover | Cửa sổ cuốn 2 giờ; filter incremental chỉ khớp các dòng từ producer chạy trực tiếp — xem Bước 1 bên dưới |
dim_taxi_zones | analytics | Table | Bước 1 | Nạp lại toàn bộ mỗi lượt chạy; là nguồn cho dictionary vùng ở Bước 2 |
dim_payment_type | analytics | Table | Bước 1 | Nạp lại toàn bộ mỗi lượt chạy |
dim_vendor | analytics | Table | Bước 1 | Nạp lại toàn bộ mỗi lượt chạy |
dim_date | analytics | Table | Bước 1 | Trục ngày tĩnh, 2009-2029; nạp lại toàn bộ mỗi lượt chạy |
mv_live_trip_feed | analytics | Materialized view (có thể refresh) | Lệnh dbt run của module 03; chu kỳ refresh chưa bao giờ được bật | Không có đối ứng ở Snowflake — xem điểm 4 ở trên |
Bước 1 — Nạp dữ liệu cho tầng analytics
Kích hoạt venv dbt-clickhouse bạn đã dựng ở module 00, rồi chạy dbt run lần thứ hai. Module
03 đã chạy nó một lần trên các bảng rỗng để tạo schema; lượt chạy này có dữ liệu thật phía sau
— trips_raw giờ chứa 50 triệu dòng.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .venv/bin/activate
source .env && source .clickhouse_state
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt rundbt đọc thông tin kết nối ClickHouse từ ~/.dbt/profiles.yml, dựa trên
dbt/nyc_taxi_dbt_ch/profiles.yml.example — cùng profile mà module 03 đã dùng để tạo các
schema rỗng.
Mong đợi: khoảng 8-12 phút (50 triệu dòng được các model incremental xử lý).
agg_hourly_zone_trips sẽ rỗng sau lượt chạy này — đó là điều được dự kiến, không phải một
lỗi. Filter incremental của nó là WHERE pickup_at >= now() - INTERVAL 2 HOUR, chỉ khớp các
dòng do producer chạy trực tiếp ghi vào. Mọi dòng bạn vừa di chuyển đều là dữ liệu lịch sử,
nên không dòng nào nằm trong cửa sổ 2 giờ tính từ lúc này. Bảng này vẫn rỗng cho đến khi bước
cutover của module 05 khởi động producer ClickHouse — đừng mất thời gian debug nó như một
pipeline bị hỏng.
Sau đó chạy bộ test:
dbt testMong đợi: mọi test đều pass.
Kiểm chứng:
SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 million
SELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0 (normal — populated after cutover in module 05)
SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265Bước 2 — Tạo dictionary vùng
analytics.dim_taxi_zones giờ đã có toàn bộ 265 vùng NYC TLC. Hãy dựng
analytics.taxi_zones_dict, một dictionary trong bộ nhớ dựa trên bảng đó, để các truy vấn
phía sau có thể tra borough của một vùng bằng dictGet() thay vì một phép JOIN.
Dictionary hơn gì so với một phép join. Một dictionary được nạp vào bộ nhớ một lần và giữ
nóng ở đó; một lượt tra cứu vào nó gần như miễn phí ở mọi truy vấn sau đó. Một phép JOIN với
dim_taxi_zones phải đọc lại và khớp lại bảng chiều mỗi lần chạy. Với một bảng tham chiếu
nhỏ, ít thay đổi như bảng này — 265 dòng, được nạp lại toàn bộ ở mỗi lượt dbt run — thì đó
là một cuộc đánh đổi một chiều. Lượt benchmark của module 05 truy vấn taxi_zones_dict trực
tiếp bằng dictGet, nên bước này là một phụ thuộc cứng của module đó, không phải một phần
thêm tùy chọn.
Hãy source thông tin kết nối, rồi nạp DDL của dictionary bằng clickhouse-client hoặc HTTP
API — chọn cái nào bạn có sẵn:
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .clickhouse_state
# Via clickhouse-client
clickhouse-client \
--host "${CLICKHOUSE_HOST}" \
--port 9440 \
--user default \
--password "${CLICKHOUSE_PASSWORD}" \
--secure \
--multiquery \
< scripts/04_create_dictionary.sql
# Or via HTTP API
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/" \
--user "default:${CLICKHOUSE_PASSWORD}" \
--data-binary @scripts/04_create_dictionary.sqlKiểm chứng:
-- Should return 'Manhattan' for zone 42
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));
-- Should show status = LOADED, element_count = 265
SELECT name, status, element_count
FROM system.dictionaries
WHERE name = 'taxi_zones_dict';Cách kiểm tra bạn đã xong
SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 millionSELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0Số không ở đây là đúng, không phải một lỗi. Filter incremental của agg_hourly_zone_trips
(WHERE pickup_at >= now() - INTERVAL 2 HOUR) chỉ khớp các dòng do producer chạy trực tiếp
ghi vào, và mọi dòng đang có trong ClickHouse lúc này đều là dữ liệu lịch sử do script di
chuyển của module 03 chuyển sang — không dòng nào mới hơn 2 giờ so với now(). Bảng này chỉ
được nạp dữ liệu sau khi bước cutover của module 05 khởi động producer ClickHouse; đến lúc đó,
bất kỳ chart dashboard nào dựa trên bảng này sẽ không hiện dữ liệu, và điều đó là dự kiến,
không phải thứ cần debug.
SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt testMong đợi: mọi test đều pass.
-- Should return a borough name, e.g. 'Manhattan'
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));Trạng thái kết thúc
Tầng analytics đã có dữ liệu và đã được test: fact_trips chứa khoảng 50 triệu dòng,
dim_taxi_zones, dim_payment_type, dim_vendor và dim_date đã được nạp đầy đủ,
dbt test pass từ đầu đến cuối, và analytics.taxi_zones_dict đang hoạt động và trả về
borough qua dictGet(). agg_hourly_zone_trips vẫn rỗng — theo thiết kế, không phải do lỗi —
và giữ nguyên như vậy cho đến bước cutover của module 05.
Dashboard, benchmark ClickHouse so với Snowflake, và cutover thuộc module 05, không phải module này.
Producer Snowflake vẫn đang chạy, và khoảng trống giữa Snowflake và ClickHouse vẫn còn mở. Không có gì trong module này chạm vào producer hay script di chuyển — module 05 đóng khoảng trống đó một cách có chủ ý, trong đúng bước hai lượt được kiểm soát mà module 03 đã giới thiệu trước. Đừng dừng producer lúc này.
03 Cấp phát và di chuyển
Cấp phát ClickHouse Cloud bằng Terraform, tạo các bảng đích từ kế hoạch của bạn, và chuyển 50 triệu dòng bằng một script di chuyển Python có thể tiếp tục sau khi gián đoạn.
05 Benchmark và cutover
Dựng lại các dashboard trên ClickHouse, benchmark cả bảy truy vấn trên cả hai engine, cutover producer, kiểm chứng tính tương đương, và dỡ bỏ.