Snowflake MigrationClickHouse Workshops

ตัวอย่างที่ทำเสร็จแล้ว: แผนฉบับสมบูรณ์

แผนการย้ายระบบสำหรับ workload NYC taxi ที่กรอกเสร็จแล้ว ใช้เทียบกับแผนของคุณเองหลังจากเขียนเสร็จ

นี่คือคำตอบที่ทำครบทั้งหมดสำหรับเวิร์กชีตทั้งห้า (1, 2, 3, 4, 5) ที่นำไปใช้กับ workload NYC Taxi ใช้มันเพื่อ:

  • ตรวจคำตอบในเวิร์กชีตของคุณหลังทำแต่ละส่วนเสร็จ
  • เข้าใจเหตุผลเบื้องหลังการตัดสินใจที่ Part 3 นำไปทำจริง
  • เทียบกับตาราง Decision Alignment ของ Part 3 ถ้าคุณเลือกต่างออกไป

นี่คือเฉลย — อย่ากรอกลงในนี้เป็นแผนของคุณ ให้กรอกใน migration-plan.md แทน


รายการตรวจความครบถ้วน

  • การเลือก engine: เสร็จแล้ว
  • การออกแบบ sort key: เสร็จแล้ว
  • การแปลง schema: เสร็จแล้ว
  • แผนการย้ายเป็นระลอก: เสร็จแล้ว
  • การออกแบบโมเดล dbt: เสร็จแล้ว

ส่วนที่ 1: สรุปโปรไฟล์

ตัวชี้วัดค่า
จำนวนตารางทั้งหมด7 (TRIPS_RAW, FACT_TRIPS, AGG_HOURLY_ZONE_TRIPS, DIM_TAXI_ZONES, DIM_PAYMENT_TYPE, DIM_VENDOR, DIM_DATE)
จำนวน view ทั้งหมด2 (STG_TRIPS, STG_TAXI_ZONES)
Streams1 (TRIPS_CDC_STREAM บน TRIPS_RAW)
Tasks2 (CDC_CONSUME_TASK, HOURLY_AGG_TASK)
จำนวนแถวทั้งหมดใน TRIPS_RAW~50,000,000
ช่วงวันที่หน้าต่างเลื่อน 4 ปี สิ้นสุดที่เวลาตอน setup
คอลัมน์ VARIANT1 (TRIPS_RAW.TRIP_METADATA)
การใช้ QUALIFY ที่ตรวจพบ1 (query Q3)
การใช้ MERGE INTO ที่ตรวจพบ2 (HOURLY_AGG_TASK, CDC_CONSUME_TASK)

ส่วนที่ 2: บัญชีรายการออบเจ็กต์

ออบเจ็กต์ชนิดSchemaจำนวนแถวระดับความซับซ้อนหมายเหตุ
trips_rawTableraw~50MBRMT ที่มีคอลัมน์ version _synced_at; การซ้อนกันของ bulk + CDC ต้องกำจัดข้อมูลซ้ำ; stg_trips ต้องใช้ FINAL
stg_tripsdbt Viewstaging—BJSONExtract สำหรับ TRIP_METADATA; ต้องทดสอบ JSON path
stg_taxi_zonesdbt Viewstaging—Aส่งผ่านตรง ๆ; ง่ายมาก
int_trips_enricheddbt Ephemeralstaging—ACTE; ความต่างของ SQL จัดการในโมเดลแม่
fact_tripsdbt Incrementalanalytics~50MCengine RMT; delete_insert; เขียน QUALIFY ใหม่; ต้องใช้ FINAL
agg_hourly_zone_tripsdbt Incrementalanalytics~140KBRMT; หน้าต่างคำนวณใหม่แบบเลื่อน 2 ชั่วโมง; ทดสอบขอบเขต partition อย่างระมัดระวัง
dim_taxi_zonesdbt Tableanalytics265Aข้อมูลอ้างอิงคงที่; โหลดใหม่ทั้งชุด; ง่ายมาก
dim_payment_typedbt Tableanalytics6Aข้อมูลอ้างอิงคงที่; ง่ายมาก
dim_vendordbt Tableanalytics3Aข้อมูลอ้างอิงคงที่; ง่ายมาก
taxi_zones_dictDictionaryanalytics265Bไวยากรณ์เฉพาะของ ClickHouse; dictGet() ตอน query
mv_hourly_revenueRefreshable MVanalytics—Bไวยากรณ์ REFRESH EVERY; ยืนยันการแทนที่แบบ atomic
TRIPS_CDC_STREAM / CDC_CONSUME_TASKSnowflake Stream + Task——Dไม่มีสิ่งเทียบเท่าใน ClickHouse; แทนด้วยการ cutover producer โดยตรงใน Part 3

