数据工程深度指南:Modern Data Stack 全栈实践

深入探索 Modern Data Stack 的每个组件:ELT/ETL 架构选型、dbt 数据转换、数据仓库/湖/湖仓一体、Airflow/Prefect/Dagster 编排、Great Expectations/Soda 数据质量、DataHub/Amundsen 数据目录与治理。附完整代码示例与对比表。

数据工程(Data Engineering)是连接原始数据与数据驱动决策的核心桥梁。Modern Data Stack(现代数据技术栈)以 ELT 为核心范式,以云数据仓库/湖仓为存储底座,以声明式工具替代手写代码,以 DataOps 重塑数据交付流程。本文从架构演进、核心组件、代码实践与数据治理四个维度,构建一份可落地的全栈深度指南。


1. Modern Data Stack 全景

Modern Data Stack 的典型分层如下:

┌─────────────────────────────────────────────────────────────┐
│                     消费层 (BI / Notebook)                    │
├─────────────────────────────────────────────────────────────┤
│               目录与治理 (DataHub / Amundsen)                 │
├─────────────────────────────────────────────────────────────┤
│         质量与可观测 (GE / Soda / Monte Carlo)                │
├─────────────────────────────────────────────────────────────┤
│              编排层 (Airflow / Prefect / Dagster)             │
├─────────────────────────────────────────────────────────────┤
│              转换层 (dbt / SQL + Jinja)                       │
├─────────────────────────────────────────────────────────────┤
│        存储与计算 (Snowflake / BigQuery / Delta Lake)         │
├─────────────────────────────────────────────────────────────┤
│         摄取层 (Fivetran / Airbyte / Kafka Connect)           │
└─────────────────────────────────────────────────────────────┘

企业架构通常经历三个阶段:传统数仓(Oracle + Informatica)→ 大数据平台(Hadoop + Hive/Spark)→ 云原生 Modern Data Stack(对象存储 + 弹性计算 + 声明式转换)。第三阶段以存算分离、按量付费与敏捷交付为核心优势。


2. ETL vs ELT:范式迁移

2.1 核心对比

对比维度ETLELT
处理顺序先抽取→外部转换→加载先抽取加载→仓库内转换
转换引擎外部 ETL 服务器 / Spark云数据仓库内置 SQL 引擎
硬件需求独立 ETL 集群,维护成本高仓库弹性计算,无额外集群
灵活性修改逻辑重跑全量,迭代慢SQL + dbt 快速迭代,支持增量
数据保留原始数据可能丢失原始数据完整保留
典型工具Informatica, Talend, DataStagedbt + Fivetran + Snowflake/BigQuery
适用场景遗留系统、严格合规清洗云原生数仓、数据民主化
成本模型前期投入高,需预留资源按需付费,存算分离

2.2 为什么 ELT 成为主流

云数据仓库的存算分离架构,使仓库内部 SQL 转换成本远低于维护独立 ETL 集群。dbt 的出现进一步放大优势:数据分析师可用纯 SQL 写出可测试、可版本控制的管道,无需 Python/Java。

-- ELT 典型管道:Fivetran 自动加载 → dbt 在仓库内转换
WITH source AS (
    SELECT * FROM raw_salesforce.opportunities
),
renamed AS (
    SELECT
        id AS opportunity_id,
        account_id,
        amount::DECIMAL(18,2) AS amount,
        close_date::DATE AS close_date,
        is_deleted = 'true' AS is_deleted
    FROM source
)
SELECT * FROM renamed WHERE NOT is_deleted

3. 数据仓库、数据湖与数据湖仓一体

3.1 核心对比

对比维度数据仓库数据湖数据湖仓一体
存储格式专有列式存储开放格式(Parquet/JSON)开放格式 + 事务层(Delta/Iceberg/Hudi)
数据类型结构化为主全类型支持全类型支持
Schema 管理写时严格 Schema读时 SchemaSchema 演进 + 约束
事务支持完整 ACID无原生事务ACID + 时间旅行 + 回滚
查询性能极高较低(需扫描大量文件)接近数仓(布局优化 + 缓存)
成本模型存储计算绑定,价格较高对象存储极低成本低成本存储 + 弹性计算
代表产品Snowflake, BigQuery, RedshiftS3 + EMR, ADLSDatabricks, Starburst, Iceberg
典型场景BI 报表、企业指标数据科学、原始日志、ML统一 BI + AI,消除数据孤岛

