Ví dụ mẫu: một bản kế hoạch đã hoàn thành
Một bản kế hoạch di trú đã điền đầy đủ cho workload NYC taxi, để bạn đối chiếu với bản của mình sau khi đã viết xong.
Đây là đáp án đầy đủ cho cả năm worksheet
(1,
2,
3,
4,
5) áp dụng cho workload NYC
Taxi. Hãy dùng nó để:
- Kiểm tra đáp án worksheet của bạn sau khi hoàn thành từng phần
- Hiểu lập luận phía sau các quyết định mà Phần 3 triển khai
- Đối chiếu với bảng Decision Alignment của Phần 3 nếu bạn chọn khác
Đây là đáp án — đừng điền vào đây như kế hoạch của bạn. Hãy điền vào migration-plan.md thay vào đó.
| Chỉ số | Giá trị |
|---|
| Tổng số bảng | 7 (TRIPS_RAW, FACT_TRIPS, AGG_HOURLY_ZONE_TRIPS, DIM_TAXI_ZONES, DIM_PAYMENT_TYPE, DIM_VENDOR, DIM_DATE) |
| Tổng số view | 2 (STG_TRIPS, STG_TAXI_ZONES) |
| Stream | 1 (TRIPS_CDC_STREAM trên TRIPS_RAW) |
| Task | 2 (CDC_CONSUME_TASK, HOURLY_AGG_TASK) |
| Tổng số dòng trong TRIPS_RAW | ~50,000,000 |
| Khoảng thời gian | Cửa sổ trượt 4 năm kết thúc tại thời điểm setup |
| Cột VARIANT | 1 (TRIPS_RAW.TRIP_METADATA) |
| Số lần dùng QUALIFY được phát hiện | 1 (truy vấn Q3) |
| Số lần dùng MERGE INTO được phát hiện | 2 (HOURLY_AGG_TASK, CDC_CONSUME_TASK) |
| Đối tượng | Loại | Schema | Số dòng | Bậc phức tạp | Ghi chú |
|---|
trips_raw | Table | raw | ~50M | B | RMT với cột version _synced_at; bulk + CDC chồng lấn nên cần dedup; stg_trips phải dùng FINAL |
stg_trips | dbt View | staging | — | B | JSONExtract cho TRIP_METADATA; cần kiểm thử các JSON path |
stg_taxi_zones | dbt View | staging | — | A | Chuyển tiếp trực tiếp; đơn giản |
int_trips_enriched | dbt Ephemeral | staging | — | A | CTE; khác biệt SQL được xử lý ở các model cha |
fact_trips | dbt Incremental | analytics | ~50M | C | Engine RMT; delete_insert; viết lại QUALIFY; cần FINAL |
agg_hourly_zone_trips | dbt Incremental | analytics | ~140K | B | RMT; cửa sổ tính lại trượt 2 giờ; hãy kiểm thử ranh giới partition thật cẩn thận |
dim_taxi_zones | dbt Table | analytics | 265 | A | Dữ liệu tham chiếu tĩnh; nạp lại toàn bộ; đơn giản |
dim_payment_type | dbt Table | analytics | 6 | A | Dữ liệu tham chiếu tĩnh; đơn giản |
dim_vendor | dbt Table | analytics | 3 | A | Dữ liệu tham chiếu tĩnh; đơn giản |
taxi_zones_dict | Dictionary | analytics | 265 | B | Cú pháp đặc thù ClickHouse; dictGet() tại thời điểm truy vấn |
mv_hourly_revenue | Refreshable MV | analytics | — | B | Cú pháp REFRESH EVERY; xác minh việc thay thế nguyên tử |
TRIPS_CDC_STREAM / CDC_CONSUME_TASK | Snowflake Stream + Task | — | — | D | Không có tương đương trong ClickHouse; được thay bằng cutover producer trực tiếp ở Phần 3 |
| Bảng | Engine | Cột version | Lập luận |
|---|
trips_raw | ReplacingMergeTree(_synced_at) | _synced_at | Script di trú Python (scripts/02_migrate_trips.py) có thể thử lại một lô và insert lại cùng một trip_id. Sau cutover, producer chạy trực tiếp cũng có thể thử lại khi có lỗi tạm thời. _synced_at DateTime DEFAULT now() được đặt tự động khi INSERT — một lần thử lại sau đó có timestamp lớn hơn, nên RMT giữ lần ghi mới nhất. stg_trips truy vấn với FINAL để cưỡng chế dedup trước khi bất kỳ model hạ nguồn nào chạy. |
fact_trips | ReplacingMergeTree(updated_at) | updated_at | Chuyến đi có thể được chỉnh sửa (điều chỉnh giá vé, thay đổi trạng thái). Cùng một trip_id được insert lại với giá trị đã cập nhật. updated_at tăng đơn điệu ở mỗi lần chỉnh sửa — giá trị lớn hơn thắng khi RMT dedup. Luôn truy vấn với FINAL. |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | dbt tính lại 2 giờ gần nhất và insert lại. Không có RMT, các bản tổng hợp cũ và mới sẽ tích tụ và bị đếm hai lần. updated_at được đặt thành now() ở mỗi lần dbt run bảo đảm giá trị mới nhất thắng. |
dim_taxi_zones | MergeTree() | — | dbt nạp lại toàn bộ (hoán đổi bảng nguyên tử (dựng lại toàn bộ)). Không có bản trùng nào tích tụ được. Không cần dedup. |
dim_payment_type | MergeTree() | — | Tương tự — nạp lại toàn bộ. |
dim_vendor | MergeTree() | — | Tương tự — nạp lại toàn bộ. |
mv_hourly_revenue | MergeTree() | — | REFRESHABLE MV thay thế nguyên tử toàn bộ tập kết quả của nó ở mỗi lần REFRESH. Không có upsert. |
| Bảng | ORDER BY | Lập luận |
|---|
trips_raw | (pickup_at, trip_id) | Các lần scan theo khoảng thời gian lọc trên pickup_at trước. trip_id là khóa dedup của RMT — nó phải nằm trong ORDER BY để RMT nhận biết được những dòng nào là bản trùng. pickup_at đứng trước vì các lần scan theo khoảng phục vụ phân tích chiếm ưu thế; trip_id đứng cuối vì nó có lực lượng cao và chỉ đóng vai trò phân biệt tính duy nhất. |
fact_trips | (toStartOfMonth(pickup_at), pickup_at, trip_id) | Cả 7 truy vấn phân tích đều lọc trên pickup_at. Tiền tố theo tháng gom dữ liệu cùng tháng dương lịch vào các block liền kề — cho phép bỏ qua block ở mức thô cho các phép tổng hợp theo tháng mà không cần thêm PARTITION BY. trip_id đứng cuối để bảo đảm tính duy nhất cho RMT mà không phá vỡ việc bỏ qua block. |
agg_hourly_zone_trips | (hour_bucket, zone_id) | Q6 (và mọi truy vấn tổng hợp) lọc trên hour_bucket và zone_id. hour_bucket có ~35K giá trị khác nhau; zone_id có 265. hour_bucket đứng trước vì scan theo khoảng thời gian là mẫu truy cập chính. zone_id đứng thứ hai để lọc phụ. |
dim_taxi_zones | (location_id) | 265 dòng = một granule. ORDER BY không liên quan tới hiệu năng. location_id làm khóa join là quy ước thông thường và giúp dễ đọc. |
| Cột | Kiểu Snowflake | Kiểu ClickHouse | Lý do quyết định |
|---|
TRIP_METADATA | VARIANT | String | Giữ nguyên chính xác JSON thô. JSONExtract* xử lý được đường dẫn bất kỳ tại thời điểm truy vấn. Map(String,String) làm mất các cấu trúc lồng nhau; Tuple đòi hỏi schema cố định. String là lựa chọn an toàn cho JSON bất kỳ. |
PICKUP_DATETIME / PICKUP_AT | TIMESTAMP_NTZ(9) | DateTime64(3, 'UTC') | Độ chính xác milli giây là đủ cho timestamp chuyến đi. Nano giây (9) là quá mức. 'UTC' làm cho múi giờ trở nên tường minh và tránh những bất ngờ liên quan tới DST trong các phép tổng hợp theo khoảng thời gian. |
PICKUP_LOCATION_ID | INTEGER | UInt16 | Giá trị 1–265. UInt8 tối đa là 255 (quá nhỏ). UInt16 tối đa là 65535 (đúng). 2 byte so với 4 byte của Int32 — tiết kiệm ~95MB chưa nén trên 50M dòng cho mỗi cột. |
VENDOR_ID | INTEGER | UInt8 | Giá trị 1–3. UInt8 tối đa là 255 — đúng. 1 byte mỗi dòng. |
DRIVER_RATING | FLOAT | Nullable(Float32) | Thường xuyên NULL (không phải chuyến đi nào cũng có điểm đánh giá). Nullable giữ đúng ngữ nghĩa null. Float32 là đủ cho khoảng 1.0–5.0. Float64 sẽ tốn thêm lưu trữ mà không thêm độ chính xác có ý nghĩa. |
UPDATED_AT | TIMESTAMP_NTZ(9) | DateTime64(3, 'UTC') | Cột version cho ReplacingMergeTree. Phải dùng DateTime64 chứ không phải DateTime — hai lần chỉnh sửa trong cùng một giây sẽ không xác định được thứ tự nếu chỉ có độ chính xác giây. Độ chính xác milli giây bảo đảm thứ tự dedup đúng. |
| Biểu thức Snowflake | Tương đương trong ClickHouse |
|---|
DATE_TRUNC('hour', pickup_at) | toStartOfHour(pickup_at) |
DATEADD('day', -7, CURRENT_DATE) | today() - 7 |
DATEDIFF('minute', pickup_at, dropoff_at) | dateDiff('minute', pickup_at, dropoff_at) |
TRIP_METADATA:driver.rating::FLOAT | JSONExtractFloat(trip_metadata, 'driver', 'rating') |
TRIP_METADATA:surge_multiplier::FLOAT | JSONExtractFloat(trip_metadata, 'surge_multiplier') |
QUALIFY ROW_NUMBER() OVER (PARTITION BY pickup_location_id ORDER BY fare_amount DESC) <= 10 | SELECT ... FROM (SELECT ..., ROW_NUMBER() OVER (...) AS rn FROM ...) WHERE rn <= 10 |
MERGE INTO fact_trips ... WHEN MATCHED THEN UPDATE | dbt delete_insert incremental — DELETE các dòng có khóa trùng khớp, rồi INSERT toàn bộ dòng mới |
| Đợt | Đối tượng | Phụ thuộc | Ghi chú |
|---|
| Đợt 0 | trips_raw (schema), dim_taxi_zones, dim_payment_type, dim_vendor | Không có | dbt tạo các bảng rỗng. Các bảng dim được nạp ngay từ dữ liệu tham chiếu tĩnh (không phụ thuộc vào trips). Chạy: dbt run --select trips_raw dim_* |
| Đợt 1 | Nạp khối lượng lớn bằng Python (scripts/02_migrate_trips.py) | Đợt 0 (schema trips_raw phải tồn tại) | 50M dòng từ Snowflake TRIPS_RAW. Có thể tiếp tục với --resume. Xác minh số dòng bằng scripts/01_verify_migration.sh. ~40-50 phút. |
| Đợt 2 | stg_trips, stg_taxi_zones, int_trips_enriched, fact_trips, agg_hourly_zone_trips | Đợt 1 hoàn tất (trips_raw đã có dữ liệu) + Đợt 0 (các bảng dim đã tồn tại) | dbt run đầy đủ. stg_trips đọc trips_raw; int_trips_enriched join với các dim; fact_trips và agg_hourly_zone_trips dựng lên trên đó. |
| Đợt 3 | taxi_zones_dict, mv_live_trip_feed | Đợt 2 (dim_taxi_zones đã có dữ liệu cho dict; fact_trips đã có dữ liệu cho MV) | Dictionary được tạo qua scripts/04_create_dictionary.sql. Refreshable MV được tạo qua dbt model; việc bật khoảng làm mới của nó sau đó là một bước ALTER TABLE ... MODIFY REFRESH thủ công, không phải việc dbt tự chạy. |
| Đợt 4 | Cutover producer (scripts/03_cutover.sh) | Đợt 1 hoàn tất (đã xác minh nạp khối lượng lớn) + Đợt 2 hoàn tất (tầng analytics đã dựng) | Dừng producer Snowflake; khởi động producer ClickHouse ghi trực tiếp vào ClickHouse Cloud; chạy dbt run để nạp agg_hourly_zone_trips với dữ liệu trực tiếp. |
| Đối tượng | Rủi ro | Cách kiểm chứng |
|---|
fact_trips | Truy vấn không có FINAL sẽ đếm vượt trong lúc merge còn trễ. Khoảng partition của delete_insert phải được giới hạn theo tiền tố ORDER BY để tránh xóa các partition không phải mục tiêu. | SELECT COUNT(*) FINAL khớp với Snowflake ± độ trễ CDC. Chạy dbt test. So sánh kết quả Q3 giữa hai hệ thống. Kiểm tra trip_id trùng: SELECT trip_id, count() FROM fact_trips GROUP BY trip_id HAVING count() > 1 LIMIT 10. |
agg_hourly_zone_trips | Cửa sổ tính lại trượt 2 giờ phải giới hạn đúng khoảng xóa. Nếu quá rộng, các bản tổng hợp cũ bị xóa; nếu quá hẹp, các bản tổng hợp lỗi thời vẫn tồn tại. | Kiểm tra ngẫu nhiên các tuple (hour_bucket, zone_id) cụ thể so với Snowflake. Xác minh tổng trip_count trên tất cả các zone khớp với AGG_HOURLY_ZONE_TRIPS của Snowflake trong cùng kỳ. |
| Cutover producer | Script di trú bị ngắt giữa lúc chạy sẽ để lại khoảng hụt về số dòng; hãy chạy lại với --resume để bù. Producer thử lại sau cutover có thể insert lại các chuyến đi đã có trong ClickHouse. | scripts/01_verify_migration.sh — kiểm tra sự tương đương số dòng giữa Snowflake và ClickHouse. ReplacingMergeTree(_synced_at) xử lý các lần insert trùng một cách idempotent. |
Di chuyển dữ liệu: Script di trú bằng Python (scripts/02_migrate_trips.py)
Tại sao dùng script Python thay vì trung chuyển qua object storage hay ClickPipes?
remoteSecure() dành cho truyền dữ liệu ClickHouse-sang-ClickHouse — không áp dụng được ở đây.
- Trung chuyển qua object storage (Snowflake → S3 → hàm table S3 của ClickHouse) thì được nhưng thêm phức tạp: cần cấp phát S3 bucket, IAM role, và Snowflake COPY INTO — overhead không cần thiết cho một lab.
- ClickPipes không hỗ trợ Snowflake làm nguồn. Các nguồn nó hỗ trợ là Kafka, S3, Kinesis, PostgreSQL CDC và MySQL CDC.
- Script Python dùng
snowflake-connector-python và clickhouse-connect — các package đã được cài cho lab. Nó hiển thị tiến trình theo thời gian thực, hỗ trợ --resume khi bị ngắt, và mã nguồn có thể xem xét đầy đủ.
Chiến lược incremental (dbt): delete_insert
Tại sao dùng delete_insert thay vì chiến lược append hay merge?
append insert các dòng mới mà không chạm tới các dòng đã có. Với fact_trips, nơi các dòng có thể bị cập nhật, cách này tạo ra bản trùng. Không đúng.
merge (nếu có) sẽ gần nhất với MERGE INTO của Snowflake, nhưng chiến lược merge của dbt-clickhouse có hạn chế khi làm việc với ReplacingMergeTree và không phải cách tiếp cận được khuyến nghị.
delete_insert xóa các dòng trong khoảng khóa của lô dữ liệu đến, rồi insert toàn bộ dòng mới. Cách này là idempotent (chạy lại cho ra cùng kết quả), xử lý được cả insert và update, và hoạt động đúng với ReplacingMergeTree. Đây là khuyến nghị chuẩn của cộng đồng dbt-clickhouse cho các mẫu upsert.
| Tiêu chí | Ngưỡng | Được đo bởi |
|---|
| Tương đương số dòng | khớp ≥ 99.9% (CH ≥ SF sau cutover là điều được kỳ vọng) | scripts/01_verify_migration.sh |
| Tương đương checksum | MD5 khớp trên mẫu 10K dòng | scripts/02_validate_parity.sql |
| Tỉ lệ dbt test đạt | 100% | dbt test trong dbt/nyc_taxi_dbt_ch |
| Tương đương kết quả truy vấn | Cả 7 truy vấn trả về cùng kết quả (trong dung sai số thực dấu phẩy động) | So sánh thủ công trong output của scripts/run_benchmark.sh |
| Model | Materialization | Tại sao |
|---|
stg_trips | view | Đọc và làm sạch trips_raw; không cập nhật gì cho model này; không tốn chi phí lưu trữ; luôn phản ánh trạng thái hiện tại của nguồn |
stg_taxi_zones | view | Tương tự — làm sạch dạng chuyển tiếp trực tiếp một bảng nguồn |
int_trips_enriched | ephemeral | Logic join thuần chỉ được fact_trips dùng; nhúng inline như một CTE giúp tránh một bảng vật lý dư thừa; không model nào truy vấn nó trực tiếp |
fact_trips | incremental | Chuyến đi có thể được chỉnh sửa sau đó; mỗi lần chạy chỉ nên xử lý các dòng mới và đã cập nhật |
agg_hourly_zone_trips | incremental | Tính lại trượt 2 giờ là một mẫu incremental — xử lý các dòng gần đây, không phải toàn bộ 50M |
dim_taxi_zones | table | 265 zone tĩnh; dựng lại toàn bộ ở mỗi lần dbt run qua hoán đổi bảng nguyên tử (dựng lại toàn bộ); không cập nhật cục bộ |
dim_payment_type | table | 6 loại tĩnh; cùng lập luận như dim_taxi_zones |
dim_vendor | table | 3 vendor; cùng lập luận |
| Model | ENGINE | Cột version | Tại sao |
|---|
fact_trips | ReplacingMergeTree(updated_at) | updated_at | Chuyến đi có thể được chỉnh sửa; updated_at được đặt thành now() ở mỗi lần insert, nghĩa là phiên bản mới nhất thắng trong quá trình dedup nền của RMT; delete_insert là đường bảo đảm tính đúng đắn chính, RMT là lưới an toàn |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | Việc tính lại theo cửa sổ trượt insert lại các bản tổng hợp cho cùng những cặp (hour_bucket, zone_id); RMT bảo đảm các bản tổng hợp lỗi thời bị loại bỏ khi merge nền |
dim_taxi_zones | MergeTree() | — | dbt nạp lại toàn bộ nghĩa là hoán đổi bảng nguyên tử (dựng lại toàn bộ) ở mỗi lần chạy; bản trùng không thể tích tụ; không cần dedup |
dim_payment_type | MergeTree() | — | Giống dim_taxi_zones |
dim_vendor | MergeTree() | — | Giống dim_taxi_zones |
| Model | unique_key | incremental_strategy | Bộ lọc incremental | Tại sao chọn bộ lọc này? |
|---|
fact_trips | trip_id | delete_insert | WHERE updated_at > (SELECT max(updated_at) FROM {{ this }}) | Mốc nước cao trên updated_at bắt được cả chuyến đi mới lẫn chuyến đi đã chỉnh sửa (điều chỉnh giá vé sẽ insert lại cùng trip_id với cùng pickup_at nhưng updated_at mới hơn); một mốc nước theo pickup_at sẽ âm thầm bỏ sót các lần chỉnh sửa |
agg_hourly_zone_trips | [hour_bucket, zone_id] | delete_insert | WHERE pickup_at >= now() - INTERVAL 2 HOUR | Cửa sổ trượt 2 giờ buộc tổng hợp lại các giờ ở ranh giới để số đếm của giờ chưa trọn vẹn luôn được sửa đúng; một mốc nước cao max(pickup_at) sẽ đếm thiếu vĩnh viễn ở giờ ranh giới |
| Model | Có FINAL trong mệnh đề FROM? | Tại sao |
|---|
stg_trips | Có — FROM trips_raw FINAL | trips_raw là ReplacingMergeTree; nó có thể có các dòng trip_id trùng do script di trú thử lại hoặc do producer thử lại sau cutover. stg_trips là điểm cưỡng chế duy nhất: hãy loại trùng ở đây để mọi model hạ nguồn (int_trips_enriched, fact_trips, agg_hourly_zone_trips) đều nhận được dữ liệu sạch |
int_trips_enriched | Không | Đọc từ stg_trips (một view), không phải một bảng RMT; FINAL không liên quan với view |
fact_trips | Không (trong thân model) | delete_insert giữ fact_trips sạch sau mỗi lần chạy hoàn tất; thêm FINAL bên trong model sẽ khiến nó bị áp một cách vô ích lên subquery is_incremental() đọc max(updated_at) từ {{ this }}. Dashboard và dbt test dùng FINAL từ bên ngoài khi truy vấn trực tiếp fact_trips |
Đây là ví dụ đã hoàn thành. migration-plan.md của bạn nên khớp với các quyết định then chốt ở đây — hoặc ghi lại tường minh lý do bạn chọn khác.