4 minutes
SQL 在 ETL 中的应用
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、清洗与性能优化。