03 จัดเตรียมและย้ายข้อมูล
จัดเตรียม ClickHouse Cloud ด้วย Terraform, สร้างตารางเป้าหมายจากแผนของคุณ และย้าย 50 ล้านแถวด้วยสคริปต์ย้ายข้อมูล Python ที่ทำงานต่อจากจุดเดิมได้
จุดเริ่มต้น
โมดูล 02 เสร็จแล้ว: migration-plan.md ถูกกรอกครบพร้อมช่องติ๊กทุกช่องใน Completion
Checklist ติ๊กแล้ว และ producer ของ Snowflake ยังทำงานอยู่ setup.sh ของโมดูลนี้ตรวจหา
ไฟล์นั้นและเตือนถ้ามันหายไปหรือไม่สมบูรณ์ แต่มันไม่เคยขัดขวาง — ไม่มีอะไรที่นี่หยุดคุณจาก
การไปต่อโดยไม่มีมัน มีแต่ความเข้าใจของคุณเองในสองโมดูลถัดไปที่จะหยุดคุณ กันเวลาไว้ประมาณ
60 นาทีรวม โดยราว 40-50 นาทีเป็นการถ่ายโอนข้อมูลที่ไม่ต้องเฝ้าและคุณปล่อยให้รันเบื้องหลังได้
นี่ยังเป็นจุดที่การใช้จ่ายเครดิตทดลองใช้ของ ClickHouse Cloud เริ่มขึ้น: การจัดเตรียมเซอร์วิส
และทำโมดูลนี้ให้จบใช้เครดิตทดลองใช้ราว $1-2 (แล็บทั้งหมดใช้ราว $2-4 รวม)
ทำไม
โมดูลนี้คือจุดที่แผนกลายเป็นของจริง ทุกการตัดสินใจที่คุณเขียนลงใน migration-plan.md
ในโมดูล 02 — เอนจิน MergeTree ตัวใดต่อตาราง, คีย์ ORDER BY ที่อนุมานจากเวิร์กโหลดคิวรี
ที่เกิดขึ้นจริง, โครงสร้างที่มีแต่ใน Snowflake แปลออกมาเป็นอะไร — จะถูกพิมพ์ลงใน DDL
ของตารางที่นี่โดยตรง ไม่ใช่อนุมานขึ้นมาใหม่จากศูนย์ ClickHouse ไม่มีอินเด็กซ์ที่คุณจะแปะเพิ่ม
ภายหลังได้: ถ้าคีย์ ORDER BY กลายเป็นว่าผิดเมื่อมี 50 ล้านแถวนั่งอยู่ในตารางแล้ว วิธีแก้
คือโหลดใหม่ทั้งหมด ไม่ใช่ ALTER แบบเร็ว ๆ
นั่นก็เป็นเหตุผลว่าทำไมด่านแบบไม่บังคับจึงสำคัญ แม้ว่ามันจะหยุดคุณไม่ได้ ถ้าคุณรันโมดูลนี้
โดยไม่มีแผนที่เสร็จ คุณก็จะยังสำเร็จในเชิงกลไก — dbt run จะยังสร้าง fact_trips เป็น
ReplacingMergeTree, สคริปต์ย้ายข้อมูลจะยังย้าย 50 ล้านแถว — แต่คุณจะไม่รู้ว่าทำไมต้องเป็น
เอนจินนั้นและไม่ใช่ MergeTree ธรรมดา, ทำไม sort key มีรูปทรงอย่างที่มันเป็น หรือจะให้
เหตุผลรองรับตัวเลขความเร็วที่เพิ่มขึ้นราว 6-9 เท่าจากเบนช์มาร์กที่โมดูล 04 แสดงให้คุณดู
ภายหลังได้อย่างไร ตารางความสอดคล้องของการตัดสินใจด้านล่างแมปทุกทางเลือกที่โมดูลนี้
ทำจริงกลับไปยังคำถามในใบงานที่มันตอบ เพื่อให้คุณเทียบแผนของคุณเองกับมันได้ก่อนจัดเตรียม
สิ่งใด
แนวคิด — เบื้องหลังการทำงาน
สถาปัตยกรรมเป้าหมาย Snowflake ยังเขียนข้อมูลการเดินทางใหม่ต่อไปผ่าน producer การเดินทาง ขณะที่สคริปต์ Python แบบครั้งเดียว backfill 50 ล้านแถวที่มีอยู่เข้าไปใน ClickHouse — สองระบบทำงานเคียงข้างกันตลอดช่วงการย้ายระบบ ไม่ใช่การตัดสวิตช์
บนฝั่ง ClickHouse trips_raw คือตารางรองรับที่สคริปต์เขียนลงไป จากนั้น dbt จะสร้าง
staging view และส่วนที่เหลือของเลเยอร์ analytics ทับบนมัน — สคีมาที่โมดูลนี้สร้าง แต่ยังไม่
ใส่ข้อมูลนอกเหนือจาก trips_raw:
คำอธิบายสีในไดอะแกรม:
- เขียว — แหล่งข้อมูล (producer การเดินทาง ทั้งก่อนและหลังตัดสวิตช์)
- น้ำเงิน — ตารางของ Snowflake
- ส้ม — โมเดลและไปป์ไลน์ของ dbt
- แดง — ตารางและ materialized view ของ ClickHouse
- ฟ้าน้ำทะเล — แดชบอร์ด Apache Superset
- ลูกศรเส้นประ — ทางไหลหลังตัดสวิตช์
ทำไมต้องเป็นสคริปต์ Python ไม่ใช่ connector ในตัว มีหลายวิธีในการย้ายข้อมูลจาก Snowflake ไป ClickHouse แล็บนี้ใช้สคริปต์ Python แบบแบตช์ — นี่คือเหตุผล เทียบกับ ทางเลือกอื่น:
| วิธี | ทำงานอย่างไร | ทำไมไม่ใช้ที่นี่ |
|---|---|---|
| ClickPipes (แหล่งข้อมูล Snowflake) | connector ในตัวของ ClickHouse Cloud — zero-ETL, UI แบบจัดการให้ | Snowflake ไม่ใช่แหล่งข้อมูลที่ ClickPipes รองรับ ClickPipes รองรับ Kafka, S3, Kinesis, PostgreSQL CDC, MySQL CDC และ object storage |
| ส่งออกไป S3 → ClickPipes S3 | COPY INTO @stage ส่งออก Parquet/CSV ไป S3; connector ClickPipes S3 โหลดเข้า ClickHouse | ต้องมี bucket ของ S3, IAM role, stage ของ Snowflake และบัญชี AWS เพิ่มขั้นตอนตั้งค่าราว 3 ขั้นก่อนที่ข้อมูลจะเริ่มย้าย ใช้ได้ในโปรดักชันแต่เป็นโครงสร้างพื้นฐานที่มากเกินไปสำหรับแล็บ |
ส่งออกไป S3 → clickhouse-client | ส่งออกไป S3 แบบเดียวกัน แต่โหลดด้วย INSERT INTO ... SELECT FROM s3(...) | ต้องมีสิ่งจำเป็นเบื้องต้นของ S3 เหมือนกัน และยังต้องให้พาร์ตเนอร์จัดการการแบ่งไฟล์และความสามารถทำงานต่อจากจุดเดิมด้วยมือ |
| Snowflake → Kafka → ClickHouse | สตรีม CDC ของ Snowflake ป้อนข้อมูลให้ Kafka topic; connector ClickPipes Kafka นำเข้ามัน | ไปป์ไลน์สตรีมมิงเต็มรูปแบบ — เหมาะกับความต้องการ latency ต่ำกว่านาทีในโปรดักชัน คลัสเตอร์ Kafka หนักเกินไปมากสำหรับสภาพแวดล้อมแล็บ |
| สคริปต์ Python (แล็บนี้) | snowflake-connector-python อ่านเป็นชุด cursor ละ 100K แถว; clickhouse-connect แทรกข้อมูลโดยตรง | ไม่ต้องมีโครงสร้างพื้นฐานเพิ่มเลยนอกจากแพ็กเกจที่แล็บต้องใช้อยู่แล้ว ทำงานต่อจากจุดเดิมได้ผ่าน --resume (watermark max(pickup_at)) แสดงความคืบหน้าแบบเรียลไทม์ ราว 40-50 นาทีสำหรับ 50 ล้านแถวที่ราว 20K แถว/วินาที — ยอมรับได้สำหรับแบบฝึกหัดย้ายระบบครั้งเดียว |
ทำไมสคริปต์ Python เป็นทางเลือกที่ถูกต้องสำหรับแล็บนี้:
- ไม่ต้องมีบัญชี AWS แนวทางที่ใช้ S3 ต้องสร้าง bucket, กำหนดนโยบาย IAM และมี external stage ของ Snowflake — สามขั้นตอนตั้งค่าที่ไม่มีอะไรเกี่ยวกับ ClickHouse
- เบ็ดเสร็จในตัว แพ็กเกจสองตัว (
snowflake-connector-python,clickhouse-connect) ถูกติดตั้งลงใน venv เดียวกับ dbt ไม่มีเซอร์วิสใหม่ ไม่มี credential ใหม่ - ทำงานต่อจากจุดเดิมได้
--resumeทำให้สคริปต์ปลอดภัยที่จะขัดจังหวะและเริ่มใหม่ReplacingMergeTree(_synced_at)รับประกันว่าการแทรกซ้ำเมื่อลองใหม่จะถูกกำจัดข้อมูลซ้ำ อัตโนมัติ - โปร่งใส พาร์ตเนอร์อ่านสคริปต์ได้, เข้าใจการแมปคอลัมน์ และดัดแปลงมันให้เข้ากับสคีมา ของตัวเองได้ — ซึ่งให้ความรู้มากกว่าการคลิกผ่าน UI wizard
การจัดการช่องว่างของการย้ายระบบ producer ของ Snowflake ยังทำงานต่อไปขณะสคริปต์ ย้ายข้อมูลรัน (~40-50 นาที) การเดินทางใดที่ถูกเขียนลง Snowflake ในช่วงเวลานั้นจะไม่อยู่ใน ClickHouse แล็บนี้ปิดช่องว่างนั้นด้วยแนวทางสองรอบตอนตัดสวิตช์ ซึ่งโมดูล 05 จะพาทำโดยตรง:
- หยุด producer ของ Snowflake เพื่อตรึงชุดข้อมูล
- รัน
python scripts/02_migrate_trips.py --resume— ถ่ายโอนเฉพาะแถวส่วนต่าง (ระดับวินาที ไม่ใช่นาที) - เริ่ม producer ของ ClickHouse
การกำจัดข้อมูลซ้ำด้วย ReplacingMergeTree(_synced_at) ตัวเดียวกันที่จัดการการลองใหม่
ตอนย้ายข้อมูล ก็จัดการเรื่องนี้ด้วย: ถ้ามีแถวใดซ้อนทับกันระหว่างการรันในโมดูลนี้กับรอบ
--resume ในภายหลัง แถวที่มี _synced_at ใหม่กว่าจะชนะ
เมื่อไรที่คุณจะเลือก S3 ในโปรดักชัน หากชุดข้อมูลมีมากกว่า 500 ล้านแถว หรือหากค่าใช้จ่าย ของคิวรีบน warehouse ของ Snowflake สำหรับการสแกนทั้งตารางมีนัยสำคัญ เส้นทางส่งออกไป S3 จะดีกว่า: Snowflake ส่งออก Parquet ที่บีบอัดแบบขนาน (เร็วกว่า cursor เดียวมาก) และ ClickHouse ก็โหลดจาก S3 แบบขนานได้ด้วย แนวทางสคริปต์ Python ที่นี่ทำงานได้ดีในสเกลของแล็บ
ความสอดคล้องของการตัดสินใจ ตารางด้านล่างคือรายการการตัดสินใจชุดเดียวกันจากใบงานที่ 1
(การเลือกเอนจิน), 2 (sort key) และ 3 (การแปลงสคีมา) ใน migration-plan.md ตรวจทานกับ
สิ่งที่แล็บนี้สร้างจริง — เทียบกับแผนของคุณเองก่อนจัดเตรียมสิ่งใด:
| การตัดสินใจ | แล็บนี้ทำอะไร | ทำไม |
|---|---|---|
เอนจินของ trips_raw | ReplacingMergeTree(_synced_at) | สคริปต์ย้ายข้อมูล Python ใช้ INSERT แบบแบตช์ที่อาจถูกลองใหม่ถ้าถูกขัดจังหวะ _synced_at DateTime DEFAULT now() ถูกตั้งค่าในทุก INSERT ดังนั้นแถวที่ลองใหม่จะมาถึงทีหลังและมีค่า _synced_at สูงกว่า — แถวที่มาทีหลังชนะตอน RMT กำจัดข้อมูลซ้ำ ทำให้การลองใหม่เป็น idempotent การลองใหม่ของ producer หลังตัดสวิตช์ก็ปลอดภัยด้วยเหตุผลเดียวกัน stg_trips คิวรีด้วย FINAL เพื่อรับประกันหนึ่งแถวต่อหนึ่งการเดินทาง |
เอนจินของ fact_trips | ReplacingMergeTree(updated_at) | การเดินทางอาจถูกแก้ไข (ปรับค่าโดยสาร); updated_at คือคอลัมน์เวอร์ชัน |
เอนจินของ agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | การคำนวณใหม่แบบเลื่อนช่วง = รูปแบบ upsert |
เอนจินของตาราง dim_* | MergeTree() | โหลดใหม่ทั้งหมดในทุกการรัน dbt; ไม่มี upsert |
ORDER BY ของ fact_trips | (toStartOfMonth(pickup_at), pickup_at, trip_id) | Q1-Q7 ทั้งหมดกรองบน pickup_at; trip_id รับประกันความไม่ซ้ำที่ระดับใบ |
ORDER BY ของ agg_hourly_zone_trips | (hour_bucket, zone_id) | ทั้งสองคอลัมน์ปรากฏในคิวรีรวมค่าทุกตัว |
| VARIANT → | String + JSONExtract* | รักษา JSON ดิบไว้; การดึงค่าเกิดขึ้นตอนคิวรี |
| QUALIFY → | subquery ที่ครอบ ROW_NUMBER() | ClickHouse มี clause QUALIFY ในตัวมาตั้งแต่ v24.5 แต่รูปแบบ subquery ถูกสอนเพราะมันพอร์ตไปใช้กับ ClickHouse เวอร์ชันและเอนจิน SQL ที่เก่ากว่าหรือไม่มี QUALIFY ได้ |
| MERGE INTO → | incremental แบบ delete_insert ใน dbt | กลยุทธ์ upsert ตามสำนวนของ dbt-clickhouse; หลีกเลี่ยงการเขียนทั้งตารางใหม่ |
ขั้นที่ 1 — จัดเตรียมคลัสเตอร์ ClickHouse
setup.sh ทำสิ่งเดียว: รัน terraform apply และเขียนรายละเอียดการเชื่อมต่อไปที่
.clickhouse_state มันยังตรวจหา migration-plan.md อีกครั้งก่อนจัดเตรียมสิ่งใด — ดูหัวข้อ
ทำไม ด้านบน — แต่การตรวจนั้นเพียงเตือน มันไม่เคยขัดขวาง
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 ถูก gitignore ไว้ — อย่า commit มันเด็ดขาด
ผลลัพธ์ที่คาดไว้: Terraform สร้าง 2 ทรัพยากร (เซอร์วิส + รายการ IP ที่เข้าถึงได้) ในเวลาประมาณ 2-3 นาที:
Apply complete! Resources: 2 added, 0 changed, 0 destroyed.
Outputs:
clickhouse_host = "abc123xyz.us-east-1.aws.clickhouse.cloud"
clickhouse_port = 8443โฮสต์และพอร์ตถูกบันทึกไว้ที่ .clickhouse_state ให้ source มันในเทอร์มินัลใดก็ได้เพื่อรับ
ค่าการเชื่อมต่อ:
source .clickhouse_stateตรวจสอบ:
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1ขั้นที่ 2 — สร้างตารางเปล่า
ก่อนอื่น สร้าง trips_raw ด้วยมือโดยใช้เอนจินที่ถูกต้อง สคริปต์ย้ายข้อมูลจะโหลดข้อมูลลง
ตารางนี้ในขั้นที่ 3 — มันต้องมีอยู่แล้วโดยเป็น ReplacingMergeTree เพื่อให้คอลัมน์เวอร์ชัน
พร้อมอยู่ก่อนที่แถวใดจะมาถึง
-- 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 ถูกตั้งค่าอัตโนมัติในทุก INSERT หากสคริปต์ย้ายข้อมูลถูกขัดจังหวะและถูกรันใหม่
ด้วย --resume แถวซ้ำสำหรับ trip_id เดียวกันอาจมีอยู่ชั่วครู่ — RMT จะเก็บแถวที่มาทีหลัง
(_synced_at สูงกว่า) stg_trips คิวรี trips_raw FINAL เพื่อบังคับการกำจัดข้อมูลซ้ำก่อนที่
โมเดลปลายน้ำใดจะเห็นข้อมูล
ถัดไป ใส่ข้อมูลอ้างอิงของโซน นี่เป็นข้อมูลคงที่ (265 โซนของ NYC TLC) ที่ stg_taxi_zones
ของ dbt อ่านเป็นแหล่งข้อมูล
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(หรือวางเนื้อหาของ scripts/00_seed_zones.sql ลงใน SQL console ของ ClickHouse โดยตรงแทน)
ตั้งค่าโปรไฟล์ dbt dbt_project.yml ของโปรเจกต์นี้ประกาศ profile: 'nyc_taxi_ch'
หากไม่มีโปรไฟล์ที่ตรงกันใน ~/.dbt/profiles.yml dbt run จะล้มเหลวทันทีด้วย
Could not find profile named 'nyc_taxi_ch' เทมเพลตอยู่ที่
workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch/profiles.yml.example
โมดูล 01 เขียน ~/.dbt/profiles.yml ไว้แล้วโดยมีโปรไฟล์ nyc_taxi: สำหรับ Snowflake และ
ลูปรีเฟรชในขั้นที่ 4 ของมันยังคิวรีกับโปรไฟล์นั้นต่อไปตราบเท่าที่ producer ของ Snowflake
ทำงานอยู่ อย่าแทนที่ไฟล์นั้นด้วยเทมเพลตของ ClickHouse — การเขียนทับด้วย
profiles.yml.example จะลบโปรไฟล์ nyc_taxi: และทำให้ลูปรีเฟรชของโมดูล 01 พัง แต่ให้
เปิดเทมเพลตและผสานบล็อก nyc_taxi_ch: ของมันเข้าไปใน ~/.dbt/profiles.yml ที่คุณมีอยู่
เป็นโปรไฟล์ระดับบนสุดตัวที่สอง เคียงข้างกับ 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: อ่าน CLICKHOUSE_HOST, CLICKHOUSE_USER และ CLICKHOUSE_PASSWORD จาก
สภาพแวดล้อมผ่าน env_var() ดังนั้นต้อง source .env และ .clickhouse_state ก่อนคำสั่ง
dbt ใดในโมดูลนี้ — dbt run ด้านล่างทำสิ่งนี้อยู่แล้ว เช่นเดียวกับโปรไฟล์ของ Snowflake
~/.dbt/profiles.yml เก็บ credential และถูก gitignore ไว้ อย่า commit มันเด็ดขาด และการ
ผสานนี้เพิ่ม credential ชุดที่สองลงในไฟล์ที่มีอยู่หนึ่งชุดแล้ว
ตรวจสอบ:
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 ClickHouseจากนั้นรัน dbt run เพื่อสร้างตาราง analytics และ staging view:
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 pointคาดว่า: สร้างโมเดลราว 8 ตัวในเวลาไม่ถึง 2 นาที (ทุกตารางว่างเปล่า)
ตรวจสอบ:
# 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: ReplacingMergeTreeตาราง analytics ทั้งหกตารางและ staging view ทั้งสองตัวมีอยู่แล้ว แต่ทุกตัวยังว่างเปล่า —
dbt run เพียงสร้างสคีมาของมัน ตารางเดียวที่มีข้อมูลหลังขั้นตอนนี้คือ trips_raw และมัน
ก็ยังไม่มีข้อมูลเช่นกัน นั่นมาในขั้นถัดไป
ขั้นที่ 3 — ย้ายข้อมูล
โหลดทุกแถวจาก Snowflake NYC_TAXI_DB.RAW.TRIPS_RAW เข้าไปใน ClickHouse
default.trips_raw โดยใช้สคริปต์ย้ายข้อมูลแบบแบตช์ที่เขียนด้วย 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.pyผลลัพธ์ที่คาดไว้ (ประมาณ 40-50 นาทีสำหรับ 50 ล้านแถว):
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
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 ยังเขียนลง TRIPS_RAW ตลอดเวลาที่สคริปต์นี้รัน ดังนั้น ClickHouse
จะตามหลังอยู่ราวเท่ากับความยาวของการถ่ายโอนนี้ — ช่องว่างนั้นเป็นเรื่องที่คาดไว้แล้วและ
ถูกจัดการในโมดูล 05 ไม่ใช่ที่นี่
หากสคริปต์ถูกขัดจังหวะ ให้รันมันใหม่ด้วย --resume เพื่อทำต่อจากจุดตรวจล่าสุด:
python scripts/02_migrate_trips.py --resume--resume อ่าน max(pickup_at) จาก ClickHouse และข้ามแถวที่โหลดไปแล้ว ดังนั้นจึงปลอดภัย
เสมอที่จะขัดจังหวะสคริปต์นี้และเริ่มมันใหม่ — คุณจะไม่ลงเอยกับการโหลดที่ไม่สมบูรณ์และ
กู้คืนไม่ได้
วิธีตรวจสอบว่าคุณทำเสร็จแล้ว
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: 1สถานะปลายทาง
เซอร์วิส ClickHouse Cloud ทำงานอยู่และเข้าถึงได้ และ .clickhouse_state ถูกเขียนลงดิสก์
โดย setup.sh พร้อม CLICKHOUSE_HOST และ CLICKHOUSE_PORT ตารางเป้าหมายและ
staging view ทั้งหมดมีอยู่ — default.trips_raw, staging view สองตัว (stg_trips,
stg_taxi_zones), ตาราง analytics หกตาราง (fact_trips, agg_hourly_zone_trips,
dim_taxi_zones, dim_payment_type, dim_vendor, dim_date) และ refreshable
materialized view ชื่อ analytics.mv_live_trip_feed — รวมเจ็ดอ็อบเจกต์ในสคีมา analytics
ซึ่งสร้างโดย dbt run ของโมดูลนี้ ทุกตารางใน analytics ยังว่างเปล่า ยกเว้น
mv_live_trip_feed ซึ่งมีแถว snapshot หนึ่งแถวที่ dbt run ผลิตขึ้นเมื่อมันสร้าง view นั้น
นอกจากนั้นมีแค่ default.trips_raw ที่มีข้อมูล: ราว 50 ล้านแถว ย้ายมาโดยสคริปต์ย้ายข้อมูล
Python ในขั้นที่ 3
producer ของ Snowflake ยังทำงานอยู่ มันไม่เคยถูกหยุดในโมดูลนี้และไม่หยุดที่นี่
การเดินทางทุกรายการที่ถูกเขียนลง TRIPS_RAW ของ Snowflake หลังจากแบตช์สุดท้ายของ
สคริปต์ย้ายข้อมูล เป็นแถวที่ ClickHouse ไม่มี ดังนั้น ClickHouse ตอนนี้ตามหลัง Snowflake
อยู่ราวเท่ากับความยาวของหน้าต่างการย้ายข้อมูล (~40-50 นาที บวกกับเวลาที่ใช้ตั้งค่าใน
โมดูลนี้) ช่องว่างนั้นมีจริงและโตขึ้นเรื่อย ๆ ตราบเท่าที่ producer ยังทำงาน อย่าปิดมันใน
โมดูลนี้ การตัดสวิตช์ในโมดูล 05 ปิดมันโดยเจตนา ในขั้นตอนสองรอบที่ควบคุมได้ ซึ่งวัดขนาด
ของช่องว่างก่อนกำจัดมัน — การหยุด producer หรือรันสคริปต์ย้ายข้อมูลใหม่ตอนนี้จะเอาสิ่งที่
โมดูล 05 ถูกสร้างขึ้นมาเพื่อสาธิตออกไปพอดี
ใบงานที่ 5: การออกแบบโมเดล dbt
ตั้งค่า materialization, เอนจิน และกลยุทธ์ incremental ให้แต่ละโมเดล dbt พร้อมผลตรวจทันทีในทุกคำตอบ
04 สร้างไปป์ไลน์ dbt ขึ้นใหม่
สร้างไปป์ไลน์ Medallion ขึ้นใหม่บน ClickHouse ด้วย dbt-clickhouse — โมเดล incremental แบบ delete_insert, ReplacingMergeTree, refreshable materialized view — และสร้าง dictionary ของโซน