Snowflake MigrationClickHouse Workshops

dbt บน Snowflake

วิธีสร้าง pipeline แบบ Medallion ต้นทาง: sources, staging views, โมเดล incremental MERGE, snapshots และการทดสอบ

เอกสารนี้อธิบายว่า dbt (data build tool) ถูกใช้อย่างไรใน NYC Taxi Snowflake Migration Lab — มันทำอะไร ทำไมแต่ละส่วนจึงต้องมี และควรคิดถึงมันอย่างไร


สิ่งที่ dbt ทำ (และไม่ทำ)

dbt แปลง ข้อมูลที่อยู่ในฐานข้อมูลของคุณแล้ว มันไม่โหลดข้อมูลจากภายนอก ไม่ย้ายไฟล์ และไม่จัดการ infrastructure หน้าที่ของมันคือรับตารางดิบมาแปลงให้เป็นตารางที่สะอาด ผ่านการทดสอบ และพร้อมใช้งานเชิงวิเคราะห์ — ด้วยการรัน SQL ที่คุณเขียน

ให้คิดว่ามันเป็น build system สำหรับ SQL ทุกไฟล์ .sql ใน models/ คือโมเดลหนึ่งตัวที่จะกลายเป็นตารางหรือ view ใน Snowflake dbt จัดการโค้ดซ้ำซาก CREATE OR REPLACE ให้ แก้ลำดับ dependency ระหว่างโมเดล และรันการทดสอบของคุณ


โครงสร้างโปรเจกต์

dbt/nyc_taxi_dbt/
├── dbt_project.yml          # Project config: name, folder layout, materialization defaults
├── profiles.yml.example     # Connection config template (copy to ~/.dbt/profiles.yml)
├── packages.yml             # Third-party dbt packages
│
├── macros/
│   ├── generate_schema_name.sql   # Overrides dbt's default schema naming logic
│   └── generate_surrogate_key.sql # Wrapper for consistent surrogate key generation
│
└── models/
    ├── sources.yml          # Declares RAW.TRIPS_RAW as an external source
    │
    ├── staging/             # Layer 1: clean and rename raw columns
    │   ├── schema.yml       # Column-level tests for staging models
    │   ├── stg_trips.sql
    │   └── stg_taxi_zones.sql
    │
    ├── intermediate/        # Layer 2: joins and enrichment (no physical table)
    │   └── int_trips_enriched.sql
    │
    └── analytics/           # Layer 3: final tables consumed by dashboards
        ├── schema.yml
        ├── fact_trips.sql
        ├── agg_hourly_zone_trips.sql
        ├── dim_date.sql
        ├── dim_payment_type.sql
        ├── dim_taxi_zones.sql
        └── dim_vendor.sql

สามชั้น (สถาปัตยกรรม Medallion)

ชั้นที่ 1 — Staging (models/staging/)

วัตถุประสงค์: รับข้อมูลดิบตามที่มาถึงจริง ๆ แล้วทำให้ใช้งานได้

โมเดลเหล่านี้ลงใน schema STAGING ในรูปของ views (ไม่มีต้นทุนพื้นที่จัดเก็บ — รันตอน query) โมเดล staging แต่ละตัวทำงานหนึ่งอย่าง:

โมเดลSourceทำอะไร
stg_tripsRAW.TRIPS_RAWเปลี่ยนชื่อคอลัมน์เป็น snake_case, เพิ่ม duration_minutes, แผ่คอลัมน์ VARIANT TRIP_METADATA ออกเป็นคอลัมน์ที่มีชนิดข้อมูลชัดเจน
stg_taxi_zonesANALYTICS.DIM_TAXI_ZONESทำความสะอาดเล็กน้อย เพิ่มการป้องกันด้วย COALESCE และสร้าง node สำหรับ lineage ของ dbt

งานที่สำคัญที่สุดที่นี่คือการแผ่คอลัมน์ VARIANT TRIP_METADATA ไวยากรณ์ colon-path ของ Snowflake ดึงฟิลด์ JSON ที่ซ้อนกันออกมา:

-- Snowflake: colon-path notation
TRIP_METADATA:driver.rating::FLOAT     AS driver_rating,
TRIP_METADATA:app.surge_multiplier::FLOAT AS surge_multiplier

