SQL 在 ETL 中的角色

很多人一提到 ETL 就想到 Python 或 Spark,但 SQL 才是 ETL 领域最基础也最强大的工具。无论你用什么框架,最终的数据操作大多会落到 SQL 上。dbt 用 SQL 做转换,Spark SQL 用 SQL 做分析,甚至 Pandas 的 DataFrame 操作在概念上也和 SQL 的 SELECT、GROUP BY、JOIN 一一对应。

SQL 的优势在于声明式:你告诉数据库"要什么",而不是"怎么做"。数据库优化器会帮你决定最佳的执行路径。对于数据量在千万级以下的 ETL 任务,纯 SQL 方案往往比 Python 方案更简洁、更高效。

CTE:构建可读的 ETL 管道

公用表表达式(CTE)是 SQL ETL 中最有用的语法之一。它让你能把复杂的转换拆成多个逻辑步骤,每个步骤清晰独立,就像管道中的一个个节点。

`sql – 一个典型的 ETL 管道,用 CTE 分步实现 WITH – 步骤 1:抽取原始数据 raw_orders AS ( SELECT * FROM orders WHERE order_date >= ‘2026-01-01’ ),

– 步骤 2:清洗数据 cleaned_orders AS ( SELECT order_id, customer_id, COALESCE(amount, 0) AS amount, COALESCE(quantity, 1) AS quantity, order_date, status FROM raw_orders WHERE status != ‘cancelled’ ),

– 步骤 3:计算衍生字段 enriched_orders AS ( SELECT *, amount * quantity AS total_amount, DATE_TRUNC(‘month’, order_date) AS order_month FROM cleaned_orders ),

– 步骤 4:聚合 monthly_summary AS ( SELECT order_month, COUNT(DISTINCT customer_id) AS active_customers, SUM(total_amount) AS revenue, COUNT(*) AS order_count FROM enriched_orders GROUP BY order_month )

– 最终输出 SELECT * FROM monthly_summary ORDER BY order_month; `

每个 CTE 只做一件事,命名清晰,后续维护的人一眼就能看懂每个步骤在做什么。这比嵌套子查询或者临时表的方式可读性强得多。

窗口函数:高级数据转换

窗口函数是 SQL ETL 中的利器。它能在不改变行数的情况下,对每一行计算其所在窗口的聚合值、排名或偏移量。

行号去重

数据抽取时经常拿到重复记录,用 ROW_NUMBER 可以精确去重:

sql WITH deduped AS ( SELECT *, ROW_NUMBER() OVER ( PARTITION BY order_id ORDER BY updated_at DESC ) AS rn FROM raw_orders ) SELECT * FROM deduped WHERE rn = 1;

这个模式在增量抽取中非常常见。当源系统没有提供变更日志时,我们只能全量拉取,然后用 ROW_NUMBER 取每个订单的最新版本。

滚动计算

sql -- 计算每个客户过去 30 天的累计消费 SELECT customer_id, order_date, amount, SUM(amount) OVER ( PARTITION BY customer_id ORDER BY order_date ROWS BETWEEN 29 PRECEDING AND CURRENT ROW ) AS rolling_30d_amount FROM orders WHERE order_date >= '2026-01-01';

同比环比

sql WITH monthly_revenue AS ( SELECT DATE_TRUNC('month', order_date) AS month, SUM(amount) AS revenue FROM orders GROUP BY 1 ) SELECT month, revenue, LAG(revenue, 1) OVER (ORDER BY month) AS prev_month_revenue, LAG(revenue, 12) OVER (ORDER BY month) AS prev_year_revenue, ROUND( (revenue - LAG(revenue, 1) OVER (ORDER BY month)) / LAG(revenue, 1) OVER (ORDER BY month) * 100, 2 ) AS mom_growth_pct FROM monthly_revenue ORDER BY month;

MERGE / UPSERT:增量加载

增量加载是生产环境 ETL 的核心需求。全量加载每天重刷整张表,数据量大时不可行。增量加载只处理新增和变更的数据,然后用 MERGE 语句合并到目标表。

PostgreSQL 的 MERGE(也叫 UPSERT)语法如下:

sql MERGE INTO target_sales AS t USING ( SELECT * FROM staging_sales WHERE load_date = CURRENT_DATE ) AS s ON t.order_id = s.order_id WHEN MATCHED THEN UPDATE SET amount = s.amount, quantity = s.quantity, status = s.status, updated_at = CURRENT_TIMESTAMP WHEN NOT MATCHED THEN INSERT (order_id, customer_id, amount, quantity, status, created_at) VALUES (s.order_id, s.customer_id, s.amount, s.quantity, s.status, CURRENT_TIMESTAMP);

MySQL 的语法略有不同,使用 ON DUPLICATE KEY UPDATE:

sql INSERT INTO target_sales (order_id, customer_id, amount, quantity, status) SELECT order_id, customer_id, amount, quantity, status FROM staging_sales WHERE load_date = CURRENT_DATE ON DUPLICATE KEY UPDATE amount = VALUES(amount), quantity = VALUES(quantity), status = VALUES(status);

数据清洗函数

SQL 内置的数据清洗函数在 ETL 中非常实用:

`sql – 处理空值 SELECT COALESCE(phone, email, ‘无联系方式’) AS contact, NULLIF(empty_string_field, ‘’) AS meaningful_null, COALESCE(CAST(age AS INTEGER), 0) AS age FROM raw_customers;

– 条件逻辑 SELECT customer_id, CASE WHEN total_spent >= 10000 THEN ‘VIP’ WHEN total_spent >= 5000 THEN ‘黄金’ WHEN total_spent >= 1000 THEN ‘白银’ ELSE ‘普通’ END AS customer_tier FROM customer_summary;

– 字符串清洗 SELECT UPPER(TRIM(email)) AS email_clean, REGEXP_REPLACE(phone, ‘[^0-9]’, ‘’, ‘g’) AS phone_digits_only FROM raw_contacts; `

