Snowflake MigrationClickHouse Workshops

dbt trên Snowflake

Cách pipeline Medallion nguồn được xây dựng: sources, staging view, model incremental MERGE, snapshot và test.

Tài liệu này giải thích cách dbt (data build tool) được dùng trong NYC Taxi Snowflake Migration Lab — nó làm gì, tại sao từng phần tồn tại, và cách tư duy về nó.


dbt Làm Gì (và Không Làm Gì)

dbt biến đổi dữ liệu đã có trong database của bạn. Nó không nạp dữ liệu từ bên ngoài, không di chuyển file, và không quản lý hạ tầng. Việc của nó là lấy các bảng thô và biến chúng thành các bảng sạch, đã được kiểm thử, sẵn sàng cho phân tích — bằng cách chạy SQL mà bạn viết.

Hãy hình dung nó như một build system cho SQL. Mỗi file .sql trong models/ là một model trở thành một bảng hoặc một view trong Snowflake. dbt xử lý phần boilerplate CREATE OR REPLACE, giải quyết các phụ thuộc giữa các model, và chạy các test của bạn.


Bố Cục Project

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

Ba Tầng (Kiến Trúc Medallion)

Tầng 1 — Staging (models/staging/)

Mục đích: Lấy dữ liệu thô đúng như lúc nó vừa tới và làm cho nó dùng được.

Các model này nằm trong schema STAGING dưới dạng view (không tốn chi phí lưu trữ — chúng chạy tại thời điểm truy vấn). Mỗi staging model làm một việc:

ModelNguồnNó làm gì
stg_tripsRAW.TRIPS_RAWĐổi tên cột sang snake_case, thêm duration_minutes, làm phẳng cột VARIANT TRIP_METADATA thành các cột có kiểu rõ ràng
stg_taxi_zonesANALYTICS.DIM_TAXI_ZONESDọn dẹp nhẹ, thêm các lớp bảo vệ COALESCE, cung cấp một node lineage cho dbt

Phần việc quan trọng nhất ở đây là làm phẳng cột VARIANT TRIP_METADATA. Cú pháp colon-path của Snowflake trích xuất các trường JSON lồng nhau:

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

Đây là một trong những thách thức di trú — ClickHouse dùng JSONExtractFloat(TRIP_METADATA, 'driver', 'rating') thay thế.

Tầng 2 — Intermediate (models/intermediate/)

Mục đích: Thực hiện toàn bộ các phép join ở một chỗ để không phải lặp lại chúng.

int_trips_enriched join stg_trips với mọi dimension (zone, payment type, vendor, date) và tạo ra một dòng rộng, đã denormalize hoàn toàn cho mỗi chuyến đi. Nó được khai báo là ephemeral, nghĩa là dbt nhúng inline SQL của nó vào bất kỳ model nào tham chiếu tới nó — không có bảng hay view vật lý nào được tạo trong Snowflake.

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

Dùng ephemeral khi kết quả trung gian chỉ cần cho một model hạ nguồn duy nhất và bạn không muốn trả chi phí lưu trữ hay overhead biên dịch truy vấn.

Tầng 3 — Analytics (models/analytics/)

Mục đích: Các bảng cuối cùng, sẵn sàng cho dashboard.

Chúng nằm trong schema ANALYTICS. Có hai loại:

Bảng dimension tĩnh — nhỏ, được nạp lại toàn bộ ở mỗi lần dbt run:

ModelSố dòngGhi chú
dim_date~7,670Date spine 2009–2029, có quý tài chính và các ngày lễ liên bang Mỹ
dim_payment_type6Chuyển tiếp trực tiếp từ dữ liệu seed
dim_vendor3Chuyển tiếp trực tiếp từ dữ liệu seed
dim_taxi_zones265Chuyển tiếp qua stg_taxi_zones

Bảng fact/aggregate incremental — lớn, được cập nhật bằng MERGE ở mỗi lần chạy:

ModelSố dòngGhi chú
fact_trips50MMột dòng cho mỗi chuyến đi, đã denormalize hoàn toàn
agg_hourly_zone_trips~9MSố lượt đếm theo giờ cho từng zone, đã tổng hợp trước

Materialization

Một materialization quyết định dbt tạo ra cái gì trong Snowflake cho một model nhất định.

MaterializationĐối tượng Snowflake nàoKhi nào dùng
viewCREATE VIEWRẻ; luôn phản ánh dữ liệu mới nhất; dùng cho staging
tableCREATE TABLE AS SELECTDựng lại toàn bộ mỗi lần chạy; dùng cho các dimension nhỏ
incrementalMERGE INTO bảng đã cóBảng lớn; chỉ xử lý các dòng mới
ephemeral(không có đối tượng — nhúng inline như CTE)Logic trung gian dùng chung bởi một model hạ nguồn