3.2 Delta Lake 实践

from delta import configure_spark_with_delta_pip
from pyspark.sql import SparkSession

builder = SparkSession.builder \
    .appName("LakehouseDemo") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")

spark = configure_spark_with_delta_pip(builder).getOrCreate()

# 写入(自动 ACID 事务)
df = spark.range(0, 1000).toDF("id")
df.write.format("delta").mode("overwrite").save("/mnt/delta/users")

# 时间旅行查询
spark.read.format("delta").option("versionAsOf", 0).load("/mnt/delta/users").show(5)

4. 声明式数据转换:dbt 实践

dbt 是 Modern Data Stack 转换层的事实标准,将 SQL 转化为可重用、可测试、可版本控制的数据模型。

4.1 项目配置

# dbt_project.yml
name: 'ecommerce_analytics'
version: '1.0.0'
config-version: 2
profile: 'snowflake_prod'
model-paths: ["models"]
test-paths: ["tests"]
macro-paths: ["macros"]
models:
  ecommerce_analytics:
    staging:
      +materialized: view
      +schema: staging
    marts:
      +materialized: table
      +schema: marts
      core:
        +tags: ["daily", "critical"]
# profiles.yml
snowflake_prod:
  target: dev
  outputs:
    dev:
      type: snowflake
      account: xy12345.us-east-1
      user: DBT_USER
      password: "{{ env_var('DBT_SNOWFLAKE_PASSWORD') }}"
      role: TRANSFORMER
      database: ANALYTICS
      warehouse: DBT_WH
      schema: staging
      threads: 8

4.2 完整模型示例

-- models/staging/stg_orders.sql
WITH source AS (SELECT * FROM raw_ecommerce.orders),
cleaned AS (
    SELECT
        order_id, customer_id, order_status, order_date,
        amount::NUMERIC(18, 2) AS order_amount, currency,
        created_at, updated_at,
        ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY updated_at DESC) AS dedup_rank
    FROM source WHERE order_id IS NOT NULL
)
SELECT order_id, customer_id, order_status, order_date, order_amount, currency, created_at, updated_at
FROM cleaned WHERE dedup_rank = 1
-- models/marts/core/fct_orders.sql
WITH orders AS (SELECT * FROM {{ ref('stg_orders') }}),
payments AS (SELECT * FROM {{ ref('stg_payments') }}),
order_payments AS (
    SELECT order_id, SUM(payment_amount) AS total_payment_amount, COUNT(*) AS payment_count
    FROM payments GROUP BY 1
)
SELECT
    orders.order_id, orders.customer_id, orders.order_date, orders.order_amount, orders.order_status,
    COALESCE(order_payments.total_payment_amount, 0) AS payment_amount,
    orders.order_amount - COALESCE(order_payments.total_payment_amount, 0) AS amount_difference
FROM orders LEFT JOIN order_payments USING (order_id)

4.3 dbt 测试与文档

# models/staging/schema.yml
version: 2
models:
  - name: stg_orders
    description: "清洗后的订单基础表,已去重"
    columns:
      - name: order_id
        description: "主键,唯一标识一笔订单"
        tests:
          - unique
          - not_null
      - name: customer_id
        tests:
          - not_null
          - relationships:
              to: ref('stg_customers')
              field: customer_id
      - name: order_amount
        tests:
          - not_null
          - dbt_utils.expression_is_true:
              expression: ">= 0"
      - name: order_status
        tests:
          - accepted_values:
              values: ['placed', 'shipped', 'completed', 'returned', 'cancelled']

4.4 增量模型与 Python 调用

-- models/marts/core/fct_orders_incremental.sql
{{ config(materialized='incremental', unique_key='order_id', on_schema_change='append_new_columns') }}