ส่วนที่ 3: การตัดสินใจเลือก Engine

ตารางEngineคอลัมน์ Versionเหตุผล
trips_rawReplacingMergeTree(_synced_at)_synced_atสคริปต์ย้ายข้อมูล Python (scripts/02_migrate_trips.py) อาจลอง batch ใหม่และ insert trip_id เดิมซ้ำ หลัง cutover producer ที่ทำงานสดก็อาจลองใหม่เมื่อเกิดความล้มเหลวชั่วคราวได้เช่นกัน _synced_at DateTime DEFAULT now() ถูกตั้งอัตโนมัติตอน INSERT — การลองใหม่ครั้งหลังจะมี timestamp สูงกว่า จึงทำให้ RMT เก็บการเขียนล่าสุดไว้ stg_trips query ด้วย FINAL เพื่อบังคับกำจัดข้อมูลซ้ำก่อนโมเดลปลายทางใดจะรัน
fact_tripsReplacingMergeTree(updated_at)updated_atทริปสามารถถูกแก้ไขได้ (ปรับค่าโดยสาร เปลี่ยนสถานะ) trip_id เดิมถูก insert ซ้ำด้วยค่าที่อัปเดตแล้ว updated_at เพิ่มขึ้นแบบ monotonic ทุกครั้งที่แก้ไข — ค่าที่สูงกว่าชนะตอน RMT กำจัดข้อมูลซ้ำ ให้ query ด้วย FINAL เสมอ
agg_hourly_zone_tripsReplacingMergeTree(updated_at)updated_atdbt คำนวณ 2 ชั่วโมงล่าสุดใหม่แล้ว insert ซ้ำ ถ้าไม่มี RMT ค่า aggregate เก่าและใหม่จะสะสมและนับซ้ำ การตั้ง updated_at เป็น now() ในทุกรอบ dbt run ทำให้ค่าล่าสุดชนะ
dim_taxi_zonesMergeTree()—dbt โหลดใหม่ทั้งชุด (สลับตารางแบบ atomic (สร้างใหม่ทั้งชุด)) ข้อมูลซ้ำสะสมไม่ได้ ไม่ต้องกำจัดข้อมูลซ้ำ
dim_payment_typeMergeTree()—เหมือนกัน — โหลดใหม่ทั้งชุด
dim_vendorMergeTree()—เหมือนกัน — โหลดใหม่ทั้งชุด
mv_hourly_revenueMergeTree()—REFRESHABLE MV แทนที่ชุดผลลัพธ์ทั้งชุดแบบ atomic ในทุกครั้งที่ REFRESH ไม่มีการ upsert

ส่วนที่ 4: การออกแบบ Sort Key

