现代开源 ETL 生态

如果你在 2015 年问一个数据工程师用什么做 ETL,答案很可能是 Informatica 或 DataStage。到了 2026 年,答案已经完全不同。开源工具链的成熟让数据团队可以用更低的成本搭建媲美商业产品的数据管道。

现代开源 ETL 栈的核心是两个工具:Airbyte 负责 Extract 和 Load,dbt 负责 Transform。这个组合被称为"ELT 模式"——先把原始数据加载到目标仓库,再在仓库内部完成转换。相比传统 ETL,ELT 利用数据仓库的计算能力做转换,避免了大量数据传输。

Airbyte:数据集成层

什么是 Airbyte?

Airbyte 是一个开源的数据集成平台,提供从各种数据源到目标存储的连接器。它的核心理念是"连接器优先"——社区已经贡献了 300 多个预建连接器,覆盖数据库、API、文件存储、SaaS 应用等常见数据源。

核心概念

Source(源):数据从哪里来。可以是 PostgreSQL、MySQL、MongoDB 等数据库,也可以是 Stripe、Salesforce、Google Analytics 等 SaaS 服务,还可以是 S3、GCS 等文件存储。

Destination(目标):数据到哪里去。常见的目标包括 Snowflake、BigQuery、Redshift、PostgreSQL、S3 等。

Connection(连接):定义 Source 和 Destination 之间的数据流。一个连接包含同步频率、同步模式、流的选择等配置。

Sync Mode(同步模式):Airbyte 支持多种同步策略:

模式 说明 适用场景
Full Refresh - Replace 每次全量覆盖 小表、字典表
Full Refresh - Append 每次全量追加 不可变日志
Incremental - Append 增量追加新数据 事件表、日志表
Incremental - Deduped 增量追加并去重 订单表、用户表

部署 Airbyte

# 使用 Docker Compose 快速部署
mkdir airbyte && cd airbyte
wget https://raw.githubusercontent.com/airbytehq/airbyte/master/docker-compose.yaml
docker compose up -d

# 访问 http://localhost:8000 进入 Web UI
[/code]

### 配置数据源

通过 Web UI 或 API 配置数据源。下面是用 API 创建 PostgreSQL 数据源的示例:

创建 Source(PostgreSQL)

curl -X POST http://localhost:8000/api/v1/sources/create
-H “Content-Type: application/json”
-d ‘{ “name”: “My Postgres DB”, “sourceDefinitionId”: “decd338e-5647-4c0b-adf4-da0e75f5a750”, “connectionConfiguration”: { “host”: “localhost”, “port”: 5432, “database”: “ecommerce”, “username”: “airbyte”, “password”: “password”, “ssl”: false } }’ [/code]

配置目标

# 创建目标(PostgreSQL 数据仓库)
curl -X POST http://localhost:8000/api/v1/destinations/create \
  -H "Content-Type: application/json" \
  -d '{
    "name": "Warehouse",
    "destinationDefinitionId": "25c5221d-dc02-4b5c-9a37-fb7002d6e2f4",
    "connectionConfiguration": {
      "host": "warehouse-host",
      "port": 5432,
      "database": "analytics",
      "username": "dw_user",
      "password": "dw_password",
      "schema": "raw_data"
    }
  }'
[/code]

### 增量同步配置

增量同步是生产环境的关键功能。Airbyte 通过游标字段(cursor field)追踪上次同步的位置:

{ “syncMode”: “incremental”, “cursorField”: [“updated_at”], “destinationSyncMode”: “append_dedup” } [/code]

每次同步时,Airbyte 会记录 updated_at 的最大值,下次只拉取大于这个值的数据。这大大减少了数据传输量和源系统的压力。

Normalization

Airbyte 内置了基本的 normalization 功能,将 JSON 格式的原始数据展开为关系表。例如,从 API 抽取的嵌套 JSON 会被自动展开成多张表,并建立外键关系。

不过,Airbyte 的 normalization 比较基础。更复杂的转换工作应该交给 dbt。

dbt:数据转换层

什么是 dbt?

dbt(data build tool)是一个专注于数据转换的开源工具。它的核心理念是"用 SQL 定义转换逻辑,用 YAML 定义配置和测试"。dbt 不处理数据抽取和加载,它只做 T(Transform)的部分。

dbt 的工作方式很简单:你在 SQL 文件中用 SELECT 语句定义模型(model),dbt 自动将这些 SELECT 语句包装成 CREATE TABLE 或 CREATE VIEW,并在目标数据库中执行。

核心概念

Model(模型):一个 .sql 文件,包含一个 SELECT 语句。dbt 会将其编译为 CREATE TABLE AS 或 CREATE VIEW AS。

-- models/monthly_sales.sql
-- 每个模型对应一个 SELECT 语句
WITH orders AS (
    SELECT * FROM {{ ref('stg_orders') }}
),
monthly AS (
    SELECT
        DATE_TRUNC('month', order_date) AS month,
        SUM(amount) AS revenue,
        COUNT(*) AS order_count
    FROM orders
    GROUP BY 1
)
SELECT * FROM monthly
[/code]

**ref() 函数**:dbt 最强大的功能之一。ref('stg_orders') 自动解析模型之间的依赖关系,dbt 会根据依赖顺序执行模型。

**Source(源)**:定义 dbt 项目外部的数据来源(通常是 Airbyte 加载的原始数据)。

sources.yml

version: 2

sources:

  • name: airbyte_raw database: analytics schema: raw_data tables:
    • name: orders
    • name: customers
    • name: products [/code]