นี่คือหนึ่งในความท้าทายของการย้ายระบบ — ClickHouse ใช้ JSONExtractFloat(TRIP_METADATA, 'driver', 'rating') แทน

ชั้นที่ 2 — Intermediate (models/intermediate/)

วัตถุประสงค์: ทำ join ทั้งหมดไว้ที่เดียว เพื่อไม่ต้องเขียนซ้ำ

int_trips_enriched join stg_trips เข้ากับทุก dimension (zones, payment types, vendors, dates) และให้ผลเป็นแถวกว้างที่ denormalize เต็มที่หนึ่งแถวต่อหนึ่งทริป มันถูกประกาศเป็น ephemeral ซึ่งหมายความว่า dbt จะฝัง SQL ของมันเข้าไปในโมเดลใดก็ตามที่อ้างถึงมัน — ไม่มีการสร้างตารางหรือ view จริงใน Snowflake

-- dbt_project.yml
intermediate:
  +materialized: ephemeral   # compiled inline, no CREATE TABLE

ใช้ ephemeral เมื่อผลลัพธ์ระดับกลางจำเป็นสำหรับโมเดลปลายทางเพียงตัวเดียว และคุณไม่ต้องการจ่ายค่าพื้นที่จัดเก็บหรือ overhead ในการคอมไพล์ query

ชั้นที่ 3 — Analytics (models/analytics/)

วัตถุประสงค์: ตารางสุดท้ายที่พร้อมใช้กับ dashboard

โมเดลเหล่านี้ลงใน schema ANALYTICS มีสองแบบ:

ตาราง dimension แบบคงที่ — ขนาดเล็ก โหลดใหม่ทั้งหมดทุกครั้งที่ dbt run:

โมเดลจำนวนแถวหมายเหตุ
dim_date~7,670Date spine 2009–2029 พร้อมไตรมาสทางบัญชีและวันหยุดราชการสหรัฐ
dim_payment_type6ส่งผ่านมาจากข้อมูล seed
dim_vendor3ส่งผ่านมาจากข้อมูล seed
dim_taxi_zones265ส่งผ่านมาทาง stg_taxi_zones

ตาราง fact/aggregate แบบ incremental — ขนาดใหญ่ อัปเดตด้วย MERGE ในแต่ละรอบ:

โมเดลจำนวนแถวหมายเหตุ
fact_trips50Mหนึ่งแถวต่อหนึ่งทริป denormalize เต็มที่
agg_hourly_zone_trips~9Mจำนวนรายชั่วโมงต่อโซนที่ aggregate ไว้ล่วงหน้า

Materializations

materialization ควบคุมว่า dbt จะสร้างอะไรใน Snowflake สำหรับโมเดลหนึ่ง ๆ

Materializationออบเจ็กต์ใน Snowflakeใช้เมื่อไร
viewCREATE VIEWราคาถูก สะท้อนข้อมูลล่าสุดตลอด ใช้กับ staging
tableCREATE TABLE AS SELECTสร้างใหม่ทั้งหมดทุกรอบ ใช้กับ dimension ขนาดเล็ก
incrementalMERGE INTO ตารางที่มีอยู่ตารางขนาดใหญ่ ประมวลผลเฉพาะแถวใหม่
ephemeral(ไม่มีออบเจ็กต์ — ฝังเป็น CTE)ตรรกะระดับกลางที่ใช้ร่วมกันโดยโมเดลปลายทางตัวเดียว

โมเดล incremental สองตัวแสดง incremental strategy ที่ต่างกัน:

fact_trips — ประมวลผลทริปใหม่ตั้งแต่รอบล่าสุด:

{% if is_incremental() %}
  WHERE pickup_at > (SELECT MAX(pickup_at) FROM {{ this }})
{% endif %}

agg_hourly_zone_trips — aggregate ใหม่ในหน้าต่างเลื่อน 2 ชั่วโมง เพื่อรับข้อมูลที่มาถึงช้า:

{% if is_incremental() %}
  WHERE pickup_at >= DATEADD('hour', -2, CURRENT_TIMESTAMP())
{% endif %}