ตารางORDER BYเหตุผล
trips_raw(pickup_at, trip_id)การสแกนช่วงเวลากรองด้วย pickup_at ก่อน trip_id คือคีย์กำจัดข้อมูลซ้ำของ RMT — มันต้องอยู่ใน ORDER BY เพื่อให้ RMT ระบุได้ว่าแถวใดซ้ำกัน วาง pickup_at ก่อนเพราะการสแกนช่วงเชิงวิเคราะห์เป็นหลัก และวาง trip_id ท้ายสุดเพราะมี cardinality สูงและทำหน้าที่เพียงเป็นตัวแยกความไม่ซ้ำ
fact_trips(toStartOfMonth(pickup_at), pickup_at, trip_id)query เชิงวิเคราะห์ทั้ง 7 ตัวกรองด้วย pickup_at คำนำหน้าระดับเดือนจับกลุ่มข้อมูลของเดือนปฏิทินไว้ในบล็อกที่อยู่ติดกัน — ทำให้ข้าม block แบบหยาบได้สำหรับการ aggregate รายเดือน โดยไม่ต้องเพิ่ม PARTITION BY วาง trip_id ท้ายสุดเพื่อความไม่ซ้ำของ RMT โดยไม่รบกวนการข้าม block
agg_hourly_zone_trips(hour_bucket, zone_id)Q6 (และ query aggregation ทุกตัว) กรองด้วย hour_bucket และ zone_id hour_bucket มีค่าที่แตกต่างกัน ~35K ค่า ส่วน zone_id มี 265 ค่า วาง hour_bucket ก่อนเพราะการสแกนช่วงเวลาเป็นรูปแบบการเข้าถึงหลัก และ zone_id เป็นลำดับที่สองสำหรับการกรองรอง
dim_taxi_zones(location_id)265 แถว = หนึ่ง granule ORDER BY ไม่มีผลต่อประสิทธิภาพ การใช้ location_id ซึ่งเป็น join key เป็นธรรมเนียมทั่วไปและช่วยให้อ่านง่าย

ส่วนที่ 5: บันทึกการแปลง Schema

คอลัมน์ชนิดใน Snowflakeชนิดใน ClickHouseเหตุผลของการตัดสินใจ
TRIP_METADATAVARIANTStringเก็บ JSON ดิบไว้ตรงตัว JSONExtract* จัดการ path ใดก็ได้ตอน query Map(String,String) ทำให้เสียโครงสร้างที่ซ้อนกัน ส่วน Tuple ต้องมี schema ตายตัว String เป็นทางเลือกที่ปลอดภัยสำหรับ JSON ที่ไม่กำหนดรูป
PICKUP_DATETIME / PICKUP_ATTIMESTAMP_NTZ(9)DateTime64(3, 'UTC')ความละเอียดระดับมิลลิวินาทีเพียงพอสำหรับ timestamp ของทริป ระดับนาโนวินาที (9) เกินความจำเป็น การใส่ 'UTC' ทำให้ timezone ชัดเจนและเลี่ยงเรื่องเซอร์ไพรส์จาก DST ในการ aggregate ตามช่วงเวลา
PICKUP_LOCATION_IDINTEGERUInt16ค่า 1–265 UInt8 สูงสุด 255 (เล็กเกินไป) UInt16 สูงสุด 65535 (ถูกต้อง) 2 ไบต์เทียบกับ 4 ไบต์ของ Int32 — ประหยัดราว 95MB แบบไม่บีบอัดตลอด 50M แถวต่อหนึ่งคอลัมน์
VENDOR_IDINTEGERUInt8ค่า 1–3 UInt8 สูงสุด 255 — ถูกต้อง 1 ไบต์ต่อแถว
DRIVER_RATINGFLOATNullable(Float32)เป็น NULL บ่อย (ไม่ใช่ทุกทริปมีคะแนน) Nullable รักษาความหมายของ null ให้ถูกต้อง Float32 เพียงพอสำหรับช่วง 1.0–5.0 Float64 จะเปลืองพื้นที่จัดเก็บโดยไม่ได้ความละเอียดที่มีความหมายเพิ่มขึ้น
UPDATED_ATTIMESTAMP_NTZ(9)DateTime64(3, 'UTC')คอลัมน์ version สำหรับ ReplacingMergeTree ต้องใช้ DateTime64 ไม่ใช่ DateTime — การแก้ไขสองครั้งในวินาทีเดียวกันจะไม่แน่นอนถ้าใช้ความละเอียดระดับวินาที ความละเอียดระดับมิลลิวินาทีทำให้ลำดับการกำจัดข้อมูลซ้ำถูกต้อง

