数据工程(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 核心对比
| 对比维度 | ETL | ELT |
|---|---|---|
| 处理顺序 | 先抽取→外部转换→加载 | 先抽取加载→仓库内转换 |
| 转换引擎 | 外部 ETL 服务器 / Spark | 云数据仓库内置 SQL 引擎 |
| 硬件需求 | 独立 ETL 集群,维护成本高 | 仓库弹性计算,无额外集群 |
| 灵活性 | 修改逻辑重跑全量,迭代慢 | SQL + dbt 快速迭代,支持增量 |
| 数据保留 | 原始数据可能丢失 | 原始数据完整保留 |
| 典型工具 | Informatica, Talend, DataStage | dbt + 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 | 读时 Schema | Schema 演进 + 约束 |
| 事务支持 | 完整 ACID | 无原生事务 | ACID + 时间旅行 + 回滚 |
| 查询性能 | 极高 | 较低(需扫描大量文件) | 接近数仓(布局优化 + 缓存) |
| 成本模型 | 存储计算绑定,价格较高 | 对象存储极低成本 | 低成本存储 + 弹性计算 |
| 代表产品 | Snowflake, BigQuery, Redshift | S3 + EMR, ADLS | Databricks, 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 Airflow | Prefect | Dagster |
|---|---|---|---|
| 架构 | Scheduler + Worker | 混合云原生,2.x 去中心化 | 数据感知型编排 |
| 核心抽象 | DAG + Operator | Flow + Task | Job + 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: 数据质量工具应该在哪个环节介入?
推荐"多层防线":
- 入库前:Fivetran/Airbyte 做 Schema 校验;
- 转换中:dbt 测试做列级校验;
- 产出后:GE/Soda 做行级与分布级校验;
- 消费端:反向 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 转换,建立
staging→intermediate→marts三级目录。 - 编排层:粒度控制在"一次原子业务目标",单个 DAG 内任务不超过 30 个。
- 质量层:关键业务表必须配置
not_null+unique+ 行数波动检测,质量失败阻断下游。 - 目录层:dbt docs、GE Data Docs 与 DataHub 统一映射,形成"代码文档即真实文档"的文化。
- 安全层:所有凭据通过环境变量或 Vault 注入,禁止硬编码。
- CI/CD:PR 阶段必须执行
dbt compile、sqlfluff lint与dbt test --select state:modified+。
数据工程的终极目标是让数据消费者在正确的时间,以可信赖的方式获取易洞察的数据。Modern Data Stack 以 ELT 范式降低转换门槛,以声明式工具提升协作效率,以 DataOps 文化保障交付质量——这正是云原生时代数据团队的核心竞争力所在。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。