Test(测试):dbt 支持在数据上定义测试,确保数据质量。

# models/schema.yml
version: 2

models:
  - name: stg_orders
    columns:
      - name: order_id
        tests:
          - unique
          - not_null
      - name: amount
        tests:
          - not_null
          - dbt_utils.accepted_range:
              min_value: 0
[/code]

### 项目结构

一个典型的 dbt 项目结构如下:

dbt_project/ ├── dbt_project.yml # 项目配置 ├── models/ │ ├── staging/ # 原始数据层 │ │ ├── stg_orders.sql │ │ ├── stg_customers.sql │ │ └── sources.yml │ ├── intermediate/ # 中间处理层 │ │ └── int_order_items.sql │ └── marts/ # 业务层 │ ├── monthly_sales.sql │ └── customer_metrics.sql ├── tests/ # 自定义测试 │ └── assert_positive_revenue.sql ├── macros/ # Jinja 宏 │ └── date_utils.sql └── analyses/ # 分析查询 └── revenue_trends.sql [/code]

分层架构

dbt 推荐的分层架构通常包含三层:

Staging 层:直接映射源数据,做基本的类型转换和字段重命名。这一层和源数据是一一对应的。

-- models/staging/stg_orders.sql
SELECT
    id AS order_id,
    customer_id,
    total_amount AS amount,
    status,
    created_at AS order_date,
    updated_at
FROM {{ source('airbyte', 'orders') }}
[/code]

**Intermediate 层**:业务逻辑处理,多表关联,数据聚合。这一层开始产生业务价值。

– models/intermediate/int_order_details.sql SELECT o.order_id, o.order_date, o.amount, c.customer_name, c.customer_tier FROM {{ ref(‘stg_orders’) }} o LEFT JOIN {{ ref(‘stg_customers’) }} c ON o.customer_id = c.customer_id [/code]

Marts 层:面向业务分析的主题表,直接供 BI 工具使用。

-- models/marts/monthly_sales.sql
SELECT
    DATE_TRUNC('month', order_date) AS month,
    customer_tier,
    COUNT(DISTINCT order_id) AS order_count,
    SUM(amount) AS revenue
FROM {{ ref('int_order_details') }}
GROUP BY 1, 2
[/code]

### Jinja 模板

dbt 使用 Jinja 模板引擎,让 SQL 具备编程能力:

– macros/pivot.sql {% macro pivot_columns(column, values) %} {% for value in values %} SUM(CASE WHEN {{ column }} = ‘{{ value }}’ THEN amount ELSE 0 END) AS {{ value }}_revenue {% if not loop.last %},{% endif %} {% endfor %} {% endmacro %}

– 使用宏 SELECT month, {{ pivot_columns(‘category’, [‘Electronics’, ‘Clothing’, ‘Food’]) }} FROM monthly_sales GROUP BY month [/code]

运行 dbt

# 安装 dbt
pip install dbt-postgres

# 初始化项目
dbt init my_project

# 运行所有模型
dbt run

# 运行指定模型
dbt run --select stg_orders

# 运行测试
dbt test

# 生成文档
dbt docs generate
dbt docs serve
[/code]

## Airbyte + dbt 集成模式

Airbyte 和 dbt 的组合是当前最流行的开源 ELT 方案。工作流程如下:

数据源 -> Airbyte (Extract & Load) -> 数据仓库 -> dbt (Transform) -> 分析表 [/code]

具体步骤:

  1. Airbyte 抽取数据:从源系统(PostgreSQL、API、CSV 等)抽取数据
  2. Airbyte 加载数据:将原始数据写入数据仓库的 raw_data schema
  3. dbt Staging:从 raw_data 读取,做基本清洗和类型转换
  4. dbt Intermediate:业务逻辑处理,多表关联
  5. dbt Marts:生成面向分析的主题表

使用 Airbyte 的 dbt 集成

Airbyte 支持在同步完成后自动触发 dbt 运行:

# Airbyte 连接配置中的 dbt 集成
dbtConfig:
  dbtCloudHost: "https://cloud.getdbt.com"
  dbtCloudAccountId: 12345
  dbtCloudJobId: 67890
  dbtCloudApiToken: "your_token"
[/code]

或者使用 Airflow 编排整个流程。我们会在下一章详细介绍 Airflow。

### 文档自动生成

dbt 一个很受欢迎的功能是自动生成文档。运行 dbt docs generate 后,会生成一个包含数据血缘关系的静态网站:

生成文档

dbt docs generate

本地预览

dbt docs serve [/code]

文档中包含每张表的描述、字段定义、数据血缘图(哪些模型依赖哪些源),以及测试结果。这让数据团队的协作效率大幅提升。

与其他工具对比

特性 Airbyte Singer Meltano Fivetran
开源
连接器数量 300+ 200+ 基于 Singer 300+
部署方式 自托管 / Cloud 自托管 自托管 SaaS
增量同步 原生支持 需配置 需配置 原生支持
Normalization 内置 内置
dbt 集成 原生 需手动 原生 原生

小结

Airbyte 和 dbt 的组合代表了现代开源 ELT 的最佳实践。Airbyte 负责从各种数据源抽取和加载原始数据,dbt 在数据仓库内完成转换、测试和文档生成。两者结合,数据团队可以用纯开源方案搭建出媲美商业产品的数据管道。下一章我们将引入工作流编排工具 Airflow,让整个 ETL 流程自动化运行。

Summary: Airbyte 数据集成、dbt 数据转换及 ELT 模式实践。