Snowflake MigrationClickHouse Workshops

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 đó.


Danh Sách Kiểm Tra Hoàn Thành

  • Chọn engine: đã hoàn thành
  • Thiết kế sort key: đã hoàn thành
  • Dịch schema: đã hoàn thành
  • Kế hoạch các đợt di trú: đã hoàn thành
  • Thiết kế dbt model: đã hoàn thành

Phần 1: Tóm Tắt Hồ Sơ

Chỉ sốGiá trị
Tổng số bảng7 (TRIPS_RAW, FACT_TRIPS, AGG_HOURLY_ZONE_TRIPS, DIM_TAXI_ZONES, DIM_PAYMENT_TYPE, DIM_VENDOR, DIM_DATE)
Tổng số view2 (STG_TRIPS, STG_TAXI_ZONES)
Stream1 (TRIPS_CDC_STREAM trên TRIPS_RAW)
Task2 (CDC_CONSUME_TASK, HOURLY_AGG_TASK)
Tổng số dòng trong TRIPS_RAW~50,000,000
Khoảng thời gianCửa sổ trượt 4 năm kết thúc tại thời điểm setup
Cột VARIANT1 (TRIPS_RAW.TRIP_METADATA)
Số lần dùng QUALIFY được phát hiện1 (truy vấn Q3)
Số lần dùng MERGE INTO được phát hiện2 (HOURLY_AGG_TASK, CDC_CONSUME_TASK)

Phần 2: Kiểm Kê Đối Tượng

Đối tượngLoạiSchemaSố dòngBậc phức tạpGhi chú
trips_rawTableraw~50MBRMT 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_tripsdbt Viewstaging—BJSONExtract cho TRIP_METADATA; cần kiểm thử các JSON path
stg_taxi_zonesdbt Viewstaging—AChuyển tiếp trực tiếp; đơn giản
int_trips_enricheddbt Ephemeralstaging—ACTE; khác biệt SQL được xử lý ở các model cha
fact_tripsdbt Incrementalanalytics~50MCEngine RMT; delete_insert; viết lại QUALIFY; cần FINAL
agg_hourly_zone_tripsdbt Incrementalanalytics~140KBRMT; 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_zonesdbt Tableanalytics265ADữ liệu tham chiếu tĩnh; nạp lại toàn bộ; đơn giản
dim_payment_typedbt Tableanalytics6ADữ liệu tham chiếu tĩnh; đơn giản
dim_vendordbt Tableanalytics3ADữ liệu tham chiếu tĩnh; đơn giản
taxi_zones_dictDictionaryanalytics265BCú pháp đặc thù ClickHouse; dictGet() tại thời điểm truy vấn
mv_hourly_revenueRefreshable MVanalytics—BCú pháp REFRESH EVERY; xác minh việc thay thế nguyên tử
TRIPS_CDC_STREAM / CDC_CONSUME_TASKSnowflake Stream + Task——DKhông có tương đương trong ClickHouse; được thay bằng cutover producer trực tiếp ở Phần 3

Phần 3: Các Quyết Định Chọn Engine

BảngEngineCột versionLập luận
trips_rawReplacingMergeTree(_synced_at)_synced_atScript 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_tripsReplacingMergeTree(updated_at)updated_atChuyế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_tripsReplacingMergeTree(updated_at)updated_atdbt 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_zonesMergeTree()—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_typeMergeTree()—Tương tự — nạp lại toàn bộ.
dim_vendorMergeTree()—Tương tự — nạp lại toàn bộ.
mv_hourly_revenueMergeTree()—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.

Phần 4: Thiết Kế Sort Key

BảngORDER BYLậ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.

Phần 5: Ghi Chú Dịch Schema