การแปลงฟังก์ชันที่จำเป็น

นิพจน์ของ Snowflakeสิ่งเทียบเท่าใน 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 — ลบแถวที่คีย์ตรงกัน แล้ว INSERT แถวใหม่ทั้งหมด

ส่วนที่ 6: ระลอกการย้ายระบบ

ระลอกออบเจ็กต์Dependencyหมายเหตุ
ระลอก 0trips_raw (schema), dim_taxi_zones, dim_payment_type, dim_vendorไม่มีdbt สร้างตารางเปล่า ตาราง dim ถูกเติมจากข้อมูลอ้างอิงคงที่ได้ทันที (ไม่ขึ้นกับ trips) รัน: dbt run --select trips_raw dim_*
ระลอก 1การโหลดข้อมูลก้อนใหญ่ด้วย Python (scripts/02_migrate_trips.py)ระลอก 0 (schema ของ trips_raw ต้องมีอยู่)50M แถวจาก Snowflake TRIPS_RAW ทำต่อจากที่ค้างได้ด้วย --resume ยืนยันจำนวนแถวด้วย scripts/01_verify_migration.sh ~40-50 นาที
ระลอก 2stg_trips, stg_taxi_zones, int_trips_enriched, fact_trips, agg_hourly_zone_tripsระลอก 1 เสร็จ (trips_raw มีข้อมูล) + ระลอก 0 (ตาราง dim มีอยู่)dbt run เต็มรูปแบบ stg_trips อ่าน trips_raw; int_trips_enriched join กับ dim; fact_trips และ agg_hourly_zone_trips สร้างต่อจากนั้น
ระลอก 3taxi_zones_dict, mv_live_trip_feedระลอก 2 (dim_taxi_zones มีข้อมูลสำหรับ dict; fact_trips มีข้อมูลสำหรับ MV)dictionary ถูกสร้างผ่าน scripts/04_create_dictionary.sql Refreshable MV ถูกสร้างผ่านโมเดล dbt ส่วนการเปิดช่วงรีเฟรชของมันภายหลังเป็นขั้นตอน ALTER TABLE ... MODIFY REFRESH ที่ทำด้วยมือ ไม่ใช่สิ่งที่ dbt รันให้อัตโนมัติ
ระลอก 4การ cutover producer (scripts/03_cutover.sh)ระลอก 1 เสร็จ (ยืนยันการโหลดก้อนใหญ่แล้ว) + ระลอก 2 เสร็จ (สร้างชั้น analytics แล้ว)หยุด producer ของ Snowflake; เริ่ม producer ของ ClickHouse ที่เขียนตรงไปที่ ClickHouse Cloud; รัน dbt run เพื่อเติม agg_hourly_zone_trips ด้วยข้อมูลสด

ทะเบียนความเสี่ยง (ออบเจ็กต์ระดับ C/D)