ในรอบแรกสุด (ตารางว่าง) is_incremental() คืนค่า false และชุดข้อมูลทั้งหมดจะถูกประมวลผล ในรอบต่อ ๆ ไปจะประมวลผลเฉพาะข้อมูลใหม่ ถ้า schema เปลี่ยนและคุณต้องสร้างใหม่ตั้งแต่ต้น ให้รัน:

dbt run --full-refresh

กลยุทธ์ MERGE (ความท้าทายสำคัญของการย้ายระบบ)

เมื่อ incremental_strategy = 'merge' dbt จะสร้างคำสั่ง MERGE INTO ของ Snowflake:

MERGE INTO ANALYTICS.FACT_TRIPS AS target
USING (SELECT ...) AS source
ON target.trip_id = source.trip_id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT ...;

นี่คือหนึ่งในความท้าทายของการย้ายระบบที่สำคัญที่สุดที่บันทึกไว้ใน lab นี้ ClickHouse ไม่มีคำสั่ง MERGE สิ่งที่เทียบเท่าใน ClickHouse คือการใช้ table engine ReplacingMergeTree และเพิ่ม FINAL ในการ query หรือใช้ CollapsingMergeTree เพื่อความหมายแบบ insert/delete ที่ระบุชัดเจน


การตั้งชื่อ Schema: มาโคร generate_schema_name

พฤติกรรมเริ่มต้นของ dbt คือต่อ target schema จาก profiles.yml เข้ากับ custom schema ใน dbt_project.yml:

target schema = STAGING  +  custom schema = ANALYTICS  →  STAGING_ANALYTICS  (wrong)

โปรเจกต์นี้เขียนทับพฤติกรรมนั้นด้วยมาโครที่กำหนดเองใน macros/generate_schema_name.sql:

{% macro generate_schema_name(custom_schema_name, node) -%}
  {%- if custom_schema_name is none -%}
    {{ target.schema | upper }}      -- no custom schema → use target schema
  {%- else -%}
    {{ custom_schema_name | upper }}  -- custom schema → use it directly
  {%- endif -%}
{%- endmacro %}

ผลลัพธ์: โมเดลที่มี +schema: ANALYTICS จะลงใน ANALYTICS ไม่ใช่ STAGING_ANALYTICS

มาโครนี้จำเป็นทุกครั้งที่คุณมีหลาย schema ในโปรเจกต์ dbt เดียว และไม่ต้องการให้ชื่อ target schema ถูกเติมไว้ข้างหน้า


การเชื่อมต่อและ Credentials (profiles.yml)

dbt เชื่อมต่อกับ Snowflake ผ่าน profile ที่กำหนดไว้ใน ~/.dbt/profiles.yml (ไม่เคย commit เข้า git) ชื่อ profile ใน dbt_project.yml ต้องตรงกัน:

# dbt_project.yml
profile: 'nyc_taxi'

# ~/.dbt/profiles.yml
nyc_taxi:
  target: dev
  outputs:
    dev:
      type: snowflake
      account: "{{ env_var('SNOWFLAKE_ORG') }}-{{ env_var('SNOWFLAKE_ACCOUNT') }}"
      role: DBT_ROLE
      database: NYC_TAXI_DB
      warehouse: TRANSFORM_WH
      schema: STAGING       # ← this is the "target schema" / default schema
      threads: 4

ประเด็นสำคัญ:

  • schema: STAGING คือ schema เริ่มต้น โมเดลที่ไม่มีการเขียนทับด้วย +schema: จะลงที่นี่
  • role: DBT_ROLE เป็น role แบบสิทธิ์น้อยที่สุดที่ Terraform สร้างขึ้น โดยมีเฉพาะสิทธิ์ที่ dbt ต้องใช้
  • threads: 4 ควบคุมว่า dbt จะ build โมเดลพร้อมกันได้กี่ตัว
  • Credentials มาจากตัวแปรสภาพแวดล้อม ซึ่งโหลดจาก .env ก่อนรัน setup

การทดสอบ

การทดสอบของ dbt มีสองรูปแบบ:

Schema tests (ประกาศใน schema.yml)