Hai model incremental minh họa hai chiến lược incremental khác nhau:

fact_trips — xử lý các chuyến đi mới kể từ lần chạy trước:

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

agg_hourly_zone_trips — tổng hợp lại một cửa sổ trượt 2 giờ để bắt được dữ liệu đến muộn:

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

Ở lần chạy đầu tiên (bảng trống), is_incremental() trả về false và toàn bộ tập dữ liệu được xử lý. Ở các lần chạy sau, chỉ dữ liệu mới được xử lý. Nếu schema thay đổi và bạn cần dựng lại từ đầu, hãy chạy:

dbt run --full-refresh

Chiến Lược MERGE (Thách Thức Di Trú Then Chốt)

Khi incremental_strategy = 'merge', dbt sinh ra một câu lệnh MERGE INTO của 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 ...;

Đây là một trong những thách thức di trú quan trọng nhất được ghi lại trong lab. ClickHouse không có câu lệnh MERGE. Tương đương trong ClickHouse là dùng table engine ReplacingMergeTree và thêm FINAL vào các truy vấn, hoặc dùng CollapsingMergeTree cho ngữ nghĩa insert/delete tường minh.


Đặt Tên Schema: Macro generate_schema_name

Hành vi mặc định của dbt là nối target schema từ profiles.yml với custom schema trong dbt_project.yml:

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

Project này ghi đè hành vi đó bằng một macro tùy chỉnh trong 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 %}

Kết quả: các model có +schema: ANALYTICS nằm trong ANALYTICS, không phải STAGING_ANALYTICS.

Macro này là bắt buộc bất cứ khi nào bạn có nhiều schema trong một dbt project và không muốn tên target schema bị thêm vào phía trước.


Kết Nối và Thông Tin Đăng Nhập (profiles.yml)

dbt kết nối tới Snowflake bằng một profile được định nghĩa trong ~/.dbt/profiles.yml (không bao giờ commit vào git). Tên profile trong dbt_project.yml phải trùng khớp:

# 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

Các điểm chính:

  • schema: STAGING là schema mặc định. Các model không có override +schema: sẽ nằm ở đây.
  • role: DBT_ROLE là một role đặc quyền tối thiểu do Terraform tạo, chỉ có các quyền mà dbt cần.
  • threads: 4 kiểm soát số model mà dbt dựng song song.
  • Thông tin đăng nhập đến từ biến môi trường, được nạp từ .env trước khi chạy setup.

Kiểm Thử

Test của dbt có hai dạng:

Schema test (khai báo trong 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 và unique là built-in. Các test dbt_expectations đến từ package calogica/dbt_expectations được khai báo trong packages.yml.

Test SQL tùy chỉnh (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

Test tùy chỉnh chỉ là các truy vấn SQL. dbt chạy chúng và thất bại nếu có bất kỳ dòng nào được trả về.

Chạy toàn bộ test với:

dbt test

Package Của Bên Thứ Ba (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"]

Cài chúng trước khi dùng lần đầu:

dbt deps

dbt_utils cung cấp bộ sinh date_spine được dùng trong dim_date.sql. dbt_expectations cung cấp các test về khoảng giá trị/phân phối, vượt xa các test built-in not_null/unique.


Đồ Thị Phụ Thuộc

dbt tự động dựng các model theo đúng thứ tự bằng cách đi theo các lệnh gọi {{ 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') }} là cách một model khai báo phụ thuộc vào một model khác. {{ source('raw', 'TRIPS_RAW') }} khai báo phụ thuộc vào một bảng bên ngoài (được định nghĩa trong sources.yml).


Các Lệnh Thường Dùng

LệnhNó làm gì
dbt depsCài các package từ packages.yml
dbt runDựng toàn bộ model (incremental khi có thể)
dbt run --full-refreshDựng lại toàn bộ model incremental từ đầu
dbt run -s fact_tripsChỉ dựng fact_trips và các phụ thuộc của nó
dbt testChạy toàn bộ schema test và test tùy chỉnh
dbt builddbt run + dbt test cùng lúc
dbt compileSinh SQL mà không thực thi (hữu ích khi debug)
dbt docs generate && dbt docs serveDựng và xem đồ thị lineage trong browser

Trong project này, dbt run --full-refresh được setup.sh kích hoạt tự động nếu fact_trips trống (lần chạy đầu hoặc sau khi tear-down).


dbt Nằm Ở Đâu Trong Toàn Bộ Quá Trình 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 nằm ở giữa pipeline. Nó không thể chạy cho tới khi các bảng thô tồn tại và có dữ liệu. Script setup.sh xử lý thứ tự này.

Trên trang này

VI