ออบเจ็กต์ความเสี่ยงวิธีตรวจสอบ
fact_tripsQuery ที่ไม่ใช้ FINAL จะนับเกินในช่วงที่ merge ล่าช้า ช่วง partition ของ delete_insert ต้องถูกจำกัดขอบเขตให้อยู่ในคำนำหน้าของ ORDER BY เพื่อไม่ให้ลบ partition ที่ไม่ใช่เป้าหมายSELECT COUNT(*) FINAL ตรงกับ Snowflake ± ความล่าช้าของ CDC รัน dbt test เทียบผลของ Q3 ระหว่างสองระบบ ตรวจหา trip_id ที่ซ้ำ: SELECT trip_id, count() FROM fact_trips GROUP BY trip_id HAVING count() > 1 LIMIT 10
agg_hourly_zone_tripsหน้าต่างคำนวณใหม่แบบเลื่อน 2 ชั่วโมงต้องกำหนดขอบเขตของช่วงที่ลบให้ถูกต้อง ถ้ากว้างเกินไป ค่า aggregate เก่าจะถูกลบ ถ้าแคบเกินไป ค่า aggregate ที่ล้าสมัยจะยังอยู่สุ่มตรวจ tuple (hour_bucket, zone_id) ที่เจาะจงเทียบกับ Snowflake ยืนยันว่า trip_count รวมของทุกโซนตรงกับ AGG_HOURLY_ZONE_TRIPS ของ Snowflake ในช่วงเวลาเดียวกัน
การ cutover producerสคริปต์ย้ายข้อมูลที่ถูกขัดจังหวะกลางทางทำให้จำนวนแถวขาด รันซ้ำด้วย --resume เพื่อเติมส่วนที่ขาด producer ที่ลองใหม่หลัง cutover อาจ insert ทริปที่มีอยู่ใน ClickHouse แล้วซ้ำscripts/01_verify_migration.sh — ตรวจความเท่ากันของจำนวนแถวระหว่าง Snowflake และ ClickHouse ReplacingMergeTree(_synced_at) จัดการการ insert ซ้ำแบบ idempotent

ส่วนที่ 7: ช่องว่างของ dialect ที่ทราบแล้ว

  • QUALIFY — กระทบ: Q3 (queries/q03_top_trips_qualify.sql)
  • VARIANT colon-path — กระทบ: Q4, Q5 (การเข้าถึง JSON ของ TRIP_METADATA)
  • LATERAL FLATTEN — ไม่ได้ใช้ใน workload นี้; VARIANT เข้าถึงผ่าน colon-path ไม่ใช่ FLATTEN
  • MERGE INTO — กระทบ: โมเดล incremental ของ dbt (fact_trips, agg_hourly_zone_trips)
  • Snowflake Streams → การ cutover producer (การเขียนสดไปที่ ClickHouse โดยตรงหลัง cutover)
  • ความต่างของฟังก์ชันวันที่ — กระทบ: Q1 (DATE_TRUNC), Q3 (DATEADD), Q4 (DATEDIFF)

ส่วนที่ 8: กลยุทธ์การย้ายระบบ

การเคลื่อนย้ายข้อมูล: สคริปต์ย้ายข้อมูล Python (scripts/02_migrate_trips.py)

ทำไมเลือกสคริปต์ Python แทนการรีเลย์ผ่าน object storage หรือ ClickPipes?

  • remoteSecure() มีไว้สำหรับการถ่ายข้อมูลจาก ClickHouse ไป ClickHouse — ใช้ไม่ได้ในกรณีนี้
  • การรีเลย์ผ่าน object storage (Snowflake → S3 → ฟังก์ชันตาราง S3 ของ ClickHouse) ใช้ได้ แต่เพิ่มความซับซ้อน: ต้องจัดเตรียม bucket S3, IAM role และ COPY INTO ของ Snowflake — เป็นภาระที่ไม่จำเป็นสำหรับ lab
  • ClickPipes ไม่รองรับ Snowflake เป็นแหล่งข้อมูล แหล่งที่รองรับคือ Kafka, S3, Kinesis, PostgreSQL CDC และ MySQL CDC
  • สคริปต์ Python ใช้ snowflake-connector-python และ clickhouse-connect — แพ็กเกจที่ติดตั้งไว้แล้วสำหรับ lab มันแสดงความคืบหน้าแบบเรียลไทม์ รองรับ --resume เมื่อถูกขัดจังหวะ และตรวจอ่านโค้ดได้ทั้งหมด

Incremental strategy (dbt): delete_insert

