- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本文是 Apache Beam 内置 SQL 方言——Beam Calcite SQL 的查询语法(Query Syntax)权威参考。Beam SQL 以 Apache Calcite 的 SQL 语法为基底,为支持 Beam 统一的批处理/流处理(Batch/Streaming)模型引入了 JOIN 与窗口/触发等扩展。读完本文,你将掌握从SELECT、FROM、JOIN、WHERE、GROUP BY、HAVING、集合运算符(UNION/INTERSECT/EXCEPT)、LIMIT/OFFSET、WITH到别名(Aliases)的完整语法体系,并理解这些语句在 Beam SDK 中如何被翻译成可执行的 Beam 变换。
本文对应仓库文档:query-syntax.md,其余 Beam Calcite SQL 参考资料包括 overview.md(方言总览与函数支持矩阵)、data-types.md、scalar-functions.md 与 aggregate-functions.md。
概述:查询语句在 Beam 中的角色
查询语句(Query statement)扫描一个或多个表(table)或表达式(expression),并返回计算后的结果行。Beam Calcite SQL 的查询语义总体上是标准的 SQL 语义,但为支撑 Beam 统一的批/流模型,提供了两类扩展:
- JOIN 扩展:详见 extensions/joins.md;
- 窗口与触发扩展(Windowing & Triggering):详见 extensions/windowing-and-triggering.md。
Beam SQL 的主要功能载体就是SELECT语句,它负责查询与关联数据。Beam Calcite SQL 支持的操作是 Apache Calcite SQL 语法的一个子集,且在原文档中明确标注:Beam Calcite SQL 是 Beam SQL 的默认方言。
从源码实现看,Beam SQL 的查询执行链路位于sdks/java/extensions/sql模块:CalciteQueryPlanner负责基于 Calcite(vendor 版本calcite.v1_28_0)解析 SQL 并产出逻辑计划,随后各类BeamRelNode(位于sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/rel/)把逻辑算子翻译成 BeamPTransform。语法层面支持的每条子句,基本都能在impl/rel目录中找到对应的 Rel 实现类(如BeamJoinRel、BeamUnionRel、BeamSortRel、BeamAggregationRel等)。
SQL 语法总览
Beam Calcite SQL 查询语句的完整语法如下(方括号[ ]表示可选子句,圆括号( )表示字面括号,竖线|表示逻辑或,花括号{ }括起一组选项,[, ...]表示前一项可以逗号分隔重复出现):
query_statement: [ WITH with_query_name AS ( query_expr ) [, ...] ] query_expr query_expr: { select | ( query_expr ) | query_expr set_op query_expr } [ LIMIT count [ OFFSET skip_rows ] ] select: SELECT [{ ALL | DISTINCT }] { [ expression. ]* [ EXCEPT ( column_name [, ...] ) ] [ REPLACE ( expression [ AS ] column_name [, ...] ) ] | expression [ [ AS ] alias ] } [, ...] [ FROM from_item [, ...] ] [ WHERE bool_expression ] [ GROUP BY { expression [, ...] } ] [ HAVING bool_expression ] set_op: UNION { ALL | DISTINCT } | INTERSECT DISTINCT | EXCEPT DISTINCT from_item: { table_name [ [ AS ] alias ] | join | ( query_expr ) [ [ AS ] alias ] with_query_name [ [ AS ] alias ] } join: from_item [ join_type ] JOIN from_item [ { ON bool_expression | USING ( join_column [, ...] ) } ] join_type: { INNER | CROSS | FULL [OUTER] | LEFT [OUTER] | RIGHT [OUTER] }query_statement以可选的WITH子句开头,随后是一个query_expr;query_expr可以是单个select、带括号的查询表达式,或由集合运算符(set_op)连接的两个查询表达式,最后还可追加LIMIT/OFFSET。
SELECT 列表(SELECT list)
SELECT 列表的语法为:
SELECT [{ ALL | DISTINCT }] { [ expression. ]* | expression [ [ AS ] alias ] } [, ...]SELECT列表定义了查询将返回的列。列表中的表达式可以引用其对应FROM子句中任意from_item的列。列表中的每一项只能是以下三种形式之一:
*expressionexpression.*
SELECT *(select star)
SELECT *为完整查询执行后可见的每一列各产生一个输出列:
SELECT * FROM (SELECT 'apple' AS fruit, 'carrot' AS vegetable); +-------+-----------+ | fruit | vegetable | +-------+-----------+ | apple | carrot | +-------+-----------+SELECT expression
SELECT 列表中的项可以是表达式。这些表达式求值为单个值,并产生一个输出列,可附带显式alias。若表达式没有显式别名,则尽可能按"隐式别名"规则获得隐式别名;否则该列是匿名的,查询的其他部分无法按名称引用它。
SELECT expression.*
SELECT 列表项还可以是expression.*形式,为expression的每一列或顶层字段各产生一个输出列。这里的expression必须是表别名。例如下面的查询为别名g的表groceries的每一列产生一个输出列:
WITH groceries AS (SELECT 'milk' AS dairy, 'eggs' AS protein, 'bread' AS grain) SELECT g.* FROM groceries AS g; +-------+---------+-------+ | dairy | protein | grain | +-------+---------+-------+ | milk | eggs | bread | +-------+---------+-------+SELECT 修饰符(modifiers)
SELECT DISTINCT:丢弃重复行,只返回剩余行。注意SELECT DISTINCT不能返回 STRUCT 和 ARRAY 类型的列。SELECT ALL:返回所有行(包括重复行),这是SELECT的默认行为。
别名
SELECT 列表中别名的语法与可见性参见下文"别名"一节。
FROM 子句
FROM子句指明要从哪些表中检索行,并规定如何将这些行连接成单一数据流,供查询的其余部分处理。
语法
from_item: { table_name [ [ AS ] alias ] | join | ( query_expr ) [ [ AS ] alias ] | with_query_name [ [ AS ] alias ] }table_name:现有表的名称(可带限定符)。例如:SELECT * FROM Roster; SELECT * FROM beam.Roster;join:见下文 JOIN 类型 及 extensions/joins.md。select:
( select ) [ [ AS ] alias ]是一个表子查询(table subquery)。with_query_name:WITH子句中定义的查询名(见 WITH 子句)就像临时表名一样,可以在FROM子句的任何位置引用。例如subQ1、subQ2都是with_query_name:WITH subQ1 AS (SELECT * FROM Roster WHERE SchoolID = 52), subQ2 AS (SELECT SchoolID FROM subQ1) SELECT DISTINCT * FROM subQ2;在查询执行期间,
WITH子句会隐藏同名永久表,除非使用限定名(如beam.Roster)。
子查询(Subqueries)
子查询是出现在另一条语句内部、写在括号中的查询,也被称为 "sub-SELECT" 或 "nested SELECT"。完整的SELECT语法在子查询中均有效。子查询分两类:
- 表达式子查询(Expression Subqueries):可在任何允许表达式的位置使用,返回单个值;
- 表子查询(Table subqueries):只能用于
FROM子句,外层查询将子查询结果当作一张表。
两种子查询都必须用括号括起来。表子查询可选地带别名。
表达式子查询示例(子查询出现在FROM中作为表达式求值来源的普通用法):
SELECT AVG ( PointsScored ) FROM ( SELECT PointsScored FROM Stats WHERE SchoolID = 77 )带别名的表子查询示例:
SELECT r.LastName FROM ( SELECT * FROM Roster) AS r;FROM子句中别名的语法与可见性同样参见下文"别名"一节。
JOIN 类型
JOIN子句将两个from_item合并,使SELECT子句能将其作为单一数据源查询。join_type以及ON/USING子句(即"连接条件")规定了如何组合两个from_item的行、如何丢弃部分行,从而形成单一数据源。
join: from_item [ join_type ] JOIN from_item [ ON bool_expression | USING ( join_column [, ...] ) ] join_type: { INNER | CROSS | FULL [OUTER] | LEFT [OUTER] | RIGHT [OUTER] }所有JOIN子句都必须指定join_type。满足以下任一条件时,JOIN子句可以不写连接条件:
join_type为CROSS;- 两个
from_item中有一个或两个不是表(例如array_path或field_path)。
关于 Beam SQL 中 JOIN 的更多限制与流式语义,参见 extensions/joins.md。
[INNER] JOIN
INNER JOIN(或简写JOIN)在效果上计算两个from_item的笛卡尔积,并丢弃不满足连接条件的行。"效果上"意味着实际实现可以不必真正计算笛卡尔积。
从源码结构看,Beam 对 JOIN 的实现集中在 BeamJoinRel.java,其类注释将支持的连接场景划分为四类:有界表连接有界表、无界表连接无界表、有界表连接无界表、可查找表(BeamSeekableTable)连接非可查找表。在连接条件的约束上,extractJoinRexNodes与extractJoinPairOfRexNodes表明:Beam Calcite SQL 目前只支持equi-join(等值连接),ON条件必须能拆解为若干个基于=的比较(支持以AND组合);"非等值连接(Non equi-join)不受支持"会在规划期抛出UnsupportedOperationException。
CROSS JOIN
CROSS JOIN在 Beam Calcite SQL 中一般尚不支持(not yet supported)。这一点与上文语法图中的join_type同时呼应:源码 BeamJoinRel.java 中,当连接条件退化为字面量(如condition == true构成的 CROSS JOIN、condition == false的 JOIN ON FALSE)时,会直接抛出UnsupportedOperationException("CROSS JOIN, JOIN ON FALSE is not supported!")。
FULL [OUTER] JOIN
FULL OUTER JOIN(或简写FULL JOIN)返回两个from_item中满足连接条件的全部字段。
FULL表示两个from_item的所有行都会被返回,即使它们不满足连接条件。对于流式任务,返回的是按默认触发器判定不迟到、且在存在非全局窗口时属于同一窗口的所有行。OUTER表示若某一from_item中的行未与另一from_item中的任何行连接上,该行将返回,且另一from_item的所有列填充为 NULL。
LEFT [OUTER] JOIN
LEFT OUTER JOIN(或简写LEFT JOIN)始终保留JOIN子句中左侧from_item的所有行,即使右侧没有任何行满足连接谓词:
LEFT表示返回左侧from_item的全部行;若左侧某行未与右侧任何行连接上,该行右侧的所有列填充为 NULL。- 右侧
from_item中未与左侧任何行连接上的行被丢弃。
RIGHT [OUTER] JOIN
RIGHT OUTER JOIN(或简写RIGHT JOIN)与LEFT OUTER JOIN类似且对称。
ON 子句
ON子句包含一个bool_expression。两行组合后(连接两行的结果)只要bool_expression返回 TRUE,即满足连接条件。
示例:
SELECT * FROM Roster INNER JOIN PlayerStats ON Roster.LastName = PlayerStats.LastName;USING 子句
USING子句需要一个column_list,其中包含一个或多个同时出现在两张输入表中的列。它对该列执行相等比较,比较结果为 TRUE 时行即满足连接条件。大多数情况下,使用USING的语句等价于使用ON。例如:
SELECT FirstName FROM Roster INNER JOIN PlayerStats USING (LastName);等价于:
SELECT FirstName FROM Roster INNER JOIN PlayerStats ON Roster.LastName = PlayerStats.LastName;但使用SELECT *时,USING与ON的结果会不同:
下面的语句返回
Roster与PlayerStats中LastName相等的行,结果中只包含一个LastName列:SELECT * FROM Roster INNER JOIN PlayerStats USING (LastName);而下面的语句同样返回
LastName相等的行,但结果中包含两个LastName列(一个来自Roster,一个来自PlayerStats):SELECT * FROM Roster INNER JOIN PlayerStats ON Roster.LastName = PlayerStats.LastName;
JOIN 序列(Sequences of JOINs)
FROM子句可以按顺序包含多个JOIN子句。例如:
SELECT * FROM a LEFT JOIN b ON TRUE LEFT JOIN c ON TRUE;其中a、b、c是任意from_item。JOIN 默认从左到右绑定,但可以用括号改变分组顺序。
WHERE 子句
语法:
WHERE bool_expressionWHERE子句通过将每一行与bool_expression求值来过滤行,丢弃所有不返回 TRUE 的行(即返回 FALSE 或 NULL 的行)。
示例:
SELECT * FROM Roster WHERE SchoolID = 52;bool_expression可以包含多个子条件:
SELECT * FROM Roster WHERE LastName LIKE 'Mc%' OR LastName LIKE 'Mac%';注意:不能在WHERE子句中引用SELECT列表中的列别名。
INNER JOIN中的表达式在WHERE子句中存在等价形式。例如,使用INNER JOIN与ON的查询,等价于使用CROSS JOIN与WHERE的查询:
SELECT * FROM Roster INNER JOIN TeamMascot ON Roster.SchoolID = TeamMascot.SchoolID;等价于:
SELECT * FROM Roster CROSS JOIN TeamMascot WHERE Roster.SchoolID = TeamMascot.SchoolID;GROUP BY 子句
语法:
GROUP BY { expression [, ...] }GROUP BY子句将表中expression值非不同的行分组:对于源表中expression值非不同的多行,GROUP BY只产生一行合并结果。GROUP BY通常用于SELECT列表中出现聚合函数的情形,或用于消除输出冗余。
示例:
SELECT SUM(PointsScored), LastName FROM PlayerStats GROUP BY LastName;在 Beam 的实现中,GROUP BY与聚合由 BeamAggregationRel.java 负责翻译(其配套测试 BeamAggregationRelTest.java 可验证GROUP BY、聚合与窗口的组合行为)。当与流式窗口结合时,GROUP BY还需配合 Beam SQL 的窗口扩展,参见 extensions/windowing-and-triggering.md。
HAVING 子句
语法:
HAVING bool_expressionHAVING子句与WHERE子句类似:过滤掉对bool_expression求值不返回 TRUE 的行。与WHERE一样,bool_expression可以是任意返回布尔值的表达式,也可包含多个子条件。
HAVING与WHERE的差异在于:
HAVING要求查询中必须存在GROUP BY或聚合;HAVING发生在GROUP BY与聚合之后,因此它对结果集中的每个聚合行求值一次;而WHERE在GROUP BY与聚合之前求值。
HAVING子句可以引用FROM子句可用的列以及SELECT列表别名。HAVING中引用的表达式要么出现在GROUP BY子句中,要么必须是聚合函数的结果:
SELECT LastName FROM Roster GROUP BY LastName HAVING SUM(PointsScored) > 15;集合运算符(Set operators)
语法:
UNION { ALL | DISTINCT } | INTERSECT DISTINCT | EXCEPT DISTINCT集合运算符将两个或多个输入查询的结果合并为单一结果集。必须显式指定ALL或DISTINCT:指定ALL则保留所有行,指定DISTINCT则丢弃重复行。
若某行 R 在第一个输入查询中出现 m 次,在第二个输入查询中出现 n 次(m ≥ 0,n ≥ 0),则:
UNION ALL:R 在结果中出现 m + n 次;UNION DISTINCT:DISTINCT在UNION计算之后计算,因此 R 只出现 1 次;INTERSECT DISTINCT:DISTINCT在上述交集结果计算之后计算;EXCEPT DISTINCT:当 m > 0 且 n = 0 时,R 在输出中出现 1 次;- 当输入查询超过两个时,上述操作可推广,输出等价于从左到右增量合并输入的结果。
适用规则如下:
- 除
UNION ALL外的集合操作,所有列类型都必须支持相等比较; - 运算符两侧的输入查询必须返回相同数量的列;
- 运算符按各自
SELECT列表中列的位置配对两个输入查询返回的列,即第一个输入查询的第一列与第二个输入查询的第一列配对; - 结果集始终使用第一个输入查询的列名;
- 结果集始终使用对应列输入类型的超类型(supertypes),因此配对的列必须具有相同数据类型或共同超类型;
- 必须使用括号分隔不同的集合操作(此规则下
UNION ALL与UNION DISTINCT被视为不同的操作);如果语句只重复相同的集合操作,则无需括号。
示例(合法):
query1 UNION ALL (query2 UNION DISTINCT query3) query1 UNION ALL query2 UNION ALL query3示例(非法):
query1 UNION ALL query2 UNION DISTINCT query3 query1 UNION ALL query2 INTERSECT ALL query3; // INVALID.UNION
UNION运算符通过将每个查询结果集的列按位置配对并垂直拼接,合并两个或多个输入查询的结果集。
INTERSECT
INTERSECT运算符返回同时出现在左右两个输入查询结果集中的行。与EXCEPT不同,输入查询相对INTERSECT运算符的位置(左侧还是右侧)无关紧要。
EXCEPT
EXCEPT运算符返回出现在左侧输入查询中、但不在右侧输入查询中的行。
集合运算符的源码实现
Beam 将集合运算符统一实现在 BeamSetOperatorRelBase.java 中,它本身是一个PTransform<PCollectionList<Row>, PCollection<Row>>,通过OpType枚举区分UNION、INTERSECT、MINUS(对应 SQL 的EXCEPT),all布尔字段标记是否带ALL。其执行过程为:
- 校验左右两个输入的行窗口函数是否兼容,若不兼容直接抛出
IllegalArgumentException(这一点对应了文档中"输入两侧属于同一窗口"的流式语义约束); - 使用
CoGroup.join(By.fieldNames("*"))将左右输入按全部字段做分组连接; - 由
BeamSetOperatorsTransforms.SetOperatorFilteringDoFn(见 BeamSetOperatorsTransforms.java)按行 R 在左右两侧的出现次数 m、n 决定输出:
UNION ALL:输出 m + n 个副本;UNION DISTINCT:只输出 1 个副本;INTERSECT:两侧均大于 0 时输出(ALL时输出min(m, n)个,DISTINCT时输出 1 个);MINUS(EXCEPT):左侧有而右侧无时输出(ALL时输出 m 个,DISTINCT时输出 1 个);两侧都有时,EXCEPT ALL输出m - n个(大于 0 时),EXCEPT DISTINCT不输出。
这些语义与文档中的 m、n 规则完全一致。对应的测试用例见 BeamUnionRelTest.java(覆盖UNION去重与UNION ALL保留重复)、BeamIntersectRelTest.java、BeamMinusRelTest.java 与 BeamSetOperatorRelBaseTest.java(含两侧窗口一致的验证)。
LIMIT 子句与 OFFSET 子句
语法:
LIMIT count [ OFFSET skip_rows ]LIMIT指定一个非负的 INTEGER 类型count,最多返回count行;LIMIT 0返回 0 行。若存在集合操作,LIMIT在集合操作求值之后应用。OFFSET指定一个非负的 INTEGER 类型skip_rows,只有表中从该偏移量开始的行才会被考虑。- 这两个子句只接受字面量或参数值(literal or parameter values)。
LIMIT与OFFSET返回的具体行是未指定的(即不保证排序稳定性)。
从实现看,BeamSortRel.java 负责Sort节点的翻译:由于 Beam 不完整支持全局排序,它使用Top变换实现排序,因此仅支持"ORDER BY与LIMIT/OFFSET组合"的形态,例如:
SELECT * FROM t ORDER BY id DESC LIMIT 10; SELECT * FROM t ORDER BY id DESC LIMIT 10 OFFSET 5;而不带LIMIT的ORDER BY不被支持(会抛出异常)。相关测试见 BeamSortRelTest.java。
WITH 子句
WITH子句包含一个或多个命名的子查询,每当后续SELECT语句引用它们时都会执行。任何子句或子查询都可以引用WITH子句中定义的子查询,包括集合运算符(如UNION)两侧的任何SELECT语句。
示例:
WITH subQ1 AS (SELECT SchoolID FROM Roster), subQ2 AS (SELECT OpponentID FROM PlayerStats) SELECT * FROM subQ1 UNION ALL SELECT * FROM subQ2;别名(Aliases)
别名是查询中为表、列或表达式起的临时名称。可以在SELECT列表或FROM子句中引入显式别名,Beam 也会为某些表达式推断隐式别名。既无显式别名也无隐式别名的表达式是匿名的,查询无法按名称引用它。
显式别名语法
在FROM子句中,可以为任何项(包括表、数组、子查询和UNNEST子句)使用[AS] alias引入显式别名,AS关键字可选:
SELECT s.FirstName, s2.SongName FROM Singers AS s JOIN Songs AS s2 ON s.SingerID = s2.SingerID;在SELECT列表中,可以为任何表达式使用[AS] alias引入显式别名,AS关键字同样可选:
SELECT s.FirstName AS name, LOWER(s.FirstName) AS lname FROM Singers s;显式别名可见性
在查询中引入显式别名后,其可被引用的位置受 Beam 的名称作用域规则限制。
FROM 子句别名
Beam 从左到右处理FROM子句中的别名,且别名只对后续的JOIN子句可见。
歧义别名
若一个名称可解析为多个不同对象(即名称歧义),Beam 会报错。例如下面的查询中Singers与Songs都有SingerID列,导致列名冲突:
SELECT SingerID FROM Singers, Songs;隐式别名
在SELECT列表中,若表达式没有显式别名,Beam 按以下规则分配隐式别名(SELECT列表中可以存在多个同名列):
- 标识符:别名即标识符本身。例如
SELECT abc隐含AS abc; - 路径表达式:别名为路径中最后一个标识符。例如
SELECT abc.def.ghi隐含AS ghi; - 使用点号成员访问运算符的字段访问:别名为字段名。例如
SELECT (struct_function()).fname隐含AS fname。
其余情况没有隐式别名,该列匿名、无法按名称引用。该列的数据仍会被返回,显示结果可能带有生成的标签,但该标签不能当作别名使用。
在FROM子句中,from_item不要求必须有别名。无显式别名时的隐式别名规则:
- 标识符:别名即标识符。例如
FROM abc隐含AS abc; - 路径表达式:别名为路径最后一个标识符。例如
FROM abc.def.ghi隐含AS ghi。
表子查询没有隐式别名;FROM UNNEST(x)也没有隐式别名。
结合 Beam 使用查询语句
Beam Calcite SQL 查询语句最终通过sdks/java/extensions/sql模块被编译进 Beam 管线。你可以通过 Beam SQL 扩展的 Java API(beam-sdks-java-extensions-sql)、JDBC 驱动(JdbcDriver.java)或 shell.md 中描述的 SQL Shell 交互式执行这些查询。无论哪种方式,查询最终都会被翻译为 Beam 的PTransform图,从而可以在 Direct Runner、Dataflow、Flink、Spark 等任意 Beam Runner 上运行,这也是"一套 SQL 同时服务批处理与流处理"的关键所在。
在流式场景下,请务必结合 extensions/joins.md 与 extensions/windowing-and-triggering.md 理解 JOIN 的窗口约束与触发语义,它们与本文的语法规则共同构成 Beam Calcite SQL 的完整查询能力。
说明:本页部分内容基于 Google 的 BigQuery 标准 SQL 查询语法文档改写,并依据 Creative Commons 3.0 Attribution License 使用。
- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam Calcite SQL 查询语法完全指南:SELECT、JOIN、聚合与集合运算实战解析
Apache Beam Calcite SQL 查询语法完全指南:SELECT、JOIN、聚合与集合运算实战解析 导读 Apache Beam 提供基于 Apa
大数据批处理流处理数据工程Apache Beam Calcite SQL 完全指南:默认 SQL 方言的语法、数据类型、函数与扩展能力解析
Apache Beam Calcite SQL 完全指南:默认 SQL 方言的语法、数据类型、函数与扩展能力解析 本文以 Apache Beam 官方文档 Be
批处理流处理大数据Apache Beam ZetaSQL Query Syntax 查询语法完全参考:SELECT、FROM、JOIN、GROUP BY 与别名规则详解
Apache Beam ZetaSQL Query Syntax 查询语法完全参考:SELECT、FROM、JOIN、GROUP BY 与别名规则详解 本文是
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考