1. 为什么大批量数据处理总卡在 PL/SQL 逐行返回上
如果你在 Oracle 里做过数据迁移、ETL 清洗或者报表预处理,大概率遇到过这种场景:源表几百万甚至上亿行,需要逐行做字段转换、拼接、过滤,再写进目标表。最直觉的写法就是INSERT INTO ... SELECT ...,简单直接,但一旦转换逻辑复杂起来,SQL 就变得又长又难维护;换成 PL/SQL 游标循环吧,又慢得让人抓狂。
问题的根子在于 PL/SQL 引擎和 SQL 引擎之间的上下文切换。普通函数每返回一行,就要在两层引擎之间来回切一次,几百万行就是几百万次切换,CPU 时间全耗在切换上了。Oracle 的 pipelined(管道)函数就是为解决这个问题设计的:它允许函数像“流水线”一样边生产边返回数据,调用方可以像查表一样用TABLE()消费结果,中间不需要把整个结果集物化到内存或临时段。
这篇内容面向需要在真实库中优化大批量数据逐行返回与流式消费的开发者。我会从最常规的写法讲起,一步步过渡到管道函数骨架、PIPE ROW与RETURN的配置、BULK COLLECT批量取值、并行度开启,最后用DBMS_OUTPUT和执行计划做验证对比。你跟着操作,能在自己的测试库里跑通并看到性能差异。
2. 前置准备:TaoToken 环境与 Oracle 连接配置
在动手写管道函数之前,先把实验环境理顺。我习惯把模型辅助和文档查询放在 TaoToken 上统一管理,这样写 PL/SQL 时遇到语法细节可以直接对话确认,不用来回翻文档。
TaoToken 的官网入口是 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= ,API 调用地址是 https://taotoken.net/api 。如果你只是想在写代码时随时问 PL/SQL 语法、执行计划解读,用模型对话就够了;如果你要长期做 Oracle 相关的编码和 Agent 任务,可以看 Coding Plan;需要自己生成和管理密钥,去 API Keys 页面;接入细节看接入文档。
Oracle 这边你需要准备:
- 一个有
CREATE TYPE、CREATE PACKAGE、CREATE TABLE权限的账号 - 测试表
T_SS_NORMAL和T_TARGET,源表建议灌入 200 万行左右的数据,方便对比 - 能执行
ALTER SESSION的权限,后面开并行要用
建表语句先跑起来:
CREATE TABLE T_SS_NORMAL ( owner VARCHAR2(30), object_name VARCHAR2(128), subobject_name VARCHAR2(30), object_id NUMBER, data_object_id NUMBER, object_type VARCHAR2(19), created DATE, last_ddl_time DATE, timestamp VARCHAR2(19), status VARCHAR2(7), temporary VARCHAR2(1), generated VARCHAR2(1), secondary VARCHAR2(1) ); CREATE TABLE T_TARGET ( owner VARCHAR2(30), object_name VARCHAR2(128), comm VARCHAR2(10) );源表数据可以从DBA_OBJECTS灌,反复插入几次凑到 200 万行:
INSERT INTO T_SS_NORMAL SELECT owner, object_name, subobject_name, object_id, data_object_id, object_type, created, last_ddl_time, timestamp, status, temporary, generated, secondary FROM DBA_OBJECTS; COMMIT; -- 重复执行若干次直到行数达到预期3. 从常规 INSERT 到管道函数:可复制的配置骨架
3.1 常规写法的瓶颈
最省事的做法就是一个INSERT INTO ... SELECT:
CREATE OR REPLACE PACKAGE BODY PKG_TEST IS PROCEDURE LOAD_TARGET_NORMAL IS BEGIN INSERT INTO T_TARGET (owner, object_name, comm) SELECT owner, object_name, 'xxx' FROM T_SS_NORMAL; COMMIT; END; END PKG_TEST; /这种写法在纯 SQL 能表达转换逻辑时没问题,但一旦转换涉及复杂分支、调用其他函数、或者需要逐行做状态判断,SQL 就会失控。而且它是一次性全量插入,redo 日志和 undo 压力都集中在一个事务里。
3.2 定义对象类型和集合类型
管道函数要作为“表函数”被TABLE()调用,返回值必须是集合类型。先建两个类型:
CREATE TYPE OBJ_TARGET AS OBJECT ( owner VARCHAR2(30), object_name VARCHAR2(128), comm VARCHAR2(10) ); / CREATE OR REPLACE TYPE TYP_ARRAY_TARGET AS TABLE OF OBJ_TARGET; /OBJ_TARGET对应目标表的一行,TYP_ARRAY_TARGET是管道函数的返回集合类型。
3.3 管道函数骨架与 PIPE ROW / RETURN
在包规范里声明函数和过程:
CREATE OR REPLACE PACKAGE PKG_TEST IS FUNCTION PIPE_TARGET(P_SOURCE_DATA IN SYS_REFCURSOR) RETURN TYP_ARRAY_TARGET PIPELINED; PROCEDURE LOAD_TARGET; END PKG_TEST; /包体里实现。核心是PIPELINED关键字、循环里的PIPE ROW、以及最后的RETURN:
CREATE OR REPLACE PACKAGE BODY PKG_TEST IS FUNCTION PIPE_TARGET(P_SOURCE_DATA IN SYS_REFCURSOR) RETURN TYP_ARRAY_TARGET PIPELINED IS R_TARGET_DATA OBJ_TARGET := OBJ_TARGET(NULL, NULL, NULL); R_SOURCE_DATA T_SS_NORMAL%ROWTYPE; BEGIN LOOP FETCH P_SOURCE_DATA INTO R_SOURCE_DATA; EXIT WHEN P_SOURCE_DATA%NOTFOUND; R_TARGET_DATA.owner := R_SOURCE_DATA.owner; R_TARGET_DATA.object_name := R_SOURCE_DATA.object_name; R_TARGET_DATA.comm := 'xxx'; PIPE ROW(R_TARGET_DATA); END LOOP; CLOSE P_SOURCE_DATA; RETURN; END; PROCEDURE LOAD_TARGET IS BEGIN INSERT INTO T_TARGET (owner, object_name, comm) SELECT owner, object_name, comm FROM TABLE(PIPE_TARGET(CURSOR(SELECT * FROM T_SS_NORMAL))); COMMIT; END; END PKG_TEST; /PIPE ROW的作用是把当前构造好的对象推入返回集合,调用方立刻就能消费到这一行,不需要等函数全部执行完。RETURN在管道函数里不带值,只表示生产结束。
3.4 用 BULK COLLECT 减少上下文切换
上面逐行FETCH的写法,每取一行就切换一次引擎。改成批量取值:
FUNCTION PIPE_TARGET_ARRAY( P_SOURCE_DATA IN SYS_REFCURSOR, P_LIMIT_SIZE IN PLS_INTEGER DEFAULT 100 ) RETURN TYP_ARRAY_TARGET PIPELINED IS R_TARGET_DATA OBJ_TARGET := OBJ_TARGET(NULL, NULL, NULL); TYPE TYP_SOURCE_DATA IS TABLE OF T_SS_NORMAL%ROWTYPE INDEX BY PLS_INTEGER; AA_SOURCE_DATA TYP_SOURCE_DATA; BEGIN LOOP FETCH P_SOURCE_DATA BULK COLLECT INTO AA_SOURCE_DATA LIMIT P_LIMIT_SIZE; EXIT WHEN AA_SOURCE_DATA.COUNT = 0; FOR i IN 1 .. AA_SOURCE_DATA.COUNT LOOP R_TARGET_DATA.owner := AA_SOURCE_DATA(i).owner; R_TARGET_DATA.object_name := AA_SOURCE_DATA(i).object_name; R_TARGET_DATA.comm := 'xxx'; PIPE ROW(R_TARGET_DATA); END LOOP; END LOOP; CLOSE P_SOURCE_DATA; RETURN; END;LIMIT 100控制每次批量取的行数,太小起不到减少切换的作用,太大又占 PGA 内存,100 到 1000 之间比较稳妥。
3.5 开启并行度与直接路径插入
数据量再大,单进程还是慢。给管道函数加PARALLEL_ENABLE,配合PARTITION子句让数据分片并行处理:
FUNCTION PIPE_TARGET_PARALLEL( P_SOURCE_DATA IN SYS_REFCURSOR, P_LIMIT_SIZE IN PLS_INTEGER DEFAULT 100 ) RETURN TYP_ARRAY_TARGET PIPELINED PARALLEL_ENABLE(PARTITION P_SOURCE_DATA BY ANY) IS R_TARGET_DATA OBJ_TARGET := OBJ_TARGET(NULL, NULL, NULL); TYPE TYP_SOURCE_DATA IS TABLE OF T_SS_NORMAL%ROWTYPE INDEX BY PLS_INTEGER; AA_SOURCE_DATA TYP_SOURCE_DATA; BEGIN LOOP FETCH P_SOURCE_DATA BULK COLLECT INTO AA_SOURCE_DATA LIMIT P_LIMIT_SIZE; EXIT WHEN AA_SOURCE_DATA.COUNT = 0; FOR i IN 1 .. AA_SOURCE_DATA.COUNT LOOP R_TARGET_DATA.owner := AA_SOURCE_DATA(i).owner; R_TARGET_DATA.object_name := AA_SOURCE_DATA(i).object_name; R_TARGET_DATA.comm := 'xxx'; PIPE ROW(R_TARGET_DATA); END LOOP; END LOOP; CLOSE P_SOURCE_DATA; RETURN; END;调用时开并行 DML,并用 hint 指定并行度:
PROCEDURE LOAD_TARGET_PARALLEL IS BEGIN EXECUTE IMMEDIATE 'ALTER SESSION ENABLE PARALLEL DML'; INSERT /*+ PARALLEL(t,4) */ INTO T_TARGET t (owner, object_name, comm) SELECT owner, object_name, comm FROM TABLE(PIPE_TARGET_PARALLEL( CURSOR(SELECT /*+ PARALLEL(s,4) */ * FROM T_SS_NORMAL s), 100)); COMMIT; END;并行度 4 意味着 4 个进程同时消费管道函数输出,插入走直接路径,redo 生成量会明显下降。
4. 验证请求与成功结果:DBMS_OUTPUT 与执行计划
4.1 用 DBMS_OUTPUT 看管道函数输出
先小批量验证函数逻辑对不对。开SERVEROUTPUT,取前 5 行看看:
SET SERVEROUTPUT ON SIZE UNLIMITED; DECLARE V_CNT PLS_INTEGER := 0; BEGIN FOR R IN ( SELECT owner, object_name, comm FROM TABLE(PKG_TEST.PIPE_TARGET(CURSOR(SELECT * FROM T_SS_NORMAL))) WHERE ROWNUM <= 5 ) LOOP DBMS_OUTPUT.PUT_LINE(R.owner || ' | ' || R.object_name || ' | ' || R.comm); V_CNT := V_CNT + 1; END LOOP; DBMS_OUTPUT.PUT_LINE('sample rows: ' || V_CNT); END; /预期输出类似:
SYS | I_OBJ1 | xxx SYS | I_OBJ2 | xxx SYS | TAB$ | xxx SYS | CLU$ | xxx SYS | I_COL1 | xxx sample rows: 5注意WHERE ROWNUM <= 5能让管道函数提前停止生产,这也是流式消费的一个好处——不需要全量算完。
4.2 对比执行计划
分别看常规 INSERT 和管道函数 INSERT 的执行计划:
EXPLAIN PLAN FOR INSERT INTO T_TARGET (owner, object_name, comm) SELECT owner, object_name, 'xxx' FROM T_SS_NORMAL; SELECT * FROM TABLE(DBMS_XPLAN.DISPLAY); EXPLAIN PLAN FOR INSERT INTO T_TARGET (owner, object_name, comm) SELECT owner, object_name, comm FROM TABLE(PKG_TEST.PIPE_TARGET(CURSOR(SELECT * FROM T_SS_NORMAL))); SELECT * FROM TABLE(DBMS_XPLAN.DISPLAY);管道函数版本的计划里会出现COLLECTION ITERATOR PICKLER FETCH或TABLE()相关的行源,说明数据是流式从函数里取出来的,而不是先物化。
4.3 实测耗时对比
在 200 万行数据上跑一遍,记录时间:
SET TIMING ON; -- 常规写法 EXEC PKG_TEST.LOAD_TARGET_NORMAL; -- 管道函数 + BULK COLLECT EXEC PKG_TEST.LOAD_TARGET; -- 管道函数 + 并行 EXEC PKG_TEST.LOAD_TARGET_PARALLEL;我实测下来,200 万行从常规的 20 多秒降到管道函数加批量取值的 8 秒左右,开并行后还能再压一截,redo 日志量的下降更明显。具体数字因机器和参数而异,但趋势是一致的。
5. 本篇常见错排查
5.1 ORA-22905:无法访问非嵌套表项
报错信息类似ORA-22905: cannot access rows from a non-nested table item。原因通常是TABLE()里传的不是集合类型,或者函数没有声明PIPELINED。检查两点:函数返回类型必须是TABLE OF定义的集合类型;调用时必须用TABLE(函数名(...))包起来。
5.2 PLS-00382:表达式类型错误
PIPE ROW里推入的对象类型必须和函数声明的返回集合元素类型完全一致。如果你OBJ_TARGET定义了三个字段,PIPE ROW里推的也必须是OBJ_TARGET实例,不能是%ROWTYPE或者别的对象类型。
5.3 并行没生效
开了PARALLEL_ENABLE但执行计划里没有并行,常见原因:没有执行ALTER SESSION ENABLE PARALLEL DML;表没有开并行;或者CURSOR里的查询本身不支持并行。另外注意,管道函数里的PARTITION ... BY ANY只是告诉优化器可以任意分区,实际并行度还是由 hint 和会话参数决定。
5.4 结果集顺序不确定
管道函数加并行后,输出顺序不再保证和源表一致。如果你的业务依赖顺序,要么在最终SELECT里加ORDER BY,要么别开并行。流式消费和有序输出本身就有取舍。
5.5 内存溢出
BULK COLLECT的LIMIT设得太大,或者管道函数里累积了太多中间集合,会导致 PGA 暴涨。保持LIMIT在合理范围,及时CLOSE游标,不要在函数里做无限制的集合追加。
6. 把管道函数接进你的日常开发流
管道函数真正的价值不只是单次迁移快,而是它能作为“数据流组件”嵌进更大的处理链路里。你可以把多个管道函数串起来,前一个的输出作为后一个的输入,形成流式管道;也可以配合外部表、数据泵做分阶段处理。
写 PL/SQL 时遇到不确定的语法或者想快速验证执行计划解读,我一般直接开模型对话问,比翻文档快。需要自己生成密钥做自动化调用,去 API Keys 页面拿;接入细节和参数说明看接入文档。长期做 Oracle 编码和 Agent 任务的话,Coding Plan 会更顺手。
最后留一个实用习惯:每次改完管道函数,先用WHERE ROWNUM <= 10小批量验证输出,再跑全量。这样能避免逻辑错误在大数据量下放大成灾难。管道函数配合BULK COLLECT和并行,基本能覆盖 Oracle 里大部分大数据量逐行转换的场景,剩下的就是根据你的实际表结构和转换规则去调整骨架了。