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.
Điểm khởi đầu
Module 02 đã hoàn tất: migration-plan.md đã được điền với mọi ô trong Completion Checklist
đã tick, và producer Snowflake vẫn đang chạy. Script setup.sh của module này kiểm tra file
đó và cảnh báo nếu nó thiếu hoặc chưa hoàn chỉnh, nhưng nó không bao giờ chặn — không có gì ở
đây ngăn bạn tiếp tục mà thiếu nó, chỉ có sự hiểu biết của chính bạn về hai module tiếp theo
là bị ngăn. Hãy dự trù khoảng 60 phút tổng cộng, trong đó khoảng 40-50 phút là một lượt truyền
dữ liệu không cần trông coi mà bạn có thể để chạy ở chế độ nền. Đây cũng là nơi bắt đầu tiêu
tốn credit dùng thử của ClickHouse Cloud: cấp phát service và làm hết module này tốn khoảng
$1-2 credit dùng thử (toàn bộ lab tốn khoảng $2-4 tổng cộng).
Vì sao
Module này là nơi kế hoạch trở thành hiện thực. Mọi quyết định bạn đã viết vào
migration-plan.md ở module 02 — engine MergeTree nào cho từng bảng, khóa ORDER BY suy ra
từ workload truy vấn thực tế, cách dịch các cấu trúc chỉ có ở Snowflake — được gõ trực tiếp
vào DDL bảng ở đây, chứ không suy lại từ đầu. ClickHouse không có index nào bạn có thể gắn
thêm về sau: nếu một khóa ORDER BY hóa ra là sai khi đã có 50 triệu dòng nằm trong bảng,
cách sửa là nạp lại toàn bộ, không phải một lệnh ALTER nhanh gọn.
Đó cũng là lý do cái cửa mềm kia vẫn quan trọng dù nó không thể chặn bạn. Nếu bạn chạy module
này mà không có một kế hoạch hoàn chỉnh, bạn vẫn sẽ thành công về mặt máy móc — dbt run vẫn
sẽ tạo fact_trips như một ReplacingMergeTree, script di chuyển vẫn sẽ chuyển 50 triệu dòng
— nhưng bạn sẽ không biết vì sao lại là engine đó chứ không phải MergeTree thuần, vì sao
sort key có hình dạng như vậy, hay làm sao biện giải cho mức tăng tốc benchmark ~6-9 lần mà
module 04 sẽ cho bạn thấy sau này. Bảng đối chiếu quyết định bên dưới ánh xạ mọi lựa chọn mà
module này hiện thực về đúng câu hỏi worksheet mà nó trả lời, để bạn có thể so kế hoạch của
mình với nó trước khi cấp phát bất cứ thứ gì.
Khái niệm — bên dưới lớp vỏ
Kiến trúc đích. Snowflake tiếp tục ghi các chuyến đi mới thông qua producer chuyến đi trong khi một script Python chạy một lần backfill 50 triệu dòng hiện có sang ClickHouse — hai hệ thống chạy song song suốt thời gian di chuyển, chứ không phải một lần cutover.
Ở phía ClickHouse, trips_raw là bảng đích mà script ghi vào. Sau đó dbt dựng các view
staging và phần còn lại của tầng analytics bên trên nó — đây là schema mà module này tạo ra,
nhưng chưa nạp dữ liệu ngoài trips_raw:
Chú giải màu của sơ đồ:
- Xanh lá — nguồn dữ liệu (producer chuyến đi, trước và sau cutover)
- Xanh dương — các bảng Snowflake
- Cam — các model dbt và pipeline
- Đỏ — các bảng ClickHouse và materialized view
- Xanh lơ — các dashboard Apache Superset
- Mũi tên nét đứt — các luồng sau cutover
Vì sao dùng script Python chứ không phải một connector native. Có nhiều phương pháp để chuyển dữ liệu từ Snowflake sang ClickHouse. Lab này dùng một script Python theo lô — đây là lý do, so với các phương án khác:
| Phương pháp | Cách hoạt động | Vì sao không dùng ở đây |
|---|---|---|
| ClickPipes (nguồn Snowflake) | Connector native của ClickHouse Cloud — zero-ETL, UI được quản lý | Snowflake không phải là nguồn ClickPipes được hỗ trợ. ClickPipes hỗ trợ Kafka, S3, Kinesis, PostgreSQL CDC, MySQL CDC và object storage. |
| Xuất ra S3 → ClickPipes S3 | COPY INTO @stage xuất Parquet/CSV ra S3; connector ClickPipes S3 nạp vào ClickHouse | Cần một S3 bucket, một IAM role, một Snowflake stage và một account AWS. Thêm ~3 bước chuẩn bị trước khi dữ liệu bắt đầu chuyển. Khả thi ở production nhưng quá nhiều hạ tầng cho một lab. |
Xuất ra S3 → clickhouse-client | Cũng xuất ra S3, nhưng nạp bằng INSERT INTO ... SELECT FROM s3(...) | Vẫn cần những điều kiện tiên quyết về S3. Ngoài ra còn buộc đối tác tự quản lý việc chia file thành từng khối và khả năng tiếp tục sau gián đoạn. |
| Snowflake → Kafka → ClickHouse | Stream CDC của Snowflake nạp vào một topic Kafka; connector ClickPipes Kafka thu nạp nó | Một pipeline streaming đầy đủ — phù hợp với yêu cầu độ trễ dưới một phút ở production. Một cluster Kafka thì quá nặng cho môi trường lab. |
| Script Python (lab này) | snowflake-connector-python đọc theo từng lô cursor 100K dòng; clickhouse-connect chèn trực tiếp | Không cần hạ tầng bổ sung nào ngoài các package mà lab đã cần. Có thể tiếp tục sau gián đoạn qua --resume (mốc nước max(pickup_at)). Có output tiến độ theo thời gian thực. ~40-50 phút cho 50M dòng ở tốc độ ~20K dòng/giây — chấp nhận được cho một bài tập di chuyển một lần. |
Vì sao script Python là lựa chọn đúng cho lab này:
- Không cần account AWS. Các cách tiếp cận dựa trên S3 đòi hỏi tạo bucket, các IAM policy và một external stage của Snowflake — ba bước chuẩn bị chẳng liên quan gì đến ClickHouse.
- Tự chứa. Hai package (
snowflake-connector-python,clickhouse-connect) được cài vào cùng venv với dbt. Không service mới, không credential mới. - Có thể tiếp tục sau gián đoạn.
--resumelàm cho script an toàn khi bị ngắt và chạy lại.ReplacingMergeTree(_synced_at)đảm bảo các lượt chèn trùng khi thử lại được tự động loại trùng. - Trong suốt. Đối tác có thể đọc script, hiểu cách ánh xạ cột, và điều chỉnh nó cho schema của riêng mình — điều này mang tính giáo dục hơn là bấm qua một wizard trên UI.
Xử lý khoảng trống khi di chuyển. Producer Snowflake tiếp tục chạy trong lúc script di chuyển chạy (~40-50 phút). Mọi chuyến đi được ghi vào Snowflake trong cửa sổ đó đều không có trong ClickHouse. Lab này đóng khoảng trống đó bằng một cách tiếp cận hai lượt vào thời điểm cutover, và module 05 sẽ dẫn bạn qua nó trực tiếp:
- Dừng producer Snowflake để đóng băng tập dữ liệu.
- Chạy
python scripts/02_migrate_trips.py --resume— chỉ các dòng delta được chuyển (tính bằng giây, không phải phút). - Khởi động producer ClickHouse.
Chính cơ chế loại trùng ReplacingMergeTree(_synced_at) xử lý các lượt thử lại khi di chuyển
cũng xử lý việc này: nếu có dòng nào trùng lặp giữa lượt chạy của module này và lượt --resume
sau đó, dòng có _synced_at muộn hơn sẽ thắng.
Khi nào bạn sẽ chọn S3 ở production. Nếu tập dữ liệu lớn hơn 500M dòng, hoặc nếu chi phí truy vấn trên warehouse Snowflake cho một lượt quét toàn bảng là đáng kể, thì đường xuất ra S3 sẽ tốt hơn: Snowflake xuất Parquet nén song song (nhanh hơn nhiều so với một cursor đơn), và ClickHouse cũng có thể nạp từ S3 song song. Cách dùng script Python ở đây hoạt động tốt ở quy mô lab.
Đối chiếu quyết định. Bảng dưới đây chính là danh sách quyết định từ Worksheet 1 (chọn
engine), 2 (sort key) và 3 (dịch schema) trong migration-plan.md, được kiểm chứng chéo với
những gì lab này thực sự dựng — hãy so nó với kế hoạch của riêng bạn trước khi cấp phát bất cứ
thứ gì:
| Quyết định | Lab này hiện thực | Vì sao |
|---|---|---|
Engine của trips_raw | ReplacingMergeTree(_synced_at) | Script di chuyển Python dùng các lệnh INSERT theo lô có thể bị thử lại nếu bị ngắt. _synced_at DateTime DEFAULT now() được đặt ở mọi lệnh INSERT, nên một dòng được thử lại sẽ đến muộn hơn và có giá trị _synced_at cao hơn — dòng muộn hơn thắng trong lúc RMT loại trùng, khiến các lượt thử lại trở nên idempotent. Các lượt thử lại của producer sau cutover cũng an toàn vì cùng lý do đó. stg_trips truy vấn với FINAL để bảo đảm mỗi chuyến đi chỉ còn một dòng. |
Engine của fact_trips | ReplacingMergeTree(updated_at) | Chuyến đi có thể được sửa (điều chỉnh giá cước); updated_at là cột phiên bản |
Engine của agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | Tính lại theo kiểu cuốn = mẫu upsert |
Engine của các bảng dim_* | MergeTree() | Nạp lại toàn bộ ở mỗi lượt dbt run; không có upsert |
ORDER BY của fact_trips | (toStartOfMonth(pickup_at), pickup_at, trip_id) | Q1-Q7 đều filter trên pickup_at; trip_id bảo đảm tính duy nhất ở mức lá |
ORDER BY của agg_hourly_zone_trips | (hour_bucket, zone_id) | Cả hai cột đều xuất hiện trong mọi truy vấn tổng hợp |
| VARIANT → | String + JSONExtract* | Giữ nguyên JSON thô; việc trích xuất diễn ra lúc truy vấn |
| QUALIFY → | Subquery bọc quanh ROW_NUMBER() | ClickHouse đã có mệnh đề QUALIFY native từ v24.5, nhưng dạng subquery được dạy vì nó mang đi được sang các phiên bản ClickHouse và các engine SQL cũ hơn hoặc không có QUALIFY |
| MERGE INTO → | Model incremental delete_insert trong dbt | Chiến lược upsert đúng phong cách của dbt-clickhouse; tránh việc ghi lại toàn bảng |
Bước 1 — Cấp phát cluster ClickHouse
setup.sh làm một việc: chạy terraform apply và ghi thông tin kết nối vào
.clickhouse_state. Nó cũng kiểm tra lại migration-plan.md trước khi cấp phát bất cứ thứ
gì — xem phần Vì sao ở trên — nhưng phép kiểm tra đó chỉ cảnh báo, nó không bao giờ chặn.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
# Configure credentials
cp .env.example .env
vim .env
# Fill in: CLICKHOUSE_ORG_ID, CLICKHOUSE_TOKEN_KEY, CLICKHOUSE_TOKEN_SECRET, CLICKHOUSE_PASSWORD
# Provision
source .env && ./setup.sh.env đã được gitignore — đừng bao giờ commit nó.
Output mong đợi: Terraform tạo 2 tài nguyên (service + danh sách IP được phép truy cập) trong khoảng 2-3 phút:
Apply complete! Resources: 2 added, 0 changed, 0 destroyed.
Outputs:
clickhouse_host = "abc123xyz.us-east-1.aws.clickhouse.cloud"
clickhouse_port = 8443Host và port được lưu vào .clickhouse_state. Hãy source nó ở bất kỳ terminal nào để lấy
thông tin kết nối:
source .clickhouse_stateKiểm chứng:
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1Bước 2 — Tạo các bảng rỗng
Trước tiên, hãy tự tay tạo trips_raw với engine đúng. Script di chuyển nạp dữ liệu vào bảng
này ở Bước 3 — nó phải đã tồn tại với ReplacingMergeTree để cột phiên bản có sẵn trước khi
bất kỳ dòng nào đến.
-- Run in the ClickHouse SQL console (cloud.clickhouse.com -> SQL console)
CREATE TABLE IF NOT EXISTS default.trips_raw (
trip_id String,
vendor_id UInt8,
pickup_at DateTime64(3, 'UTC'),
dropoff_at DateTime64(3, 'UTC'),
passenger_count UInt8,
trip_distance_miles Float32,
pickup_location_id UInt16,
dropoff_location_id UInt16,
payment_type_id UInt8,
rate_code_id UInt8,
store_fwd_flag String,
fare_amount_usd Float32,
extra_amount_usd Float32,
mta_tax_usd Float32,
tip_amount_usd Float32,
tolls_amount_usd Float32,
total_amount_usd Float32,
ingested_at DateTime64(3, 'UTC'),
trip_metadata String,
_synced_at DateTime DEFAULT now()
)
ENGINE = ReplacingMergeTree(_synced_at)
ORDER BY (pickup_at, trip_id);_synced_at được đặt tự động ở mọi lệnh INSERT. Nếu script di chuyển bị ngắt và chạy lại với
--resume, các dòng trùng cho cùng một trip_id có thể tồn tại trong thời gian ngắn — RMT
giữ dòng muộn hơn (_synced_at cao hơn). stg_trips truy vấn trips_raw FINAL để buộc loại
trùng trước khi bất kỳ model phía sau nhìn thấy dữ liệu.
Tiếp theo, nạp dữ liệu tham chiếu về vùng. Đây là dữ liệu tĩnh (265 vùng NYC TLC) mà
stg_taxi_zones của dbt đọc như một source.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
clickhouse-client --host "${CLICKHOUSE_HOST}" --port 9440 --secure \
--user default --password "${CLICKHOUSE_PASSWORD}" \
--multiquery < scripts/00_seed_zones.sql(Hoặc thay vào đó dán trực tiếp nội dung của scripts/00_seed_zones.sql vào SQL console của
ClickHouse.)
Cấu hình profile dbt. File dbt_project.yml của dự án này khai báo
profile: 'nyc_taxi_ch'. Nếu không có profile khớp trong ~/.dbt/profiles.yml, dbt run sẽ
lỗi ngay với Could not find profile named 'nyc_taxi_ch'. Template nằm ở
workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch/profiles.yml.example.
Module 01 đã ghi ~/.dbt/profiles.yml với một profile nyc_taxi: cho Snowflake, và vòng lặp
refresh ở Bước 4 của nó vẫn tiếp tục truy vấn với profile đó suốt thời gian producer Snowflake
chạy. Đừng thay thế file đó bằng template ClickHouse — ghi đè nó bằng
profiles.yml.example sẽ xóa profile nyc_taxi: và làm hỏng vòng lặp refresh của module 01.
Thay vào đó, hãy mở template và trộn khối nyc_taxi_ch: của nó vào file
~/.dbt/profiles.yml hiện có như một profile cấp cao thứ hai, bên cạnh nyc_taxi::
nyc_taxi: # from module 01 — leave this one alone
target: dev
outputs:
dev:
type: snowflake
# ...
nyc_taxi_ch: # add this block
target: dev
outputs:
dev:
type: clickhouse
schema: nyc_taxi_ch
host: "{{ env_var('CLICKHOUSE_HOST') }}"
port: 8443
user: "{{ env_var('CLICKHOUSE_USER', 'default') }}"
password: "{{ env_var('CLICKHOUSE_PASSWORD') }}"
secure: truenyc_taxi_ch: đọc CLICKHOUSE_HOST, CLICKHOUSE_USER và CLICKHOUSE_PASSWORD từ môi trường
qua env_var(), nên .env và .clickhouse_state phải được source trước bất kỳ lệnh dbt nào
trong module này — lệnh dbt run bên dưới đã làm việc đó. Giống như profile Snowflake,
~/.dbt/profiles.yml chứa credential và đã được gitignore; đừng bao giờ commit nó, và lần
trộn này thêm một bộ credential thứ hai vào một file đã có một bộ.
Kiểm chứng:
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.clickhouse_state"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.env"
dbt debug
# Expected: "All checks passed!" — confirms dbt found the nyc_taxi_ch profile and
# connected to ClickHouseSau đó chạy dbt run để tạo các bảng analytics và các view staging:
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.clickhouse_state"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.env"
dbt deps # install packages (first run only)
dbt run # creates analytics tables and staging views; all empty at this pointMong đợi: ~8 model được tạo trong dưới 2 phút (mọi bảng đều rỗng).
Kiểm chứng:
# Check analytics tables were created
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SHOW+TABLES+IN+analytics" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: agg_hourly_zone_trips, dim_date, dim_payment_type, dim_vendor, dim_taxi_zones, fact_trips
# Check staging views were created
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SHOW+TABLES+IN+staging" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: stg_trips, stg_taxi_zones
# Check trips_raw exists with the correct engine
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+engine+FROM+system.tables+WHERE+database%3D%27default%27+AND+name%3D%27trips_raw%27" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: ReplacingMergeTreeCả sáu bảng analytics và cả hai view staging giờ đã tồn tại, nhưng mọi thứ trong số đó vẫn
còn rỗng — dbt run chỉ tạo schema của chúng. Bảng duy nhất có dữ liệu sau bước này là
trips_raw, và nó cũng chưa có dòng nào; việc đó đến ngay sau đây.
Bước 3 — Di chuyển dữ liệu
Nạp toàn bộ dòng từ Snowflake NYC_TAXI_DB.RAW.TRIPS_RAW vào ClickHouse
default.trips_raw bằng script di chuyển theo lô viết bằng Python.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
source .venv/bin/activate
python scripts/02_migrate_trips.pyOutput mong đợi (khoảng 40-50 phút cho 50M dòng):
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
NYC Taxi Migration: Snowflake -> ClickHouse
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Rows to migrate: 50,000,000
Batch size: 100,000
Rows inserted Elapsed ETA Rate
-------------------- ------------ ---------------------- ---------------
100,000 0m 07s 56m 14s remaining 13,945 rows/s
200,000 0m 14s 55m 28s remaining 14,021 rows/s
...Producer Snowflake tiếp tục ghi vào TRIPS_RAW suốt thời gian script này chạy, nên ClickHouse
tụt lại phía sau khoảng bằng độ dài của lượt truyền này — khoảng trống đó là điều được dự
kiến, và nó được xử lý ở module 05, không phải ở đây.
Nếu script bị ngắt, hãy chạy lại nó với --resume để tiếp tục từ checkpoint cuối:
python scripts/02_migrate_trips.py --resume--resume đọc max(pickup_at) từ ClickHouse và bỏ qua các dòng đã nạp, nên luôn an toàn khi
ngắt script này và chạy lại — bạn sẽ không bao giờ kết thúc với một lượt nạp dở dang không thể
khôi phục.
Cách kiểm tra bạn đã xong
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
# Row count in trips_raw
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+count()+FROM+default.trips_raw" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: approximately 50000000
# .clickhouse_state was written by setup.sh
ls -la .clickhouse_state
# Service is reachable
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1Trạng thái kết thúc
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 — tổng cộng bảy đối tượng trong schema analytics, được dựng
bởi lệnh dbt run của module nà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 mà dbt run sinh ra khi dựng view. Ngoài ra chỉ có
default.trips_raw là có dữ liệu: khoảng 50 triệu dòng, do script di chuyển Python chuyển
sang ở Bước 3.
Producer Snowflake vẫn đang chạy. Nó chưa bao giờ bị dừng trong module này và cũng không
dừng ở đây. Mọi chuyến đi được ghi vào TRIPS_RAW của Snowflake sau lô cuối cùng của script
di chuyển đều là một dòng mà ClickHouse không có, nên ClickHouse giờ tụt lại sau Snowflake
khoảng bằng độ dài của cửa sổ di chuyển (~40-50 phút, cộng thêm thời gian chuẩn bị của module
này). Khoảng trống đó là thật và tiếp tục lớn lên chừng nào producer còn chạy. Đừng đóng nó
trong module này. Bước cutover của module 05 đóng nó một cách có chủ ý, trong một bước hai
lượt được kiểm soát nhằm đo độ lớn của khoảng trống trước khi triệt tiêu nó — dừng producer
hoặc chạy lại script di chuyển lúc này sẽ xóa đi chính thứ mà module 05 được dựng ra để minh
họa.
Worksheet 5: Thiết kế model dbt
Cấu hình materialization, engine và chiến lược incremental cho từng model dbt, với phản hồi ngay lập tức cho mọi câu trả lời.
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.