Snowflake MigrationClickHouse Workshops

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 — สองระบบทำงานเคียงข้างกันตลอดช่วงการย้ายระบบ ไม่ใช่การตัดสวิตช์

ทางไหลของข้อมูลในการย้ายระบบ: producer การเดินทางของ Snowflake ยังเขียนลง TRIPS_RAW ต่อไปขณะที่สคริปต์ Python แบบครั้งเดียวย้าย 50 ล้านแถวเป็นชุดละ 100,000 แถวเข้าไปใน ClickHouse Cloud

บนฝั่ง ClickHouse trips_raw คือตารางรองรับที่สคริปต์เขียนลงไป จากนั้น dbt จะสร้าง staging view และส่วนที่เหลือของเลเยอร์ analytics ทับบนมัน — สคีมาที่โมดูลนี้สร้าง แต่ยังไม่ ใส่ข้อมูลนอกเหนือจาก trips_raw:

ฝั่งเป้าหมาย ClickHouse: trips_raw บน ReplacingMergeTree ป้อนข้อมูลให้ staging view ที่ dbt สร้าง, ตาราง fact และมิติ, ค่ารวมรายชั่วโมง และ dictionary ของโซน

คำอธิบายสีในไดอะแกรม:

  • เขียว — แหล่งข้อมูล (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 S3COPY 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 จะพาทำโดยตรง:

  1. หยุด producer ของ Snowflake เพื่อตรึงชุดข้อมูล
  2. รัน python scripts/02_migrate_trips.py --resume — ถ่ายโอนเฉพาะแถวส่วนต่าง (ระดับวินาที ไม่ใช่นาที)
  3. เริ่ม 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_rawReplacingMergeTree(_synced_at)สคริปต์ย้ายข้อมูล Python ใช้ INSERT แบบแบตช์ที่อาจถูกลองใหม่ถ้าถูกขัดจังหวะ _synced_at DateTime DEFAULT now() ถูกตั้งค่าในทุก INSERT ดังนั้นแถวที่ลองใหม่จะมาถึงทีหลังและมีค่า _synced_at สูงกว่า — แถวที่มาทีหลังชนะตอน RMT กำจัดข้อมูลซ้ำ ทำให้การลองใหม่เป็น idempotent การลองใหม่ของ producer หลังตัดสวิตช์ก็ปลอดภัยด้วยเหตุผลเดียวกัน stg_trips คิวรีด้วย FINAL เพื่อรับประกันหนึ่งแถวต่อหนึ่งการเดินทาง
เอนจินของ fact_tripsReplacingMergeTree(updated_at)การเดินทางอาจถูกแก้ไข (ปรับค่าโดยสาร); updated_at คือคอลัมน์เวอร์ชัน
เอนจินของ agg_hourly_zone_tripsReplacingMergeTree(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: true

nyc_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 ถูกสร้างขึ้นมาเพื่อสาธิตออกไปพอดี

ในหน้านี้

Track your progress?

Optional. We email a link to confirm your address; progress records once you open it.

Please use your work email address, not a personal one.

Progress tracking also requires accepting the current Terms of Service in Privacy settings.

TH