Apache Spark SQL FETCH 语句完全指南:游标逐行取值、变量绑定与 NOT FOUND 处理机制
2026/9/20 4:54:43 网站建设 项目流程
  • 大数据
  • 数据分析
  • 批处理
  • 流处理
  • 机器学习
  • 图计算

【免费下载链接】spark

Apache Spark - A unified analytics engine for large-scale data processing

项目地址:https://gitcode.com/gh_mirrors/sp/spark
点击查看免费下载

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—— 可选关键字

NEXTFROM均为语法糖,不影响执行行为。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')分别赋给xy

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

游标返回idname两列,但INTO只给了一个STRUCT<id: INT, name: STRING>变量:两列按位置分别落入result.idresult.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; -- 结果: 0

FETCH NEXT FROM cursor1 INTO xFETCH 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 FOUNDEXIT 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; -- 结果: 3

CURSOR_NO_MORE_ROWS是比NOT FOUND更精确的条件名,二者在'02000'上是同一条件。循环逐行计数,最终统计出行数 3。

四、底层实现:FetchCursorExec 如何工作

FETCH的物理执行节点是 FetchCursorExec.scala,其run()方法揭示了完整的内部流程:

1. 获取脚本执行上下文。通过 CursorCommandUtils.scala 中的getScriptingContextSqlScriptingContextManager取出当前脚本执行上下文;若游标在脚本上下文之外被使用,会抛出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

项目地址:https://gitcode.com/gh_mirrors/sp/spark
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询