04 สร้างไปป์ไลน์ dbt ขึ้นใหม่
สร้างไปป์ไลน์ Medallion ขึ้นใหม่บน ClickHouse ด้วย dbt-clickhouse — โมเดล incremental แบบ delete_insert, ReplacingMergeTree, refreshable materialized view — และสร้าง dictionary ของโซน
จุดเริ่มต้น
โมดูล 03 เสร็จแล้ว: เซอร์วิส 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 —
dbt run ของโมดูล 03 สร้างทั้งเจ็ดตัวไว้แล้ว ทุกตารางใน analytics ยังว่างเปล่า ยกเว้น
mv_live_trip_feed ซึ่งมีแถว snapshot หนึ่งแถวจากการ build นั้น นอกจากนั้นมีแค่
default.trips_raw ที่มีข้อมูล: ราว 50 ล้านแถว producer ของ Snowflake ยังทำงานอยู่ ดังนั้น
ClickHouse ตามหลัง Snowflake อยู่ราวเท่ากับความยาวของหน้าต่างการย้ายข้อมูล กันเวลาไว้
ประมาณ 30 นาที
ทำไม
โมดูล 03 พิสูจน์ว่า ClickHouse เก็บ 50 ล้านแถวได้ มันไม่ได้พิสูจน์ว่า ไปป์ไลน์ รันบน ClickHouse ได้ — staging view, ตาราง fact แบบ incremental, การโหลดตารางมิติใหม่, เทสต์ที่ จับโมเดลที่พังได้ก่อนที่พาร์ตเนอร์จะเห็น นั่นคือสิ่งที่โมดูลนี้สร้างขึ้นใหม่: โมเดล Medallion ชุดเดียวกันจากโมดูล 01 แสดงออกด้วย dbt-clickhouse แทน dbt-snowflake รันกับตารางที่ โมดูล 03 สร้างไว้แล้ว
ไม่มีอะไรเกี่ยวกับตรรกะของโมเดลที่เปลี่ยน — stg_trips ยังแปลงชนิดข้อมูลและดึงค่า JSON,
int_trips_enriched ยัง join ตารางมิติ, fact_trips ยังลงเอยเป็นหนึ่งแถวต่อหนึ่งการเดินทาง
สิ่งที่เปลี่ยนคือเลเยอร์ materialization ที่อยู่ข้างใต้: ไม่มี MERGE INTO, ไม่มี Snowflake
Task, ไม่มี cluster_by นี่คือโมดูลที่แสดงว่าการย้ายระบบไม่ใช่การเทข้อมูลครั้งเดียวจบ —
ไปป์ไลน์ที่ทีมของพาร์ตเนอร์รันทุกวัน ตามกำหนดเวลาเดิม โดยมีการรัน dbt test ชุดเดิม
เป็นด่านกั้น ยังทำงานได้ต่อไปเมื่อ warehouse ที่อยู่ข้างใต้เปลี่ยนไป
แนวคิด — เบื้องหลังการทำงาน
ชุดโมเดลต่างจากไปป์ไลน์ของ Snowflake ในสี่แง่ เอกสารอ้างอิงคอนฟิกฉบับเต็มอยู่ที่
dbt บน ClickHouse — นี่คือฉบับย่อ
ที่คุณต้องรู้ก่อนรัน dbt run ในขั้นที่ 1 สำหรับไปป์ไลน์ต้นทางที่โมเดลเหล่านี้แทนที่ ดู
dbt บน Snowflake
1. delete_insert แทน MERGE ClickHouse ไม่มีคำสั่ง MERGE INTO ในจุดที่ไปป์ไลน์ของ
Snowflake ใช้ incremental_strategy: merge เพื่อ upsert fact_trips และ
agg_hourly_zone_trips โมเดลของ ClickHouse ใช้ incremental_strategy: delete_insert:
dbt ลบแถวที่ตรงกับ unique_key ของแบตช์ที่เข้ามา แล้วแทรกแบตช์นั้น สำหรับ fact_trips
unique_key คือ trip_id และฟิลเตอร์ incremental ทำ watermark บน updated_at ไม่ใช่
pickup_at — การแก้ค่าโดยสารจะแทรก trip_id เดิมซ้ำโดยมี pickup_at เดิมแต่
updated_at ใหม่กว่า ดังนั้นการ watermark บน pickup_at จะพลาดมันไปอย่างเงียบ ๆ
2. ReplacingMergeTree เป็นตะแกรงกันตกใต้ delete_insert ไม่ใช่ตัวแทนของมัน
โมเดล incremental ทั้งสองตัวถูกประกาศเป็น ReplacingMergeTree(updated_at) หากการรัน
delete_insert เสร็จตามปกติ ตารางจะมีหนึ่งแถวต่อคีย์อยู่แล้วและเอนจินไม่มีอะไรต้องเก็บกวาด
หากการรันถูกขัดจังหวะกลางทาง — พังหลังลบ ก่อนแทรก — การ merge เบื้องหลังจะกำจัดแถว
ตกค้างในที่สุด โดยเก็บแถวที่มี updated_at สูงสุด อย่าพึ่ง ReplacingMergeTree เพียงลำพัง
ให้ทำงานกำจัดข้อมูลซ้ำที่ delete_insert ควรทำ: การ merge เบื้องหลังเป็นแบบอะซิงโครนัส
และอาจล่าช้าเป็นนาทีถึงชั่วโมงกับตารางขนาดนี้
3. refreshable materialized view แทน task ตามกำหนดเวลา ไปป์ไลน์ของ Snowflake ใช้
Task ตามกำหนดเวลาที่รัน stored procedure เพื่อทำให้ค่ารวมแบบเลื่อนช่วงเป็นปัจจุบัน
โปรเจกต์ dbt ของ ClickHouse ประกาศ mv_live_trip_feed ด้วย
materialized = 'materialized_view' และ engine = 'ReplacingMergeTree(refreshed_at)'
แทน — สร้างโดย dbt run เป็น refreshable materialized view เทียบกับ DDL
CREATE TASK ของ Snowflake สามสิบกว่าบรรทัดเพื่อผลลัพธ์เดียวกัน โมดูลนี้ไม่ได้เปิด
ช่วงเวลารีเฟรช การทำเช่นนั้นเป็นคำสั่ง
ALTER TABLE analytics.mv_live_trip_feed MODIFY REFRESH EVERY ... ด้วยมือที่แล็บ
ไม่ได้เขียนสคริปต์ไว้ (ดูโมดูล 05 ว่าทำไม)
4. mv_live_trip_feed ไม่มีคู่เทียบใน Snowflake เลย มันไม่ใช่การแปลโมเดลที่มีอยู่ —
มันเป็นความสามารถใหม่ที่การย้ายระบบเพิ่มเข้ามา materialized view มาตรฐานของ ClickHouse
ทำงานหนึ่งครั้งต่อหนึ่ง INSERT และเห็นเพียงแถวในแบตช์นั้น จึงคำนวณค่ารวมตลอดอายุอย่าง
จำนวนการเดินทางรวมหรือค่าโดยสารเฉลี่ยได้ไม่ถูกต้อง แต่ materialized view แบบ REFRESHABLE
จะรันคิวรีทั้งตัวใหม่ — ที่นี่คือ SELECT ... FROM {{ ref('fact_trips') }} — ตามกำหนดเวลา
ดังนั้นการรีเฟรชทุกครั้งจึงเห็นตารางทั้งตาราง ฝั่ง Snowflake ของเวิร์กช็อปนี้ไม่เคยมีตัวเลือกนี้
โมเดลที่ dbt สร้างจริง และแต่ละตัวได้ข้อมูลเมื่อไร:
| โมเดล | เลเยอร์ | Materialization | ได้ข้อมูลเมื่อ | หมายเหตุ |
|---|---|---|---|---|
stg_trips | staging | View | ขั้นที่ 1 (ทุกการรัน) | แปลงชนิดข้อมูล, JSONExtract* สำหรับ trip_metadata |
stg_taxi_zones | staging | View | ขั้นที่ 1 (ทุกการรัน) | ส่งผ่านตารางมิติของโซน |
int_trips_enriched | staging | Ephemeral | — (แทรกเข้าเป็น CTE) | join ตารางมิติทั้งหมด; ไม่มีตารางจริง |
fact_trips | analytics | Incremental | ขั้นที่ 1 | delete_insert ใช้คีย์ trip_id, watermark บน updated_at; ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id) |
agg_hourly_zone_trips | analytics | Incremental | โมดูล 05 หลังตัดสวิตช์ | หน้าต่างเลื่อน 2 ชั่วโมง; ฟิลเตอร์ incremental ตรงกับแถวจาก producer สดเท่านั้น — ดูขั้นที่ 1 ด้านล่าง |
dim_taxi_zones | analytics | Table | ขั้นที่ 1 | โหลดใหม่ทั้งหมดต่อการรัน; เป็นแหล่งข้อมูลของ dictionary โซนในขั้นที่ 2 |
dim_payment_type | analytics | Table | ขั้นที่ 1 | โหลดใหม่ทั้งหมดต่อการรัน |
dim_vendor | analytics | Table | ขั้นที่ 1 | โหลดใหม่ทั้งหมดต่อการรัน |
dim_date | analytics | Table | ขั้นที่ 1 | แกนวันที่คงที่ 2009-2029; โหลดใหม่ทั้งหมดต่อการรัน |
mv_live_trip_feed | analytics | Materialized view (refreshable) | dbt run ของโมดูล 03; ช่วงเวลารีเฟรชไม่เคยถูกเปิด | ไม่มีคู่เทียบใน Snowflake — ดูข้อ 4 ด้านบน |
ขั้นที่ 1 — ใส่ข้อมูลให้เลเยอร์ analytics
Activate venv ของ dbt-clickhouse ที่คุณสร้างในโมดูล 00 แล้วรัน dbt run เป็นครั้งที่สอง
โมดูล 03 รันมันไปหนึ่งครั้งกับตารางเปล่าเพื่อสร้างสคีมาแล้ว การรันครั้งนี้มีข้อมูลจริงอยู่
ข้างหลัง — trips_raw ตอนนี้มี 50 ล้านแถว
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 อ่านการเชื่อมต่อ ClickHouse จาก ~/.dbt/profiles.yml ซึ่งอ้างจาก
dbt/nyc_taxi_dbt_ch/profiles.yml.example — โปรไฟล์เดียวกับที่โมดูล 03 ใช้สร้างสคีมาเปล่า
คาดว่า: ประมาณ 8-12 นาที (โมเดล incremental ประมวลผล 50 ล้านแถว)
agg_hourly_zone_trips จะว่างเปล่าหลังการรันนี้ — นั่นเป็นสิ่งที่คาดไว้ ไม่ใช่ความ
ล้มเหลว ฟิลเตอร์ incremental ของมันคือ WHERE pickup_at >= now() - INTERVAL 2 HOUR
ซึ่งตรงกับแถวที่ producer สดเขียนเท่านั้น ทุกแถวที่คุณเพิ่งย้ายมาเป็นข้อมูลย้อนหลัง จึงไม่มี
แถวใดตกอยู่ในหน้าต่าง 2 ชั่วโมงที่วัดจากตอนนี้ ตารางนี้จะว่างเปล่าไปจนกว่าการตัดสวิตช์ใน
โมดูล 05 จะเริ่ม producer ของ ClickHouse — อย่าเสียเวลาดีบักเรื่องนี้ว่าเป็นไปป์ไลน์ที่พัง
จากนั้นรันชุดเทสต์:
dbt testคาดว่า: เทสต์ผ่านทั้งหมด
ตรวจสอบ:
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: 265ขั้นที่ 2 — สร้าง dictionary ของโซน
analytics.dim_taxi_zones ตอนนี้มีข้อมูลโซนของ NYC TLC ครบทั้ง 265 โซน ให้สร้าง
analytics.taxi_zones_dict ซึ่งเป็น dictionary ในหน่วยความจำที่หนุนด้วยตารางนั้น เพื่อให้
คิวรีปลายน้ำค้นหาเขต (borough) ของโซนได้ด้วย dictGet() แทนการ JOIN
dictionary ให้อะไรมากกว่าการ join dictionary โหลดเข้าหน่วยความจำครั้งเดียวและอยู่ร้อน
ตลอด การค้นหาในนั้นแทบไม่มีต้นทุนเลยในทุกคิวรีถัดไป ส่วนการ JOIN กับ dim_taxi_zones
ต้องอ่านและจับคู่ตารางมิติใหม่ทุกครั้งที่มันรัน สำหรับตารางอ้างอิงที่เล็กและเปลี่ยนน้อยแบบนี้ —
265 แถว โหลดใหม่ทั้งหมดโดย dbt run ทุกครั้ง — การแลกนั้นเอียงไปข้างเดียวชัดเจน
การรันเบนช์มาร์กในโมดูล 05 คิวรี taxi_zones_dict โดยตรงด้วย dictGet ดังนั้นขั้นตอนนี้
เป็น dependency ที่จำเป็นสำหรับโมดูลนั้น ไม่ใช่ของแถมที่เลือกได้
Source รายละเอียดการเชื่อมต่อ แล้วโหลด DDL ของ dictionary ผ่าน clickhouse-client
หรือ HTTP API — เลือกอันที่คุณมี:
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.sqlตรวจสอบ:
-- 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';วิธีตรวจสอบว่าคุณทำเสร็จแล้ว
SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 millionSELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0ศูนย์ถูกต้องแล้วในที่นี้ ไม่ใช่ความล้มเหลว ฟิลเตอร์ incremental ของ
agg_hourly_zone_trips (WHERE pickup_at >= now() - INTERVAL 2 HOUR) ตรงกับแถวที่
producer สดเขียนเท่านั้น และทุกแถวใน ClickHouse ตอนนี้เป็นข้อมูลย้อนหลังที่สคริปต์ย้ายข้อมูล
ของโมดูล 03 ย้ายมา — ไม่มีแถวใดอายุน้อยกว่า 2 ชั่วโมงเมื่อเทียบกับ now() ตารางนี้จะมีข้อมูล
หลังจากการตัดสวิตช์ในโมดูล 05 เริ่ม producer ของ ClickHouse เท่านั้น จนกว่าจะถึงตอนนั้น
ชาร์ตแดชบอร์ดใดที่หนุนด้วยตารางนี้จะแสดงว่าไม่มีข้อมูล และนั่นเป็นสิ่งที่คาดไว้ ไม่ใช่สิ่งที่
ต้องดีบัก
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 testคาดว่า: เทสต์ผ่านทั้งหมด
-- Should return a borough name, e.g. 'Manhattan'
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));สถานะปลายทาง
เลเยอร์ analytics มีข้อมูลและผ่านการทดสอบแล้ว: fact_trips มีราว 50 ล้านแถว,
dim_taxi_zones, dim_payment_type, dim_vendor และ dim_date โหลดครบแล้ว,
dbt test ผ่านตั้งแต่ต้นจนจบ และ analytics.taxi_zones_dict พร้อมใช้งานและคืนค่าเขต
ผ่าน dictGet() agg_hourly_zone_trips ยังว่างเปล่า — โดยการออกแบบ ไม่ใช่ข้อบกพร่อง —
และจะเป็นเช่นนั้นไปจนถึงการตัดสวิตช์ในโมดูล 05
แดชบอร์ด, เบนช์มาร์ก ClickHouse เทียบกับ Snowflake และการตัดสวิตช์ อยู่ในโมดูล 05 ไม่ใช่โมดูลนี้
producer ของ Snowflake ยังทำงานอยู่ และช่องว่างระหว่าง Snowflake กับ ClickHouse ยังเปิดอยู่ ไม่มีอะไรในโมดูลนี้ที่แตะ producer หรือสคริปต์ย้ายข้อมูล — โมดูล 05 ปิด ช่องว่างนั้นโดยเจตนา ในขั้นตอนสองรอบที่ควบคุมได้แบบเดียวกับที่โมดูล 03 เกริ่นไว้ อย่าหยุด producer ตอนนี้
03 จัดเตรียมและย้ายข้อมูล
จัดเตรียม ClickHouse Cloud ด้วย Terraform, สร้างตารางเป้าหมายจากแผนของคุณ และย้าย 50 ล้านแถวด้วยสคริปต์ย้ายข้อมูล Python ที่ทำงานต่อจากจุดเดิมได้
05 เบนช์มาร์กและตัดสวิตช์
สร้างแดชบอร์ดขึ้นใหม่บน ClickHouse, เบนช์มาร์กคิวรีทั้งเจ็ดตัวกับทั้งสองเอนจิน, ตัดสวิตช์ producer, ตรวจสอบความเท่าเทียม และรื้อระบบ