日期时间处理

日期和时间转换是 ETL 中最频繁的操作之一。不同源系统的日期格式五花八门,需要在加载前统一。

`sql – 字符串转日期 SELECT TO_DATE(‘2026-01-15’, ‘YYYY-MM-DD’) AS date1, TO_TIMESTAMP(‘2026/01/15 14:30:00’, ‘YYYY/MM/DD HH24:MI:SS’) AS ts1;

– 提取日期部分 SELECT order_date, EXTRACT(YEAR FROM order_date) AS year, EXTRACT(MONTH FROM order_date) AS month, EXTRACT(DOW FROM order_date) AS day_of_week, DATE_TRUNC(‘week’, order_date) AS week_start, order_date - INTERVAL ‘1 month’ AS one_month_ago FROM orders;

– 日期范围生成(用于补全缺失日期) SELECT generate_series( ‘2026-01-01’::DATE, ‘2026-12-31’::DATE, ‘1 day’::INTERVAL )::DATE AS calendar_date; `

临时表与 staging 表

在复杂的 ETL 流程中,中间结果需要暂存。临时表和 staging 表是两种常见方案。

`sql – 会话级临时表(仅当前连接可见) CREATE TEMP TABLE temp_cleaned_orders AS SELECT * FROM raw_orders WHERE amount > 0;

– 事务级临时表(事务结束后自动删除) CREATE TEMPORARY TABLE temp_aggregated ( order_month DATE, revenue NUMERIC(12,2) ) ON COMMIT DROP;

– Staging 表(持久化,用于多步骤 ETL) CREATE TABLE staging.daily_orders AS SELECT * FROM raw_orders WHERE load_date = CURRENT_DATE; `

Staging 表的好处是支持断点续传。如果 ETL 在转换步骤失败,下次运行时可以跳过已经加载到 staging 的数据,只处理失败的部分。

SQL 错误处理

SQL 层面的错误处理虽然不如 Python 灵活,但也有一些实用技巧:

`sql – 使用事务保证原子性 BEGIN;

– 数据校验:如果发现异常则回滚 DO DECLARE bad_count INTEGER; BEGIN SELECT COUNT(*) INTO bad_count FROM staging_orders WHERE amount < 0;

IF bad_count > 0 THEN
    RAISE EXCEPTION '发现 % 条金额为负的记录,回滚', bad_count;
END IF;

END ;

COMMIT;

– 使用 ASSERT 做前置检查 DO BEGIN ASSERT (SELECT COUNT() FROM staging_orders) > 0, ‘Staging 表为空’; ASSERT (SELECT COUNT() FROM staging_orders WHERE amount IS NULL) = 0, ‘存在空金额’; END ; `

SQL vs Python:如何选择?

这是一个常见问题。下面给出一些选择建议:

场景 推荐方案 原因
简单过滤聚合 SQL 一行 SQL 搞定,无需引入 Python
多表关联 SQL 数据库 JOIN 经过数十年优化,比 Pandas merge 快
复杂业务逻辑 Python 条件分支、循环、异常处理更灵活
机器学习特征工程 Python scikit-learn、numpy 等库不可替代
调用外部 API Python SQL 无法直接调用 HTTP 接口
超大数据集 SQL / Spark 数据库内计算避免数据传输开销
数据质量检查 SQL 声明式约束,简洁明了

一个实用的原则是:能在数据库里做的就在数据库里做。把数据拉到应用层处理再写回去,中间的网络传输成本往往比计算成本高得多。

SQL ETL 性能优化

索引策略

sql -- 为 ETL 查询创建合适的索引 CREATE INDEX idx_staging_load_date ON staging_orders(load_date); CREATE INDEX idx_staging_status ON staging_orders(status); CREATE INDEX idx_staging_customer_date ON staging_orders(customer_id, order_date);

分区表

对于大表,分区可以显著提升 ETL 性能:

`sql – 按月分区 CREATE TABLE sales ( order_id BIGINT, order_date DATE, amount NUMERIC(12,2) ) PARTITION BY RANGE (order_date);

CREATE TABLE sales_2026_01 PARTITION OF sales FOR VALUES FROM (‘2026-01-01’) TO (‘2026-02-01’);

CREATE TABLE sales_2026_02 PARTITION OF sales FOR VALUES FROM (‘2026-02-01’) TO (‘2026-03-01’); `

批量操作

`sql – 批量插入,避免逐行 INSERT INSERT INTO target_sales (order_id, amount, order_date) SELECT order_id, amount, order_date FROM staging_sales WHERE load_date = CURRENT_DATE;

– 使用 COPY 命令(比 INSERT 快 5-10 倍) COPY target_sales(order_id, amount, order_date) FROM ‘/path/to/clean_data.csv’ DELIMITER ‘,’ CSV HEADER; `

小结

SQL 是 ETL 工程师的看家本领。CTE 让管道逻辑清晰可读,窗口函数处理复杂的行级计算,MERGE 语句实现高效的增量加载,日期函数统一各种时间格式。掌握这些 SQL 技巧,你就能在数据库层面完成大部分 ETL 工作,既高效又可靠。下一章我们将跳出 SQL,看看开源 ETL 工具链如何进一步简化数据集成工作。

Summary: SQL 在 ETL 中的 CTE、窗口函数、MERGE、清洗与性能优化。