WITH new_orders AS (
    SELECT * FROM {{ ref('stg_orders') }}
    {% if is_incremental() %}
        WHERE updated_at > (SELECT MAX(updated_at) FROM {{ this }})
    {% endif %}
)
SELECT * FROM new_orders
import subprocess
import os
os.environ["DBT_SNOWFLAKE_PASSWORD"] = "your_secure_password"

subprocess.run(["dbt", "run", "--select", "tag:daily", "--profiles-dir", "."], check=True)
subprocess.run(["dbt", "test", "--select", "stg_orders"], check=True)

5. 数据管道编排:Airflow、Prefect 与 Dagster

5.1 三框架速览

维度Apache AirflowPrefectDagster
架构Scheduler + Worker混合云原生,2.x 去中心化数据感知型编排
核心抽象DAG + OperatorFlow + TaskJob + Op + Graph
动态工作流TaskFlow API / AIP-42原生支持Software-Defined Assets
本地开发较重重flow.run() 轻量dagster dev 体验优秀
数据血缘依赖插件内置部分原生一流
测试支持中等良好极强
社区生态最大增长迅速工程师口碑高

5.2 Airflow 完整生产级 DAG

以下 DAG 实现 S3→Snowflake 加载、dbt 转换与 GE 质量校验的完整链路。

# dags/prod_ecommerce_pipeline.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
from airflow.utils.task_group import TaskGroup
from datetime import datetime, timedelta
import subprocess

default_args = {
    'owner': 'data-engineering',
    'depends_on_past': False,
    'email': ['data-alerts@company.com'],
    'email_on_failure': True,
    'retries': 2,
    'retry_delay': timedelta(minutes=5),
}

with DAG(
    dag_id='prod_ecommerce_daily_etl',
    default_args=default_args,
    description='每日电商数据管道:S3 → Snowflake → dbt → 质量校验',
    schedule_interval='0 6 * * *',
    start_date=datetime(2026, 1, 1),
    catchup=False,
    tags=['ecommerce', 'production', 'dbt'],
    max_active_runs=1,
) as dag:

    with TaskGroup(group_id='ingestion') as ingestion:
        load_orders = SnowflakeOperator(
            task_id='load_orders_from_s3',
            snowflake_conn_id='snowflake_default',
            sql="""
                COPY INTO raw_ecommerce.orders
                FROM @s3_stage/orders/{{ ds_nodash }}/
                FILE_FORMAT = (TYPE = 'CSV', SKIP_HEADER = 1)
                ON_ERROR = 'SKIP_FILE_3';
            """
        )
        load_payments = SnowflakeOperator(
            task_id='load_payments_from_s3',
            snowflake_conn_id='snowflake_default',
            sql="""
                COPY INTO raw_ecommerce.payments
                FROM @s3_stage/payments/{{ ds_nodash }}/
                FILE_FORMAT = (TYPE = 'CSV', SKIP_HEADER = 1);
            """
        )

    def run_dbt_models(**context):
        cmd = ['dbt', 'run', '--select', 'tag:daily',
               '--vars', f'{{"execution_date": "{context["ds"]}"}}',
               '--profiles-dir', '/opt/airflow/dbt']
        subprocess.run(cmd, check=True, cwd='/opt/airflow/dbt/ecommerce_analytics')

    dbt_run = PythonOperator(task_id='dbt_run_daily', python_callable=run_dbt_models, provide_context=True)

    def run_dbt_tests(**context):
        cmd = ['dbt', 'test', '--select', 'tag:daily', '--profiles-dir', '/opt/airflow/dbt']
        subprocess.run(cmd, check=True, cwd='/opt/airflow/dbt/ecommerce_analytics')

    dbt_test = PythonOperator(task_id='dbt_test_daily', python_callable=run_dbt_tests, provide_context=True)

    def run_ge_checkpoint(**context):
        import great_expectations as gx
        ctx = gx.get_context(context_root_dir='/opt/airflow/great_expectations')
        result = ctx.run_checkpoint(checkpoint_name="fct_orders_checkpoint")
        if not result.success:
            raise ValueError("GE checkpoint failed!")

    ge_validation = PythonOperator(task_id='validate_fct_orders', python_callable=run_ge_checkpoint)
    ingestion >> dbt_run >> dbt_test >> ge_validation