CộtKiểu SnowflakeKiểu ClickHouseLý do quyết định
TRIP_METADATAVARIANTStringGiữ 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_ATTIMESTAMP_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_IDINTEGERUInt16Giá 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_IDINTEGERUInt8Giá trị 1–3. UInt8 tối đa là 255 — đúng. 1 byte mỗi dòng.
DRIVER_RATINGFLOATNullable(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_ATTIMESTAMP_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.

Các Hàm Cần Dịch

Biểu thức SnowflakeTươ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::FLOATJSONExtractFloat(trip_metadata, 'driver', 'rating')
TRIP_METADATA:surge_multiplier::FLOATJSONExtractFloat(trip_metadata, 'surge_multiplier')
QUALIFY ROW_NUMBER() OVER (PARTITION BY pickup_location_id ORDER BY fare_amount DESC) <= 10SELECT ... FROM (SELECT ..., ROW_NUMBER() OVER (...) AS rn FROM ...) WHERE rn <= 10
MERGE INTO fact_trips ... WHEN MATCHED THEN UPDATEdbt 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

Phần 6: Các Đợt Di Trú

ĐợtĐối tượngPhụ thuộcGhi chú
Đợt 0trips_raw (schema), dim_taxi_zones, dim_payment_type, dim_vendorKhô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 1Nạ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 2stg_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 3taxi_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 4Cutover 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.

Sổ Đăng Ký Rủi Ro (các đối tượng bậc C/D)

Đối tượngRủi roCách kiểm chứng
fact_tripsTruy 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_tripsCử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 producerScript 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.

Phần 7: Các Khoảng Trống Dialect Đã Biết

  • QUALIFY — ảnh hưởng tới: Q3 (queries/q03_top_trips_qualify.sql)
  • Colon-path của VARIANT — ảnh hưởng tới: Q4, Q5 (truy cập JSON TRIP_METADATA)
  • LATERAL FLATTEN — không dùng trong workload này; VARIANT được truy cập qua colon-path, không qua FLATTEN
  • MERGE INTO — ảnh hưởng tới: các dbt incremental model (fact_trips, agg_hourly_zone_trips)
  • Snowflake Streams → cutover producer (sau cutover, các lần ghi trực tiếp đi thẳng vào ClickHouse)
  • Khác biệt hàm ngày tháng — ảnh hưởng tới: Q1 (DATE_TRUNC), Q3 (DATEADD), Q4 (DATEDIFF)

Phần 8: Chiến Lược Di Trú

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.

Phần 9: Tiêu Chí Cutover

Tiêu chíNgưỡngĐược đo bởi
Tương đương số dòngkhớp ≥ 99.9% (CH ≥ SF sau cutover là điều được kỳ vọng)scripts/01_verify_migration.sh
Tương đương checksumMD5 khớp trên mẫu 10K dòngscripts/02_validate_parity.sql
Tỉ lệ dbt test đạt100%dbt test trong dbt/nyc_taxi_dbt_ch
Tương đương kết quả truy vấnCả 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


Phần 10: Thiết Kế dbt Model

Chọn Materialization

ModelMaterializationTại sao
stg_tripsviewĐọ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_zonesviewTương tự — làm sạch dạng chuyển tiếp trực tiếp một bảng nguồn
int_trips_enrichedephemeralLogic 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_tripsincrementalChuyế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_tripsincrementalTí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_zonestable265 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_typetable6 loại tĩnh; cùng lập luận như dim_taxi_zones
dim_vendortable3 vendor; cùng lập luận

Cấu Hình Engine

ModelENGINECột versionTại sao
fact_tripsReplacingMergeTree(updated_at)updated_atChuyế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_tripsReplacingMergeTree(updated_at)updated_atViệ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_zonesMergeTree()—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_typeMergeTree()—Giống dim_taxi_zones
dim_vendorMergeTree()—Giống dim_taxi_zones

Chiến Lược Incremental

Modelunique_keyincremental_strategyBộ lọc incrementalTại sao chọn bộ lọc này?
fact_tripstrip_iddelete_insertWHERE 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_insertWHERE pickup_at >= now() - INTERVAL 2 HOURCử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

Đặt FINAL

ModelCó FINAL trong mệnh đề FROM?Tại sao
stg_tripsCó — FROM trips_raw FINALtrips_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_enrichedKhô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_tripsKhô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.

Trên trang này

VI