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;
/

详细见:Oracle PL/SQL 内置包大全详解


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. 最佳实践

  1. 外部表:灵活
  2. 直接路径:性能
  3. 分区交换:高效
  4. 并行:吞吐
  5. MERGE:UPSERT
  6. BULK:批量
  7. 错误处理:完整
  8. 监控:进度
  9. 数据质量:校验
  10. 文档化:流程

16. 参考资料

[1] Oracle Database Data Warehousing Guide 19c https://docs.oracle.com/en/database/oracle/oracle-database/19/dwh/