5.3 Prefect 3.x 示例

# prefect/etl_flow.py
from prefect import flow, task
import pandas as pd
from sqlalchemy import create_engine

@task
def extract_from_api(url: str) -> pd.DataFrame:
    return pd.read_json(url)

@task
def transform_orders(df: pd.DataFrame) -> pd.DataFrame:
    df["order_amount"] = df["order_amount_cents"] / 100
    df["order_date"] = pd.to_datetime(df["order_date"])
    return df.dropna(subset=["order_id"])

@task
def load_to_warehouse(df: pd.DataFrame, table_name: str):
    engine = create_engine("snowflake://user:pass@account/db/schema")
    df.to_sql(table_name, engine, if_exists="append", index=False)

@flow(name="ecommerce_daily_sync", log_prints=True)
def ecommerce_sync(url: str, table_name: str = "raw_orders"):
    raw_df = extract_from_api(url)
    transformed_df = transform_orders(raw_df)
    load_to_warehouse(transformed_df, table_name)

5.4 Dagster 数据资产优先

# dagster/assets.py
from dagster import asset, AssetIn, Definitions, ScheduleDefinition, define_asset_job
import pandas as pd

@asset(key_prefix=["raw"], io_manager_key="snowflake_io_manager")
def raw_orders() -> pd.DataFrame:
    return pd.read_csv("s3://bucket/orders.csv")

@asset(ins={"raw_orders": AssetIn(key_prefix=["raw"])}, io_manager_key="snowflake_io_manager")
def stg_orders(raw_orders: pd.DataFrame) -> pd.DataFrame:
    df = raw_orders.copy()
    df["order_amount"] = df["amount_cents"] / 100
    return df

@asset(ins={"stg_orders": AssetIn(key_prefix=["raw"])}, io_manager_key="snowflake_io_manager")
def daily_order_metrics(stg_orders: pd.DataFrame) -> pd.DataFrame:
    return stg_orders.groupby("order_date").agg(
        total_orders=("order_id", "count"),
        total_revenue=("order_amount", "sum")
    ).reset_index()

daily_job = define_asset_job("daily_job", selection="*")
daily_schedule = ScheduleDefinition(job=daily_job, cron_schedule="0 6 * * *")

defs = Definitions(
    assets=[raw_orders, stg_orders, daily_order_metrics],
    jobs=[daily_job],
    schedules=[daily_schedule]
)

6. 数据质量与可观测性

6.1 Great Expectations 完整测试套件

# great_expectations/suite_definition.py
import great_expectations as gx
from great_expectations.core.expectation_suite import ExpectationSuite
from great_expectations.expectations import (
    ExpectColumnValuesToNotBeNull, ExpectColumnValuesToBeBetween,
    ExpectColumnValuesToBeUnique, ExpectTableRowCountToBeBetween,
    ExpectColumnValuesToBeInSet, ExpectColumnValuesToMatchRegex
)

context = gx.get_context()
suite_name = "fct_orders_suite"
try:
    suite = context.suites.add(ExpectationSuite(name=suite_name))
except Exception:
    suite = context.suites.get(name=suite_name)

suite.add_expectation(ExpectTableRowCountToBeBetween(min_value=1000, max_value=10000000))
suite.add_expectation(ExpectColumnValuesToNotBeNull(column="order_id"))
suite.add_expectation(ExpectColumnValuesToBeUnique(column="order_id"))
suite.add_expectation(ExpectColumnValuesToNotBeNull(column="order_amount"))
suite.add_expectation(ExpectColumnValuesToBeBetween(column="order_amount", min_value=0, max_value=1000000))
suite.add_expectation(ExpectColumnValuesToBeInSet(
    column="order_status", value_set=["placed", "shipped", "completed", "returned", "cancelled"]
))
suite.add_expectation(ExpectColumnValuesToMatchRegex(column="customer_id", regex=r"^CUST-[0-9]{6}$"))
suite.save()
# 运行 Checkpoint
checkpoint = context.add_or_update_checkpoint(
    name="fct_orders_checkpoint",
    validations=[{
        "batch_request": {
            "datasource_name": "snowflake_datasource",
            "data_asset_name": "fct_orders"
        },
        "expectation_suite_name": "fct_orders_suite"
    }],
    action_list=[
        {"name": "store_validation_result", "action": {"class_name": "StoreValidationResultAction"}},
        {"name": "update_data_docs", "action": {"class_name": "UpdateDataDocsAction"}},
        {"name": "send_slack_notification", "action": {
            "class_name": "SlackNotificationAction",
            "slack_webhook": "${SLACK_WEBHOOK_URL}",
            "notify_on": "failure"
        }}
    ]
)
result = checkpoint.run()
if not result.success:
    raise RuntimeError("数据质量校验未通过")

