- 大数据
- 数据分析
- 批处理
- 流处理
- 机器学习
- 图计算
【免费下载链接】spark
Apache Spark - A unified analytics engine for large-scale data processing
FETCH是 Apache Spark SQL 脚本编程(SQL Scripting)中用于从已打开游标(Cursor)逐行取出结果集并写入变量的核心语句。本文围绕 docs/control-flow/fetch-stmt.md 展开,完整覆盖其语法、参数规则、类型兼容约束与异常处理语义,并结合 sql/core/src/main/scala/org/apache/spark/sql/execution/command/v2/FetchCursorExec.scala 等源码说明底层实现。读完本文,你将掌握如何在复合语句中声明、打开、逐行消费并关闭游标,如何用REPEAT/WHILE循环配合CONTINUE/EXIT HANDLER优雅处理结果集耗尽,以及STRUCT变量一键接收多列的 SQL 标准特例。
一、FETCH 语句是什么
在 SQL 脚本编程中,OPEN负责执行游标的查询并把结果集物化在内存中,而FETCH则负责把结果集里的行"一行一行"地取出来赋值给变量。FETCH从打开的游标中取出下一行,并将列值依次赋给指定的变量;当结果集中已无更多行可取时,会抛出CURSOR_NO_MORE_ROWS条件(SQLSTATE'02000')。
它是游标生命周期DECLARE CURSOR → OPEN → FETCH... → CLOSE的中间枢纽:
- OPEN 语句 执行游标查询并定位到第一行之前;
FETCH每执行一次,游标位置前进一行;- CLOSE 语句 释放游标关联的结果集资源。
这三个语句以及游标声明、异常处理器共同组成 复合语句(Compound Statement),进而服务于 SQL 脚本编程(如 Spark SQL CLI、Thrift Server 中执行的BEGIN...END脚本块)。
二、语法与参数详解
FETCH的完整语法为:
FETCH [ [ NEXT ] FROM ] cursor_name INTO variable_name [, ...]1.cursor_name—— 游标名
一个处于打开状态的游标名。游标名可以带复合语句标签进行限定,用于引用外层作用域中声明的游标,例如outer_label.my_cursor。这一能力允许嵌套的BEGIN...END块在内部访问并消费外层游标的行。
2.NEXT FROM—— 可选关键字
NEXT与FROM均为语法糖,不影响执行行为。Spark 目前只支持正向(forward-only)取数,不支持回退、跳转或按绝对/相对偏移取行,这与其他数据库引擎中复杂的游标滚动语义不同,但足以覆盖绝大多数逐行遍历场景。
3.variable_name—— 目标变量
接收列值的本地变量或会话变量。变量与游标结果集列数必须匹配,只有一种例外情况:
- 如果只指定一个变量,且该变量的类型是
STRUCT,而游标返回多列,则各列值按位置依次赋值到结构体(struct)的字段中。
列的数据类型必须与目标变量(或结构体字段)兼容,兼容性依据**存储赋值规则(store assignment rules)**判定。
目标变量的来源有两种,见 compound-stmt.md 与 SQL Scripting 文档:
- 本地变量:在复合语句内用
DECLARE variable_name datatype [DEFAULT expr]声明; - 会话变量:在会话级别用
DECLARE VARIABLE创建,例如DECLARE VARIABLE x INT;。
三、完整示例与运行结果
以下示例均来自官方文档,可以直接在支持 SQL 脚本的 Spark SQL 会话中运行(脚本以>提示符标记输入,其后一行是查询结果)。
1. 基础取行
BEGIN DECLARE x INT; DECLARE y STRING; DECLARE my_cursor CURSOR FOR SELECT id, 'row_' || id FROM range(3); OPEN my_cursor; FETCH my_cursor INTO x, y; VALUES (x, y); CLOSE my_cursor; END; -- 结果: 0|row_0先声明两个变量与一个游标,OPEN执行SELECT id, 'row_' || id FROM range(3)生成 3 行数据,第一次FETCH取出第一行(0, 'row_0')分别赋给x、y。
2. 用 REPEAT 循环消费全部行
BEGIN DECLARE x INT; DECLARE done BOOLEAN DEFAULT false; DECLARE total INT DEFAULT 0; DECLARE sum_cursor CURSOR FOR SELECT id FROM range(5); DECLARE CONTINUE HANDLER FOR NOT FOUND SET done = true; OPEN sum_cursor; REPEAT FETCH sum_cursor INTO x; IF NOT done THEN SET total = total + x; END IF; UNTIL done END REPEAT; CLOSE sum_cursor; VALUES (total); END; -- 结果: 10这是最典型的"遍历游标求和"模式:先声明CONTINUE HANDLER FOR NOT FOUND SET done = true;,当FETCH越界抛出NOT FOUND(SQLSTATE'02xxx')时,处理器把done置为true且不中断脚本,REPEAT...UNTIL done END REPEAT随之退出。0+1+2+3+4 = 10。
3. 取多列到一个 STRUCT 变量
BEGIN DECLARE result STRUCT<id: INT, name: STRING>; DECLARE struct_cursor CURSOR FOR SELECT id, 'name_' || id FROM range(3); OPEN struct_cursor; FETCH struct_cursor INTO result; VALUES (result.id, result.name); CLOSE struct_cursor; END; -- 结果: 0|name_0游标返回id、name两列,但INTO只给了一个STRUCT<id: INT, name: STRING>变量:两列按位置分别落入result.id与result.name,免去了逐一声明多个变量的繁琐。
4. 使用NEXT FROM可选语法
BEGIN DECLARE x INT; DECLARE cursor1 CURSOR FOR SELECT id FROM range(3); OPEN cursor1; FETCH NEXT FROM cursor1 INTO x; VALUES (x); CLOSE cursor1; END; -- 结果: 0FETCH NEXT FROM cursor1 INTO x与FETCH cursor1 INTO x完全等价。
5. 带标签的限定游标名
BEGIN outer_lbl: BEGIN DECLARE outer_cur CURSOR FOR SELECT id FROM range(5); DECLARE x INT; OPEN outer_cur; inner_lbl: BEGIN FETCH outer_lbl.outer_cur INTO x; VALUES (x); END; CLOSE outer_cur; END; END; -- 结果: 0内层BEGIN...END块通过outer_lbl.outer_cur访问外层游标并取出一行,展示了标签限定在嵌套作用域中的作用。
6. 用 EXIT HANDLER 处理结果集耗尽
BEGIN DECLARE x INT; DECLARE my_cursor CURSOR FOR SELECT id FROM range(2); DECLARE EXIT HANDLER FOR NOT FOUND BEGIN VALUES ('No more rows'); END; OPEN my_cursor; FETCH my_cursor INTO x; FETCH my_cursor INTO x; FETCH my_cursor INTO x; -- 触发 EXIT HANDLER VALUES ('This will not execute'); CLOSE my_cursor; END; -- 结果: No more rows游标只有 2 行,第三次FETCH触发NOT FOUND,EXIT HANDLER打印'No more rows'后整个复合语句立即退出,其后的VALUES ('This will not execute')不再执行。
7. 针对CURSOR_NO_MORE_ROWS的专用处理器
BEGIN DECLARE x INT DEFAULT 0; DECLARE done BOOLEAN DEFAULT false; DECLARE count INT DEFAULT 0; DECLARE my_cursor CURSOR FOR SELECT id FROM range(3); DECLARE CONTINUE HANDLER FOR CURSOR_NO_MORE_ROWS SET done = true; OPEN my_cursor; WHILE NOT done DO FETCH my_cursor INTO x; IF NOT done THEN SET count = count + 1; END IF; END WHILE; CLOSE my_cursor; VALUES (count); END; -- 结果: 3CURSOR_NO_MORE_ROWS是比NOT FOUND更精确的条件名,二者在'02000'上是同一条件。循环逐行计数,最终统计出行数 3。
四、底层实现:FetchCursorExec 如何工作
FETCH的物理执行节点是 FetchCursorExec.scala,其run()方法揭示了完整的内部流程:
1. 获取脚本执行上下文。通过 CursorCommandUtils.scala 中的getScriptingContext从SqlScriptingContextManager取出当前脚本执行上下文;若游标在脚本上下文之外被使用,会抛出CURSOR_OUTSIDE_SCRIPT错误。
2. 游标状态机驱动。游标状态记录在脚本上下文中,从源码可看到两类状态:
CursorOpened(iter, schema):OPEN后的初始状态,迭代器在 OPEN 时已创建;首次FETCH会将其转换为CursorFetching(iter, schema);CursorFetching(iter, schema):后续FETCH直接复用已有迭代器。
若状态既非Opened也非Fetching(例如游标尚未打开就FETCH),则抛出CURSOR_NOT_OPEN错误——这正是文档 Notes 中"从未打开游标取数会报错"的实现依据。
3. 迭代器越界判定。if (!iterator.hasNext)时抛出CURSOR_NO_MORE_ROWS错误(SQLSTATE'02000'),对应文档中"无更多行"的语义。
4. 两种赋值路径。shouldFetchIntoStruct判断是否命中"单变量 +STRUCT+ 多列"特例:
- 普通路径
fetchIntoVariables:先做元数校验,变量数必须等于列数,否则抛ASSIGNMENT_ARITY_MISMATCH;随后逐列取出InternalRow的值,若源类型与目标变量类型不同,通过Cast表达式(ansiEnabled = true)按 ANSI 存储赋值规则做隐式转换后再赋值; - STRUCT 特例
fetchIntoStruct:同样校验结构体字段数与游标列数一致,然后逐字段构造Cast表达式,最终用CreateStruct组装成结构体整体赋给变量。
5. 变量落盘。assignToVariable根据变量所属目录选择管理器:本地变量(FakeLocalCatalog)交给ScriptingVariableManager,会话变量(FakeSystemCatalog)交给tempVariableManager,与 SetVariableExec.scala 的取值逻辑保持一致;大小写敏感性遵循spark.sql.caseSensitiveAnalysis配置。
五、注意事项与最佳实践
1. 游标必须先 OPEN 再 FETCH
FETCH之前必须已执行 OPEN 语句。从未打开(或已关闭)的游标上执行FETCH会抛出CURSOR_NOT_OPEN错误。同时注意一个游标只能被打开一次,重复OPEN会得到CURSOR_ALREADY_OPEN。
2. 每执行一次 FETCH,游标前进一行
FETCH是"取下一行"语义,不会重复返回同一行,因此遍历完整结果集通常需要与循环(REPEAT/WHILE/LOOP)配合,并通过条件处理器感知终点。
3. 无更多行时触发 CURSOR_NO_MORE_ROWS
当结果集被消费完毕再FETCH时:
- SQLSTATE:
'02000' - 错误条件名:
CURSOR_NO_MORE_ROWS - 该条件属于
NOT FOUND处理器覆盖范围(NOT FOUND捕获所有 SQLSTATE'02xxx'类条件),相关语义见 compound-stmt.md 中declare_handler的说明
4. 未声明处理器时静默忽略
如果既没有为NOT FOUND声明CONTINUE HANDLER,也没有声明EXIT HANDLER,则结果集耗尽这一完成条件(completion condition)会被静默忽略,脚本继续执行后续语句。这一设计让"不加处理器、单纯取到数据耗尽为止"的脚本也能正常运行。
5. 类型兼容遵循存储赋值规则
- 类型一致时直接赋值;
- 类型不一致但可隐式转换时,会先做强制类型转换(实现上
ansiEnabled = true,即按 ANSI 模式转换); - 完全不兼容的类型会抛出类型不匹配错误。
6. 变量作用域
FETCH INTO的目标既可以是复合语句内DECLARE的本地变量,也可以是会话级DECLARE VARIABLE创建的会话变量,二者分别由脚本上下文变量管理器与会话临时变量管理器维护(见上文实现分析)。
7. 用显式 CLOSE 释放资源
游标在复合语句正常退出、EXIT处理器触发或未处理异常退出时会被隐式关闭(见 close-stmt.md),但文档建议不再需要时显式调用 CLOSE 语句:既能提前释放与结果集关联的内存资源,也让代码意图更清晰。关闭后的游标仍可再次OPEN以新参数重新执行查询。
六、相关文档导航
- Compound Statement(复合语句):
FETCH的宿主环境,变量、游标与处理器都在此声明 - OPEN Statement:执行游标查询、绑定参数并物化结果集
- CLOSE Statement:释放游标资源、支持参数化重开
- WHILE Statement / REPEAT Statement:与
FETCH组合实现逐行遍历的循环语句 - SQL Scripting:Spark SQL 脚本编程总览
- 核心实现:FetchCursorExec.scala(FETCH 物理执行)、CursorCommandUtils.scala(脚本上下文获取)
- 大数据
- 数据分析
- 批处理
- 流处理
- 机器学习
- 图计算
【免费下载链接】spark
Apache Spark - A unified analytics engine for large-scale data processing
相关推荐
Apache Spark SQL Compound Statement 详解:BEGIN...END 复合语句、变量、游标与异常处理实战指南
Apache Spark SQL Compound Statement 详解:BEGIN...END 复合语句、变量、游标与异常处理实战指南 导读 Compou
大数据数据分析批处理流处理机器学习图计算Apache Spark SQL 变量赋值语句 SET VAR 完全指南:语法、规则与源码实现解析
Apache Spark SQL 变量赋值语句 SET VAR 完全指南:语法、规则与源码实现解析 导读 SET VAR (也可写作 SET VARIABLE
大数据数据分析批处理流处理机器学习图计算Apache Spark SQL 脚本化中的 OPEN 语句:游标开启、参数绑定与状态机详解
Apache Spark SQL 脚本化中的 OPEN 语句:游标开启、参数绑定与状态机详解 导读 OPEN 语句是 Spark SQL Scripting(S
大数据数据分析批处理流处理机器学习图计算
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考