- name: trip_id
  tests:
    - not_null
    - unique
- name: total_amount_usd
  tests:
    - dbt_expectations.expect_column_values_to_be_between:
        min_value: 0
        max_value: 1000

not_null และ unique เป็นแบบ built-in การทดสอบของ dbt_expectations มาจากแพ็กเกจ calogica/dbt_expectations ที่ประกาศไว้ใน packages.yml

การทดสอบ SQL แบบกำหนดเอง (tests/)

-- tests/assert_revenue_positive.sql
-- A passing test returns 0 rows
SELECT trip_id, total_amount_usd
FROM {{ ref('fact_trips') }}
WHERE total_amount_usd < 0

การทดสอบแบบกำหนดเองก็เป็นเพียง query SQL dbt จะรันมันและ ล้มเหลวถ้ามีแถวใดถูกคืนกลับมา

รันการทดสอบทั้งหมดด้วย:

dbt test

แพ็กเกจจากภายนอก (packages.yml)

packages:
  - package: dbt-labs/dbt_utils
    version: [">=1.0.0", "<2.0.0"]
  - package: calogica/dbt_expectations
    version: [">=0.10.0", "<1.0.0"]

ติดตั้งก่อนใช้งานครั้งแรก:

dbt deps

dbt_utils ให้ตัวสร้าง date_spine ที่ใช้ใน dim_date.sql ส่วน dbt_expectations ให้การทดสอบช่วงค่า/การกระจายตัวที่เกินกว่า not_null/unique แบบ built-in


กราฟ Dependency

dbt build โมเดลตามลำดับที่ถูกต้องโดยอัตโนมัติ ด้วยการไล่ตามการเรียก {{ ref() }}:

RAW.TRIPS_RAW (source — not managed by dbt)
    └── stg_trips (view)
            └── int_trips_enriched (ephemeral)
                    ├── fact_trips (incremental table)
                    └── agg_hourly_zone_trips (incremental table)

ANALYTICS.DIM_TAXI_ZONES (seeded by SQL script)
    └── stg_taxi_zones (view)
            ├── int_trips_enriched
            └── dim_taxi_zones (table)

dbt_utils.date_spine
    └── dim_date (table)

{{ ref('stg_trips') }} คือวิธีที่โมเดลหนึ่งประกาศ dependency บนอีกโมเดลหนึ่ง ส่วน {{ source('raw', 'TRIPS_RAW') }} ประกาศ dependency บนตารางภายนอก (กำหนดไว้ใน sources.yml)


คำสั่งที่ใช้บ่อย

คำสั่งทำอะไร
dbt depsติดตั้งแพ็กเกจจาก packages.yml
dbt runbuild โมเดลทั้งหมด (แบบ incremental ที่ทำได้)
dbt run --full-refreshสร้างโมเดล incremental ทั้งหมดใหม่ตั้งแต่ต้น
dbt run -s fact_tripsbuild เฉพาะ fact_trips และ dependency ของมัน
dbt testรันการทดสอบทั้ง schema และแบบกำหนดเองทั้งหมด
dbt builddbt run + dbt test รวมกัน
dbt compileสร้าง SQL โดยไม่รัน (มีประโยชน์ตอน debug)
dbt docs generate && dbt docs servebuild และดูกราฟ lineage ในเบราว์เซอร์

ในโปรเจกต์นี้ setup.sh จะสั่ง dbt run --full-refresh ให้อัตโนมัติถ้า fact_trips ว่าง (รอบแรก หรือหลังการรื้อระบบ)


dbt อยู่ตรงไหนในการ setup ทั้งหมด

terraform apply          → creates warehouses, database, schemas, roles
scripts/01_create_tables.sql → creates raw tables, seeds dimension data
scripts/02_seed_data.sql → loads 50M synthetic trip rows
dbt deps && dbt build    → transforms raw data into analytics-ready tables
scripts/03_create_streams_tasks.sql → creates CDC stream and scheduled task

dbt อยู่กลาง pipeline มันไม่สามารถรันได้จนกว่าตารางดิบจะมีอยู่และมีข้อมูล สคริปต์ setup.sh จัดการลำดับนี้ให้

ในหน้านี้

TH