6.2 Soda:声明式质量检查

# checks/fct_orders_checks.yml
checks for fct_orders:
  - row_count > 1000
  - duplicate_count(order_id) = 0
  - missing_count(order_id) = 0
  - invalid_count(order_status) = 0:
      valid values: [placed, shipped, completed, returned, cancelled]
  - min(order_amount) >= 0
  - max(order_amount) < 1000000
  - freshness(create_date) < 1d
from soda.scan import Scan
scan = Scan()
scan.set_data_source_name("snowflake")
scan.add_configuration_yaml_file("config.yml")
scan.add_sodacl_yaml_file("checks/fct_orders_checks.yml")
scan.set_scan_definition_name("daily_quality_scan")
scan.execute()
if scan.has_failures():
    raise ValueError("Soda 检测到数据质量问题")

7. 数据目录与治理

数据目录解决三大问题:数据发现(找数)、血缘理解(懂数)、信任建立(信数)。

7.1 DataHub 摄取配置

# datahub/snowflake_recipe.yml
source:
  type: snowflake
  config:
    account_id: xy12345.us-east-1
    username: datahub_ingest
    password: "${SNOWFLAKE_PASSWORD}"
    warehouse: DATAHUB_WH
    database: ANALYTICS
    schema_pattern:
      allow: ["marts", "staging"]
    profiling:
      enabled: true
sink:
  type: datahub-rest
  config:
    server: "http://datahub-datahub-gms:8080"
    token: "${DATAHUB_ACCESS_TOKEN}"
datahub ingest -c datahub/snowflake_recipe.yml
datahub ingest list-runs

7.2 DataHub REST API 标注资产

from datahub.emitter.rest_emitter import DatahubRestEmitter
from datahub.metadata.schema_classes import TagAssociationClass, GlobalTagsClass

emitter = DatahubRestEmitter("http://datahub-datahub-gms:8080", token="your-token")
dataset_urn = "urn:li:dataset:(urn:li:dataPlatform:snowflake,analytics.marts.fct_orders,PROD)"

tags = GlobalTagsClass(tags=[TagAssociationClass(tag="urn:li:tag:pii")])
emitter.emit_mcp({
    "entityUrn": dataset_urn,
    "aspectName": "globalTags",
    "aspect": tags,
    "changeType": "UPSERT"
})

7.3 工具选型建议

  • DataHub:功能最全面,血缘精度高,企业级特性(策略管理、标签体系),社区最活跃。
  • Amundsen:由 Lyft 开源,UI 轻量,适合已有复杂搜索需求但血缘要求不高的团队。

两者均支持通过 dbt 元数据自动生成表文档与列描述。


8. DataOps 与 CI/CD

8.1 dbt + GitHub Actions CI

# .github/workflows/dbt_ci.yml
name: dbt CI
on:
  pull_request:
    branches: [main]
    paths: ['models/**', 'tests/**', 'dbt_project.yml']
jobs:
  dbt-ci:
    runs-on: ubuntu-latest
    env:
      DBT_SNOWFLAKE_PASSWORD: ${{ secrets.DBT_SNOWFLAKE_PASSWORD }}
    steps:
      - uses: actions/checkout@v4
      - uses: actions/setup-python@v5
        with: { python-version: '3.11' }
      - run: pip install dbt-snowflake dbt-utils sqlfluff
      - run: sqlfluff lint models/ --dialect snowflake
      - run: dbt deps --profiles-dir ./ci_profiles
      - run: dbt compile --profiles-dir ./ci_profiles --target ci
      - run: dbt test --select state:modified+ --defer --state target --profiles-dir ./ci_profiles --target ci

