ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Apache Spark SQL Pipe Syntax 完全指南:使用 `|>` 链式组合查询操作符

Apache Spark SQL Pipe Syntax 完全指南:使用 `|>` 链式组合查询操作符 Apache Spark SQL Pipe Syntax 完全指南使用|链式组合查询操作符【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/sparkSQL Pipe Syntax 是 Apache Spark SQL 提供的声明式查询组合语法它允许开发者使用|管道符把一个查询的输出依次传递给后续操作符从而以数据流管道的方式逐步完成过滤、投影、聚合、连接、排序、采样与行列转置等变换。本指南以 docs/sql-pipe-syntax.md 为骨架结合sql/api与sql/catalyst中的 ANTLR 文法、AstBuilder解析实现与SQLConf配置完整讲解 Pipe Syntax 的语法规则、全部受支持操作符、与标准 SQL 的等价关系及互操作方式读者学完后可以立即在 Spark SQL 中编写、调试与迁移管道式查询。语法总览OverviewApache Spark 支持 SQL Pipe Syntax允许通过组合操作符来编写查询。其核心语法约定如下任何查询都可以在末尾追加零个或多个管道操作符以管道字符|分隔每个管道操作符以一个或多个 SQL 关键字开头后跟各自的专属语法详见下文操作符表大多数操作符复用标准 SQL 子句的既有文法操作符可以以任意顺序、任意次数应用。与此同时FROM tableName现在也是一个受支持的独立查询语句行为与TABLE tableName完全一致。这为启动一条管道式查询提供了便捷的起点当然也可以在任何合法的 Spark SQL 查询末尾追加一个或多个管道操作符其行为保持一致。一个直观的对比TPC-H Query 13以 TPC-H 基准中的第 13 号查询为例传统的标准 SQL 写法为SELECT c_count, COUNT(*) AS custdist FROM (SELECT c_custkey, COUNT(o_orderkey) c_count FROM customer LEFT OUTER JOIN orders ON c_custkey o_custkey AND o_comment NOT LIKE %unusual%packages% GROUP BY c_custkey ) AS c_orders GROUP BY c_count ORDER BY custdist DESC, c_count DESC;使用 SQL Pipe Operators 表达同样的逻辑就变成了自左向右、自顶向下的线性管道FROM customer | LEFT OUTER JOIN orders ON c_custkey o_custkey AND o_comment NOT LIKE %unusual%packages% | AGGREGATE COUNT(o_orderkey) c_count GROUP BY c_custkey | AGGREGATE COUNT(*) AS custdist GROUP BY c_count | ORDER BY custdist DESC, c_count DESC;可以看到嵌套的子查询被展开成了逐步的数据流变换阅读顺序与执行顺序一致层次更加扁平、可读性更强。解析层面的实现依据从语法树看 Pipe SyntaxPipe Syntax 并不是一条伪语法糖它在 Spark 的 SQL 解析器中是一等公民。在 SqlBaseParser.g4 中可以看到对应的 ANTLR 文法规则| leftqueryTerm OPERATOR_PIPE operatorPipeRightSide #operatorPipeStatement | leftqueryTerm {isOperatorPipeStart()}? PIPE operatorPipeRightSide #operatorPipeStatement即operatorPipeStatement允许在任意queryTerm查询项之后接一个管道标记与operatorPipeRightSide管道操作符右部。文法中还提供了一个前瞻判断函数isOperatorPipeStart()int la _input.LA(2); // Look ahead 2 tokens (current is PIPE, check what follows) return la SELECT || la EXTEND || la SET || la DROP || la AS || la WHERE || la PIVOT || la UNPIVOT || ...它的作用是消解|的歧义|在表达式里是位或bitwise OR而在查询级可以是管道操作符。该函数通过向前看一个 token 来判断紧随|之后的关键字是否是管道操作符关键字从而在不破坏既有表达式语义的前提下兼容单字符|。同时词法层面在 SqlBaseLexer.g4 中定义了CONCAT_PIPE即|等 token。在 AST 构建阶段AstBuilder.scala 会把管道表达式包装为PipeExpression节点例如PipeExpression(node.child, isAggregate false, PipeOperators.selectClause)每个操作符会基于其输入关系前一个操作符产生的行集生成新的表达式树从而在逻辑计划层面真正实现上一个操作符的输出是下一个操作符的输入这一语义。相关配置项Pipe Syntax 的启用与单字符管道符支持由 SQLConf.scala 中的两个配置控制spark.sql.operatorPipeSyntaxEnabled布尔开关开启后 Spark SQL 使用|标记来表示 SQL 子句之间的分隔从而以数据管道的方式描述查询配置说明原文This uses the operator pipe marker|to indicate separation between clauses of SQL in a manner that describes the flow of dataspark.sql.parser.singleCharacterPipeOperator.enabled当其为true时单字符管道 token|可以作为|的替代用于 SQL 管道操作符为false时仅|被识别为管道操作符|仍然按位或运算符解释。上述配置可通过SET语句或spark-defaults.conf调整具体示例SET spark.sql.operatorPipeSyntaxEnabled true;起始查询Source Tables要开始一条使用 Pipe Syntax 的新查询使用FROM tableName或TABLE tableName子句它创建一个由源表全部行构成的关系然后可以在其后追加一个或多个管道操作符继续变换FROM tableNameTABLE tableName返回源表的全部输出行不做任何修改。例如CREATE TABLE t AS VALUES (1, 2), (3, 4) AS t(a, b); TABLE t; ------ | a| b| ------ | 1| 2| | 3| 4| ------投影类操作ProjectionsPipe Syntax 在投影方面提供了高度可组合的表达式求值方式。一个重要优势是它支持基于前一步计算出来的新表达式做增量计算。由于每个操作符都独立作用于其输入表因此无需横向列引用lateral column references无论操作符以何种顺序出现每个操作符都作用于输入表而每个新计算的列都会对后续操作符可见。SELECT求值表达式生成新表| SELECT expr [[AS] alias], ...对输入表的每一行求值给定表达式。一般情况下Pipe Syntax 并不总是需要SELECT操作符它通常用于查询末尾附近来求值表达式或指定输出列列表。由于最终查询结果总是由最后一个管道操作符返回的列构成当查询中没有出现SELECT时输出将包含整行的所有列这与标准 SQL 中的SELECT *行为类似。可以使用DISTINCT和*其语义相当于常规 Spark SQL 中表子查询最外层的SELECT。SELECT列表中同样支持窗口函数此时必须提供OVER子句也可以在WINDOW子句中给出窗口定义。需要注意的是该操作符不支持聚合函数需要聚合请使用AGGREGATE操作符。CREATE TABLE t AS VALUES (0), (1) AS t(col); FROM t | SELECT col * 2 AS result; ------ |result| ------ | 0| | 2| ------EXTEND追加新列| EXTEND expr [[AS] alias], ...通过逐行求值指定表达式向输入表追加新列。在EXTEND之后顶层列名会更新但表别名仍指向原始行值例如两个表lhs、rhs内连接后进行EXTEND再SELECT lhs.col, rhs.col仍能正确引用。它等价于常规 Spark SQL 中的SELECT *, new_column。VALUES (0), (1) tab(col) | EXTEND col * 2 AS result; --------- |col|result| --------- | 0| 0| | 1| 2| ---------SET按表达式替换列值| SET column expression, ...用指定表达式的求值结果替换输入表的列。每个被引用的列必须在输入表中恰好出现一次。它类似于常规 Spark SQL 的SELECT * EXCEPT (column), expression AS column。可以在单个SET子句中执行多次赋值每次赋值都可以引用此前赋值的结果。赋值后顶层列名更新但表别名仍指向原始行值。VALUES (0), (1) tab(col) | SET col col * 2; --- |col| --- | 0| | 2| ---DROP按列名删除列| DROP column, ...按名称删除输入表的列每个被引用的列必须在输入表中恰好出现一次。它类似于常规 Spark SQL 的SELECT * EXCEPT (column)。同样DROP之后顶层列名更新但表别名仍指向原始行值。VALUES (0, 1) tab(col1, col2) | DROP col1; ---- |col2| ---- | 1| ----AS为输入表引入新别名| AS alias保留输入表的行与列名但为表引入新的别名该别名可在后续操作符中被引用原有的表别名被新别名替换。在通过SELECT或EXTEND添加新列之后或在AGGREGATE聚合之后使用该操作符非常有用它让后续JOIN操作符引用列时更简单查询也更易读。VALUES (0, 1) tab(col1, col2) | AS new_tab | SELECT col1 col2 FROM new_tab; ----------- |col1 col2| ----------- | 1| -----------聚合操作AGGREGATE总体而言Pipe Syntax 中的聚合方式与常规 Spark SQL 不同。AGGREGATE是执行聚合的唯一途径其语法有两种形式-- 全表聚合 | AGGREGATE agg_expr [[AS] alias], ... -- 分组聚合 | AGGREGATE [agg_expr [[AS] alias], ...] GROUP BY grouping_expr [AS alias], ...全表聚合不带GROUP BY时对整张输入表聚合输出表返回一行每个聚合表达式一列分组聚合带GROUP BY时对每个唯一的分组表达式取值组合返回一行。输出表先包含求值后的分组表达式列再包含聚合函数列。agg_expr可以包含任意 Spark SQL 支持的聚合函数COUNT、SUM、AVG、MIN等也可以在其上或下叠加其他表达式例如MIN(FLOOR(col)) 1每个agg_expr必须至少包含一个聚合函数否则查询报错。每个agg_expr可以通过AS alias起别名也可以使用DISTINCT关键字在聚合前去重例如COUNT(DISTINCT col)。GROUP BY子句可以包含任意数量的分组表达式。分组表达式支持直接赋别名以便后续操作符引用——这就意味着无需在GROUP BY与SELECT之间重复整个表达式因为AGGREGATE是一个同时完成分组与聚合的单一操作符。关于序数ordinal的一个关键差异GROUP BY表达式支持从 1 开始的序数。与常规 SQL 中序数指向伴随SELECT子句中的表达式不同在 Pipe Syntax 中序数指向前一个操作符产生的关系的列。例如在TABLE t | AGGREGATE COUNT(*) GROUP BY 2中2引用的是输入表t的第二列。同理在AGGREGATE之后通常也无需再追加SELECT因为它已经一次性返回了分组列与聚合列。-- 全表聚合 VALUES (0), (1) tab(col) | AGGREGATE COUNT(col) AS count; ----- |count| ----- | 2| ----- -- 分组聚合 VALUES (0, 1), (0, 2) tab(col1, col2) | AGGREGATE COUNT(col2) AS count GROUP BY col1; --------- |col1|count| --------- | 0| 2| ---------其他转换操作Other Transformations其余操作符用于过滤、连接、排序、采样与集合运算等转换总体上与常规 Spark SQL 的行为一致。WHERE条件过滤| WHERE condition返回满足条件的输入行子集。由于该操作符可以出现在管道中的任意位置无需单独的HAVING或QUALIFY语法——聚合前过滤用WHERE聚合后过滤直接再追加一个WHERE即可。VALUES (0), (1) tab(col) | WHERE col 1; --- |col| --- | 1| ---LIMIT 与 OFFSET限制行数| [LIMIT n] [OFFSET m]返回指定数量的输入行并保持已有排序如果有。LIMIT与OFFSET可以同时使用也可以各自单独使用LIMIT不带OFFSET或OFFSET不带LIMIT。VALUES (0), (0) tab(col) | LIMIT 1; --- |col| --- | 0| ---JOIN连接两路输入| [LEFT | RIGHT | FULL | CROSS | SEMI | ANTI | NATURAL | LATERAL] JOIN table [ON condition | USING(col, ...)]连接来自两侧输入的行返回管道输入表与JOIN关键字后的表表达式的过滤叉积。其行为与常规 SQL 的JOIN子句类似其中管道操作符的输入表成为连接的左端表参数成为连接的右端。LEFT、RIGHT、FULL等标准连接修饰符均受支持。当连接谓词需要同时引用两侧输入、且两侧存在同名列时需要借助表别名来区分。此时AS操作符非常有用——它可以为作为连接左端的管道输入表引入新别名右端表参数则使用标准语法赋别名。SELECT 0 AS a, 1 AS b | AS lhs | JOIN VALUES (0, 2) rhs(a, b) ON (lhs.a rhs.a); ------------ | a| b| c| d| ------------ | 0| 1| 0| 2| ------------ VALUES (apples, 3), (bananas, 4) t(item, sales) | AS produce_sales | LEFT JOIN (SELECT apples AS item, 123 AS id) AS produce_data USING (item) | SELECT produce_sales.item, sales, id; ---------------------- | item | sales | id | ---------------------- | apples | 3 | 123 | | bananas | 4 | NULL | ----------------------ORDER BY排序| ORDER BY expr [ASC | DESC], ...按要求对输入行排序后返回支持标准修饰符包括NULLS FIRST/NULLS LAST。VALUES (0), (1) tab(col) | ORDER BY col DESC; --- |col| --- | 1| | 0| ---UNION、INTERSECT、EXCEPT集合运算| {UNION | INTERSECT | EXCEPT} {ALL | DISTINCT} (query)对输入表或子查询与集合运算参数的行做并、交、差等集合运算。VALUES (0), (1) tab(a, b) | UNION ALL VALUES (2), (3) tab(c, d); ------- | a| b| ------- | 0| 1| | 2| 3| -------说明上述示例来自官方文档原文。读者在实际编写时建议让VALUES子句的列数与最终列结构保持一致以获得确定的输出列。TABLESAMPLE采样| TABLESAMPLE method(size {ROWS | PERCENT})返回由给定采样算法选出的输入行子集。VALUES (0), (0), (0), (0) tab(col) | TABLESAMPLE (1 ROWS); --- |col| --- | 0| --- VALUES (0), (0) tab(col) | TABLESAMPLE (100 PERCENT); --- |col| --- | 0| | 0| ---PIVOT行转列| PIVOT (agg_expr FOR col IN (val1, ...))返回一个新表将输入行透视pivot为列。VALUES (dotNET, 2012, 10000), (Java, 2012, 20000), (dotNET, 2012, 5000), (dotNET, 2013, 48000), (Java, 2013, 30000) courseSales(course, year, earnings) | PIVOT ( SUM(earnings) FOR COURSE IN (dotNET, Java) ) ---------------- |year|dotNET| Java| ---------------- |2012| 15000| 20000| |2013| 48000| 30000| ----------------UNPIVOT列转行| UNPIVOT (value_col FOR key_col IN (col1, ...))返回一个新表将输入列反透视unpivot为行。VALUES (dotNET, 2012, 10000), (Java, 2012, 20000), (dotNET, 2012, 5000), (dotNET, 2013, 48000), (Java, 2013, 30000) courseSales(course, year, earnings) | UNPIVOT ( earningsYear FOR year IN (2012, 2013, 2014) ---------------------- | course| year|earnings| ---------------------- | Java| 2012| 20000| | Java| 2013| 30000| | dotNET| 2012| 15000| | dotNET| 2013| 48000| | dotNET| 2014| 22500| ----------------------说明UNPIVOT 示例在官方文档原文中括号未闭合实际书写时需要补齐)并确保IN列表中的列名与输入表的列对应。独立性与互操作性Independence and InteroperabilityPipe Syntax 与既有 SQL 查询不存在向后兼容性问题任何查询都可以用常规 Spark SQL、Pipe Syntax或两者的组合来编写。因此以下不变式始终成立每个管道操作符接收一个输入表并且无论输入表是如何计算出来的它对其行都执行相同的操作前缀闭包对任意合法的 N 个管道操作符链前 MM N个操作符构成的任何子集也构成一条合法查询。这一性质对检查与调试非常有用——例如可以在 Jupyter 等 SQL 编辑器中选中部分行使用run highlighted text运行高亮文本功能逐段验证可附加性可以给任何用常规 Spark SQL 编写的合法查询追加管道操作符。启动 Pipe Syntax 查询的规范方式是用FROM tableName子句注意它本身是一条合法独立查询可以被任何其他 Spark SQL 查询替换而不失一般性子查询互通表子查询既可以用常规 Spark SQL 语法编写也可以用 Pipe Syntax 编写并且可以出现在用任意一种语法编写的外层查询中语句互通其他 Spark SQL 语句如视图、DDL、DML 命令也可以包含用任意一种语法编写的查询。受支持的管道操作符总表下表汇总了所有受支持的管道操作符及其输出行语义。注意每个操作符都接收一个输入关系该关系由|符号之前的查询生成的行构成。操作符输出行FROM 或 TABLE原样返回源表全部输出行SELECT对输入表的每一行求值指定表达式EXTEND对每个输入行求值指定表达式向输入表追加新列SET用指定表达式的求值结果替换输入表的列DROP按名称删除输入表的列AS保留输入表的行与列名但赋予新的表别名WHERE返回满足条件的输入行子集LIMIT返回指定数量的输入行保持已有排序如有AGGREGATE带或不带分组地执行聚合JOIN连接两侧输入的行返回输入表与表参数的过滤叉积ORDER BY按要求排序后返回输入行UNION ALL对输入表与其他表参数的行做并集或其他集合运算TABLESAMPLE返回给定采样算法选出的输入行子集PIVOT返回将输入行透视pivot为列的新表UNPIVOT返回将输入列反透视unpivot为行的新表小结SQL Pipe Syntax 把传统嵌套的 SQL 查询重写为线性、自左向右的数据流管道显著提升了复杂查询尤其是多步聚合、连接与列变换组合的可读性与可维护性同时它与常规 Spark SQL 完全互操作既支持在既有查询后追加管道操作符也支持在子查询、视图、DDL/DML 中混用两种语法。从 SqlBaseParser.g4 的文法规则到 AstBuilder.scala 中PipeExpression的构建再到 SQLConf.scala 中的开关配置整个特性链路完整可查。若希望在环境中启用或调整该语法请确认spark.sql.operatorPipeSyntaxEnabled与spark.sql.parser.singleCharacterPipeOperator.enabled两个配置项符合你的预期。【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表