Oracle 数据库 ETL 设计详解
Oracle 数据库 ETL 设计详解
适用版本:Oracle Database 10g / 11g / 12c / 19c / 23ai 文档版本:v1.0 / 2026-07
1. 概述
ETL(Extract-Transform-Load)设计[1]:
详细见:Oracle 数据仓库 ETL 详解。
2. 架构
2.1 模型
- 源系统
- Staging
- ODS(操作数据存储)
- Data Warehouse
- Data Mart
2.2 模式
- ELT:先加载后转换
- ETL:先转换后加载
- Hybrid
3. Extract
3.1 全量
INSERT INTO stg_emp SELECT * FROM source_emp@src_db;
3.2 增量
-- 时间戳
INSERT INTO stg_emp
SELECT * FROM source_emp@src_db
WHERE updated_at > :last_extract;
-- CDC
-- LogMiner
-- GoldenGate
3.3 外部表
CREATE TABLE ext_sales (...)
ORGANIZATION EXTERNAL (
TYPE ORACLE_LOADER
DEFAULT DIRECTORY data_dir
LOCATION ('sales.csv')
);
INSERT INTO stg_sales SELECT * FROM ext_sales;
详细见:Oracle 数据库外部表详解。
4. Transform
4.1 清洗
-- 去重
DELETE FROM stg_emp WHERE ROWID IN (
SELECT rid FROM (
SELECT ROWID rid, ROW_NUMBER() OVER (PARTITION BY email ORDER BY id) rn
FROM stg_emp
) WHERE rn > 1
);
-- 标准化
UPDATE stg_emp SET
name = TRIM(UPPER(name)),
email = LOWER(email),
phone = REGEXP_REPLACE(phone, '[^0-9]', '');
-- 默认值
UPDATE stg_emp SET
status = NVL(status, 'ACTIVE'),
hire_date = NVL(hire_date, SYSDATE);
4.2 转换
-- 类型
UPDATE stg_emp SET salary = TO_NUMBER(salary_str);
UPDATE stg_emp SET hire_date = TO_DATE(hire_str, 'YYYY-MM-DD');
-- 拆分
INSERT INTO emp_addr (emp_id, addr) SELECT id, address FROM stg_emp;
-- 合并
INSERT INTO emp_full
SELECT e.id, e.name, a.address, p.phone
FROM stg_emp e
LEFT JOIN emp_addr a ON e.id = a.emp_id
LEFT JOIN emp_phone p ON e.id = p.emp_id;
4.3 查找
-- 维度查找
MERGE INTO fact_sales f
USING dim_dept d
ON (f.dept_name = d.name)
WHEN MATCHED THEN UPDATE SET f.dept_id = d.id;
详细见:Oracle MERGE 语句详解。
4.4 聚合
INSERT INTO agg_sales_monthly
SELECT EXTRACT(YEAR FROM sale_date), EXTRACT(MONTH FROM sale_date),
dept_id, SUM(amount)
FROM fact_sales
GROUP BY EXTRACT(YEAR FROM sale_date), EXTRACT(MONTH FROM sale_date), dept_id;
5. Load
5.1 直接路径
INSERT /*+ APPEND */ INTO sales SELECT * FROM stg_sales;
5.2 分区交换
-- 准备
CREATE TABLE sales_stg AS SELECT * FROM sales WHERE 1=0;
-- 加载
INSERT /*+ APPEND */ INTO sales_stg SELECT * FROM source;
-- 交换
ALTER TABLE sales EXCHANGE PARTITION p2025_07 WITH TABLE sales_stg;
5.3 批量
DECLARE
TYPE sales_tab IS TABLE OF sales%ROWTYPE;
v_sales sales_tab;
BEGIN
SELECT * BULK COLLECT INTO v_sales FROM stg_sales LIMIT 10000;
FORALL i IN 1..v_sales.COUNT
INSERT INTO sales VALUES v_sales(i);
END;
/
详细见:Oracle BULK COLLECT 与 FORALL 详解。
6. SQL*Loader
6.1 控制
LOAD DATA
INFILE 'sales.csv'
APPEND INTO TABLE sales
FIELDS TERMINATED BY ','
(id, sale_date DATE 'YYYY-MM-DD', amount)
6.2 直接路径
sqlldr scott/tiger control=sales.ctl direct=true parallel=true
详细见:Oracle SQL Loader 详解。
7. Data Pump
7.1 导出
expdp scott/tiger DIRECTORY=dp DUMPFILE=sales.dmp TABLES=sales
7.2 导入
impdp scott/tiger DIRECTORY=dp DUMPFILE=sales.dmp TABLE_EXISTS_ACTION=APPEND
详细见:Oracle Data Pump 全集。
8. 调度
8.1 Scheduler
BEGIN
DBMS_SCHEDULER.CREATE_JOB(
job_name => 'etl_daily',
job_type => 'PLSQL_BLOCK',
job_action => 'BEGIN etl_pkg.run; END;',
start_date => SYSTIMESTAMP,
repeat_interval => 'FREQ=DAILY; BYHOUR=2',
enabled => TRUE
);
END;
/
8.2 Chain
BEGIN
DBMS_SCHEDULER.CREATE_CHAIN('etl_chain');
DBMS_SCHEDULER.DEFINE_PROGRAM_ACTION(...);
DBMS_SCHEDULER.DEFINE_CHAIN_STEP(...);
DBMS_SCHEDULER.DEFINE_CHAIN_RULE(...);
DBMS_SCHEDULER.ENABLE('etl_chain');
END;
/
9. 错误处理
9.1 日志表
CREATE TABLE etl_log (
id NUMBER GENERATED ALWAYS AS IDENTITY,
job_name VARCHAR2(100),
step VARCHAR2(100),
status VARCHAR2(20),
error_msg VARCHAR2(4000),
rows_processed NUMBER,
start_time TIMESTAMP,
end_time TIMESTAMP
);
9.2 SAVE EXCEPTIONS
FORALL i IN 1..v_data.COUNT SAVE EXCEPTIONS
INSERT INTO target VALUES v_data(i);
EXCEPTION
WHEN OTHERS THEN
FOR i IN 1..SQL%BULK_EXCEPTIONS.COUNT LOOP
INSERT INTO etl_error_log VALUES (
SQL%BULK_EXCEPTIONS(i).ERROR_INDEX,
SQL%BULK_EXCEPTIONS(i).ERROR_CODE
);
END LOOP;
9.3 错误表
-- LOG ERRORS
INSERT INTO target (...) VALUES (...)
LOG ERRORS INTO etl_errors ('batch_1') REJECT LIMIT UNLIMITED;
10. 性能
10.1 并行
ALTER SESSION ENABLE PARALLEL DML;
INSERT /*+ PARALLEL(s 8) APPEND */ INTO sales s
SELECT /*+ PARALLEL(st 8) */ * FROM stg_sales st;
10.2 直接路径
- 跳过 SQL 处理
- 高速
- 索引维护
10.3 分区
- 分区交换
- 局部索引
- 增量
详细见:Oracle 并行查询详解。
11. 监控
11.1 进度
SELECT sid, serial#, opname, sofar, totalwork,
ROUND(sofar/totalwork*100, 2) AS pct
FROM v$session_longops
WHERE opname LIKE 'ETL%';
11.2 日志
SELECT job_name, step, status, rows_processed,
start_time, end_time
FROM etl_log
WHERE job_name = 'etl_daily'
ORDER BY start_time DESC;
12. 数据质量
12.1 校验
-- 完整性
SELECT COUNT(*) FROM fact_sales WHERE dept_id IS NULL;
-- 一致性
SELECT COUNT(*) FROM fact_sales f
WHERE NOT EXISTS (SELECT 1 FROM dim_dept d WHERE d.id = f.dept_id);
-- 范围
SELECT COUNT(*) FROM sales WHERE amount < 0;
12.2 报告
-- 数据质量报告
SELECT 'TOTAL' AS metric, COUNT(*) AS value FROM sales
UNION ALL
SELECT 'NULL_DEPT', COUNT(*) FROM sales WHERE dept_id IS NULL
UNION ALL
SELECT 'NEGATIVE_AMOUNT', COUNT(*) FROM sales WHERE amount < 0;
13. 应用场景
13.1 数据仓库
- 每日 ETL
- 星型模式
- 聚合
- 报表
13.2 数据迁移
- 系统
- 平台
- 版本
13.3 数据同步
- 跨系统
- 实时
- 批量
14. 常见坑与排错
14.1 数据倾斜
- 并行不均
- 分区
- Skew
14.2 约束
- 检查
- 异常表
- 清洗
14.3 性能
- 索引
- 触发器
- 直接路径
15. 最佳实践
- 外部表:灵活
- 直接路径:性能
- 分区交换:高效
- 并行:吞吐
- MERGE:UPSERT
- BULK:批量
- 错误处理:完整
- 监控:进度
- 数据质量:校验
- 文档化:流程
16. 参考资料
[1] Oracle Database Data Warehousing Guide 19c https://docs.oracle.com/en/database/oracle/oracle-database/19/dwh/