8.2 Airflow DAG 自动化测试

# tests/test_dag_integrity.py
import pytest
from airflow.models import DagBag

@pytest.fixture
def dag_bag():
    return DagBag(dag_folder="dags", include_examples=False)

def test_no_import_errors(dag_bag):
    assert len(dag_bag.import_errors) == 0

def test_dag_has_valid_schedule(dag_bag):
    for dag_id, dag in dag_bag.dags.items():
        assert dag.schedule_interval is not None

def test_dag_retry_config(dag_bag):
    for dag_id, dag in dag_bag.dags.items():
        assert dag.default_args.get("retries", 0) > 0

9. 常见问题解答(FAQ)

Q1: ELT 是否意味着完全不需要数据清洗?

不是。ELT 将主要转换移至仓库内执行,但轻度的格式校验、去重与类型转换仍需在加载阶段完成。实践中通常会保留一个轻量"置备层"(Landing / Staging),dbt 的 staging 模型正是承担这一角色。

Q2: dbt 能处理非 SQL 的数据转换(如 Python ML 特征工程)吗?

dbt Core 原生支持 SQL + Jinja。dbt 1.3 起在 Snowflake/BigQuery/Databricks 上支持 Python models,但复杂 ML 工程更推荐在编排层(Airflow/Dagster)中串联 Python 任务与 dbt 模型,各司其职。

Q3: 数据质量工具应该在哪个环节介入?

推荐"多层防线":

  1. 入库前:Fivetran/Airbyte 做 Schema 校验;
  2. 转换中:dbt 测试做列级校验;
  3. 产出后:GE/Soda 做行级与分布级校验;
  4. 消费端:反向 ETL 校验输出一致性。越早发现问题,修复成本越低。

Q4: 小企业是否也需要数据目录?

20 张表以下、团队小于 5 人时,dbt docs 静态站点已足够。当表超 50 张、血缘复杂或需跨团队协作与合规审计时,正式的数据目录投资回报率显著提升。

Q5: Airflow、Prefect、Dagster 该如何选择?

  • Airflow:已有大量 Hadoop/Spark Operator 需求,社区生态最丰富。
  • Prefect:更现代化的 Python 体验,原生支持动态工作流,减少样板代码。
  • Dagster:以数据资产为核心,数据血缘一流,本地调试体验极佳,适合 Asset-oriented 团队。

10. 总结与最佳实践 checklist

Modern Data Stack 是以云数据仓库/湖仓为底座、以 dbt 为转换中枢、以编排工具为动脉、以质量与目录为治理保障的协作体系。

  • 存储层:优先选择支持开放表格式(Delta Lake / Iceberg / Hudi)的湖仓方案,避免供应商锁定。
  • 转换层:统一使用 dbt 管理 SQL 转换,建立 stagingintermediatemarts 三级目录。
  • 编排层:粒度控制在"一次原子业务目标",单个 DAG 内任务不超过 30 个。
  • 质量层:关键业务表必须配置 not_null + unique + 行数波动检测,质量失败阻断下游。
  • 目录层:dbt docs、GE Data Docs 与 DataHub 统一映射,形成"代码文档即真实文档"的文化。
  • 安全层:所有凭据通过环境变量或 Vault 注入,禁止硬编码。
  • CI/CD:PR 阶段必须执行 dbt compilesqlfluff lintdbt test --select state:modified+

数据工程的终极目标是让数据消费者在正确的时间,以可信赖的方式获取易洞察的数据。Modern Data Stack 以 ELT 范式降低转换门槛,以声明式工具提升协作效率,以 DataOps 文化保障交付质量——这正是云原生时代数据团队的核心竞争力所在。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「data-engineering」更多文章

  1. 数据平台工程:Data Mesh、FinOps 与 DataOps 生产实践
  2. Kafka Connect CDC 实战:Debezium 数据同步与变更捕获
  3. 数据仓库建模深度指南:从 Kimball 到 Data Mesh