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_trips | RAW.TRIPS_RAW | เปลี่ยนชื่อคอลัมน์เป็น snake_case, เพิ่ม duration_minutes, แผ่คอลัมน์ VARIANT TRIP_METADATA ออกเป็นคอลัมน์ที่มีชนิดข้อมูลชัดเจน |
stg_taxi_zones | ANALYTICS.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,670 | Date spine 2009–2029 พร้อมไตรมาสทางบัญชีและวันหยุดราชการสหรัฐ |
dim_payment_type | 6 | ส่งผ่านมาจากข้อมูล seed |
dim_vendor | 3 | ส่งผ่านมาจากข้อมูล seed |
dim_taxi_zones | 265 | ส่งผ่านมาทาง stg_taxi_zones |
ตาราง fact/aggregate แบบ incremental — ขนาดใหญ่ อัปเดตด้วย MERGE ในแต่ละรอบ:
| โมเดล | จำนวนแถว | หมายเหตุ |
|---|---|---|
fact_trips | 50M | หนึ่งแถวต่อหนึ่งทริป denormalize เต็มที่ |
agg_hourly_zone_trips | ~9M | จำนวนรายชั่วโมงต่อโซนที่ aggregate ไว้ล่วงหน้า |
Materializations
materialization ควบคุมว่า dbt จะสร้างอะไรใน Snowflake สำหรับโมเดลหนึ่ง ๆ
| Materialization | ออบเจ็กต์ใน Snowflake | ใช้เมื่อไร |
|---|---|---|
view | CREATE VIEW | ราคาถูก สะท้อนข้อมูลล่าสุดตลอด ใช้กับ staging |
table | CREATE TABLE AS SELECT | สร้างใหม่ทั้งหมดทุกรอบ ใช้กับ dimension ขนาดเล็ก |
incremental | MERGE 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: 1000not_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 depsdbt_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 run | build โมเดลทั้งหมด (แบบ incremental ที่ทำได้) |
dbt run --full-refresh | สร้างโมเดล incremental ทั้งหมดใหม่ตั้งแต่ต้น |
dbt run -s fact_trips | build เฉพาะ fact_trips และ dependency ของมัน |
dbt test | รันการทดสอบทั้ง schema และแบบกำหนดเองทั้งหมด |
dbt build | dbt run + dbt test รวมกัน |
dbt compile | สร้าง SQL โดยไม่รัน (มีประโยชน์ตอน debug) |
dbt docs generate && dbt docs serve | build และดูกราฟ 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 taskdbt อยู่กลาง pipeline มันไม่สามารถรันได้จนกว่าตารางดิบจะมีอยู่และมีข้อมูล สคริปต์ setup.sh จัดการลำดับนี้ให้