ทำไมเลือก delete_insert แทนกลยุทธ์ append หรือ merge?

  • append insert แถวใหม่โดยไม่แตะแถวที่มีอยู่ สำหรับ fact_trips ที่แถวสามารถถูกอัปเดตได้ วิธีนี้สร้างข้อมูลซ้ำ ไม่ถูกต้อง
  • merge (ถ้ามี) จะใกล้เคียง MERGE INTO ของ Snowflake ที่สุด แต่กลยุทธ์ merge ของ dbt-clickhouse มีข้อจำกัดเมื่อใช้กับ ReplacingMergeTree และไม่ใช่วิธีที่แนะนำ
  • delete_insert ลบแถวในช่วงคีย์ของ batch ที่เข้ามา แล้ว insert แถวใหม่ทั้งหมด วิธีนี้เป็น idempotent (รันซ้ำได้ผลเหมือนเดิม) จัดการได้ทั้งการ insert และการอัปเดต และทำงานถูกต้องกับ ReplacingMergeTree มันเป็นคำแนะนำมาตรฐานของชุมชน dbt-clickhouse สำหรับรูปแบบ upsert

ส่วนที่ 9: เกณฑ์การ Cutover

เกณฑ์ค่าเกณฑ์วัดด้วย
ความเท่ากันของจำนวนแถวตรงกัน ≥ 99.9% (คาดว่า CH ≥ SF หลัง cutover)scripts/01_verify_migration.sh
ความเท่ากันของ checksumMD5 ตรงกันบนตัวอย่าง 10K แถวscripts/02_validate_parity.sql
อัตราการผ่านของ dbt test100%dbt test ใน dbt/nyc_taxi_dbt_ch
ความเท่ากันของผล queryquery ทั้ง 7 ตัวคืนผลเหมือนกัน (ภายในความคลาดเคลื่อนของ floating-point)เทียบด้วยมือจากผลลัพธ์ของ scripts/run_benchmark.sh


ส่วนที่ 10: การออกแบบโมเดล dbt

การเลือก Materialization

โมเดลMaterializationทำไม
stg_tripsviewอ่านและทำความสะอาด trips_raw; ไม่มีการอัปเดตโมเดลนี้; ไม่มีต้นทุนพื้นที่จัดเก็บ; สะท้อนสภาพต้นทางปัจจุบันตลอด
stg_taxi_zonesviewเหมือนกัน — ทำความสะอาดแบบส่งผ่านจากตารางต้นทาง
int_trips_enrichedephemeralตรรกะ join ล้วน ๆ ที่ใช้โดย fact_trips เท่านั้น; การฝังเป็น CTE เลี่ยงตารางจริงที่ซ้ำซ้อน; ไม่มีโมเดลใด query มันโดยตรง
fact_tripsincrementalทริปสามารถถูกแก้ไขภายหลังได้; ควรประมวลผลเฉพาะแถวใหม่และแถวที่อัปเดตในแต่ละรอบ
agg_hourly_zone_tripsincrementalการคำนวณใหม่แบบเลื่อน 2 ชั่วโมงเป็นรูปแบบ incremental — ประมวลผลแถวล่าสุด ไม่ใช่ทั้ง 50M แถว
dim_taxi_zonestable265 โซนคงที่; สร้างใหม่ทั้งชุดในทุกรอบ dbt run ผ่านการสลับตารางแบบ atomic (สร้างใหม่ทั้งชุด); ไม่มีการอัปเดตบางส่วน
dim_payment_typetable6 ประเภทคงที่; เหตุผลเดียวกับ dim_taxi_zones
dim_vendortable3 ผู้ให้บริการ; เหตุผลเดียวกัน

การตั้งค่า Engine

โมเดลENGINEคอลัมน์ Versionทำไม
fact_tripsReplacingMergeTree(updated_at)updated_atทริปสามารถถูกแก้ไขได้; การตั้ง updated_at เป็น now() ในทุกการ insert หมายความว่าเวอร์ชันล่าสุดชนะตอน RMT กำจัดข้อมูลซ้ำเบื้องหลัง; delete_insert เป็นเส้นทางความถูกต้องหลัก และ RMT เป็นตาข่ายรองรับ
agg_hourly_zone_tripsReplacingMergeTree(updated_at)updated_atการคำนวณใหม่แบบเลื่อน insert ค่า aggregate ซ้ำสำหรับคู่ (hour_bucket, zone_id) เดิม; RMT ทำให้ค่า aggregate ที่ล้าสมัยถูกลบตอน merge เบื้องหลัง
dim_taxi_zonesMergeTree()—การโหลดใหม่ทั้งชุดโดย dbt หมายถึงการสลับตารางแบบ atomic (สร้างใหม่ทั้งชุด) ในทุกรอบ; ข้อมูลซ้ำสะสมไม่ได้; ไม่ต้องกำจัดข้อมูลซ้ำ
dim_payment_typeMergeTree()—เหมือนกับ dim_taxi_zones
dim_vendorMergeTree()—เหมือนกับ dim_taxi_zones

Incremental Strategy

โมเดลunique_keyincremental_strategyตัวกรอง incrementalทำไมต้องใช้ตัวกรองนี้
fact_tripstrip_iddelete_insertWHERE updated_at > (SELECT max(updated_at) FROM {{ this }})high-watermark บน updated_at จับได้ทั้งทริปใหม่และทริปที่ถูกแก้ไข (การปรับค่าโดยสาร insert trip_id เดิมซ้ำด้วย pickup_at เดิมแต่ updated_at ใหม่กว่า); watermark ที่ใช้ pickup_at จะพลาดการแก้ไขไปเงียบ ๆ
agg_hourly_zone_trips[hour_bucket, zone_id]delete_insertWHERE pickup_at >= now() - INTERVAL 2 HOURหน้าต่างเลื่อน 2 ชั่วโมงบังคับให้ aggregate ชั่วโมงที่อยู่ตรงขอบใหม่ เพื่อให้จำนวนของชั่วโมงที่ยังไม่ครบถูกแก้ไขเสมอ; high-watermark ที่ใช้ max(pickup_at) จะนับชั่วโมงตรงขอบต่ำกว่าจริงอย่างถาวร

การวาง FINAL

โมเดลมี FINAL ในประโยค FROM หรือไม่ทำไม
stg_tripsมี — FROM trips_raw FINALtrips_raw เป็น ReplacingMergeTree; มันอาจมีแถว trip_id ซ้ำจากการลองใหม่ของสคริปต์ย้ายข้อมูลหรือการลองใหม่ของ producer หลัง cutover stg_trips เป็นจุดบังคับใช้จุดเดียว: กำจัดข้อมูลซ้ำที่นี่เพื่อให้ทุกโมเดลปลายทาง (int_trips_enriched, fact_trips, agg_hourly_zone_trips) ได้รับข้อมูลที่สะอาด
int_trips_enrichedไม่มีอ่านจาก stg_trips (ซึ่งเป็น view) ไม่ใช่ตาราง RMT; FINAL ไม่เกี่ยวกับ view
fact_tripsไม่มี (ในตัวโมเดล)delete_insert ทำให้ fact_trips สะอาดหลังทุกรอบที่ทำงานจบ; การใส่ FINAL ไว้ในโมเดลจะทำให้มันถูกนำไปใช้กับ subquery ของ is_incremental() ที่อ่าน max(updated_at) จาก {{ this }} อย่างสิ้นเปลือง dashboard และการทดสอบ dbt ใช้ FINAL จากภายนอกเมื่อ query fact_trips โดยตรง

นี่คือตัวอย่างที่ทำเสร็จแล้ว migration-plan.md ของคุณควรตรงกับการตัดสินใจสำคัญที่นี่ — หรือบันทึกอย่างชัดเจนว่าทำไมคุณเลือกต่างออกไป

ในหน้านี้

TH