
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载本篇指南基于 Flink 官方文档 Lateral View Clause 展开讲解如何在 Flink 的 Hive 方言中编写LATERAL VIEW查询包括语法与参数、OUTER关键字的补行语义、多个 Lateral View 的级联用法以及explode()拆分行数组的完整示例读完后你还能了解该语法在 Flink SQL 解析器与 Calcite 规划器中的落地位置以及 Flink 默认方言中等价的LATERAL TABLE/UNNEST写法。Lateral View 是什么Lateral view子句与用户自定义表生成函数UDTF配合使用典型代表就是explode()。与标量 UDF 不同一个 UDTF 对每一行输入会生成零行或多行输出因此它无法直接出现在SELECT列表中而必须作为表参与 FROM 子句——这正是 Lateral View 解决的核心问题。按照文档描述Lateral View 的处理过程是先将 UDTF 逐行应用于基表base的每一个行再把 UDTF 输出的行与对应的输入行做连接join用给定的表别名table alias形成一个虚拟表供查询引用。也就是说Lateral View 在逻辑上等价于基表 × 每个基表行展开后的 UDTF 结果表的一次连接。这一语义与标准 SQL 中的 Lateral Join相关连接同源从源码结构看Flink 的 Calcite 校验器 SqlValidatorImpl 在校验 FROM 子句时就携带lateral标记允许左侧表项参与右侧表达式求值SqlToRelConverter 则负责把 lateral join 的谓词下推到内部子节点完成从 SQL 关系到物理计划节点的转换。Lateral View 属于 Hive 方言查询语法的一部分同系列的文档还包括 Join、CTE、子查询、窗口函数 等。使用前提如何启用 Hive 方言LATERAL VIEW是 Hive 语法默认方言下不可用。根据 Hive Dialect 总览Flink 支持default与hive两种 SQL 方言且可以按语句动态切换无需重启会话。通过 SQL Client 切换方言SQL 方言由table.sql-dialect会话属性指定Flink SQL SET table.sql-dialect hive; -- 切换到 Hive 方言 [INFO] Session property has been set. Flink SQL SET table.sql-dialect default; -- 切回 Flink 默认方言 [INFO] Session property has been set.通过 SQL GatewayHiveServer2 Endpoint使用配置了 HiveServer2 Endpoint 的 SQL Gateway 时方言默认为 Hive 方言无需额外设置如需切回默认方言同样执行SET table.sql-dialect default;。通过 Table API 设置// Java EnvironmentSettings settings EnvironmentSettings.inStreamingMode(); TableEnvironment tableEnv TableEnvironment.create(settings); // 使用 Hive 方言 tableEnv.getConfig().setSqlDialect(SqlDialect.HIVE); // 使用默认方言 tableEnv.getConfig().setSqlDialect(SqlDialect.DEFAULT);# Python from pyflink.table import * settings EnvironmentSettings.in_batch_mode() t_env TableEnvironment.create(settings) # 使用 Hive 方言 t_env.get_config().set_sql_dialect(SqlDialect.HIVE)启用 Hive 方言还需要满足以下前提引自 Hive Dialect 总览文档需引入与 Hive 相关的依赖见文档中 Hive 依赖说明当前 catalog 必须是 HiveCatalog否则会回退到 Flink 的default方言使用 HiveServer2 Endpoint 的 SQL Gateway 时当前 catalog 默认为 HiveCatalog建议加载 HiveModule 并置于模块列表首位以便函数解析时优先命中 Hive 内置函数使用 HiveServer2 Endpoint 的 SQL Gateway 时会自动加载Hive 方言仅支持两段式标识符2-part identifier无法在标识符中指定 catalog具体特性是否可用取决于你所使用的 Hive 版本Hive 方言主要用于批处理模式部分 Hive 语法如 Sort/Cluster/Distributed BY、Transform尚未在流模式下支持。LATERAL VIEW 语法文档给出的完整语法定义如下lateralView: LATERAL VIEW [ OUTER ] udtf( expression ) tableAlias AS columnAlias [, ... ] fromClause: FROM baseTable lateralView [, ... ]各部分含义组成部分说明LATERAL VIEW子句关键字声明对基表应用 UDTF[ OUTER ]可选控制 UDTF 无输出时是否保留基表行见下文udtf( expression )UDTF 调用及其参数expression可引用基表任意左侧表中的列tableAlias为 UDTF 结果生成的虚拟表命名AS columnAlias [, ...]UDTF 输出列的别名列表可以省略此时别名继承自 UDTF 返回的 StructObjectInspector 的字段名参数详解OUTER 关键字文档对OUTER的说明如下用户可以指定可选的OUTER关键字使得在LATERAL VIEW通常不产生行的情况下也能生成行。这种情况发生在所使用的 UDTF 未产生任何行时——最常见的场景就是对一个空数组做explode()。此时源行source row将不会出现在结果中使用OUTER可以防止这一点对应行会以 UDTF 各列填NULL的方式保留在结果里。从连接语义上看这正是两种 Join 的差异不带OUTER的LATERAL VIEW是内连接语义UDTF 无输出则基表行被丢弃带OUTER则是左外连接语义基表行恒保留UDTF 列补NULL。多个 Lateral View一个 FROM 子句中可以出现多个LATERAL VIEW子句并且后续的 LATERAL VIEW 可以引用该 LATERAL VIEW 左侧出现的任意表中的列。这意味着你可以对同一基表的不同数组列分别展开见下文示例二让第二个 UDTF 引用第一个 Lateral View 产生的虚拟表列形成级联展开。完整示例示例一用 explode() 将数组列拆成多行假设有一张表CREATE TABLE pageAds(pageid string, addid_list arrayint);表中包含两行数据front_page, [1, 2, 3]; contact_page, [3, 4, 5];使用LATERAL VIEW把列addid_list拆成独立行SELECT pageid, adid FROM pageAds LATERAL VIEW explode(adid_list) adTable AS adid; -- 结果 front_page, 1 front_page, 2 front_page, 3 contact_page, 3 contact_page, 4 contact_page, 5执行过程对应前文的语义描述explode()对每行输入产生 3 行输出数组有 3 个元素再与基表的 2 行分别连接最终得到 6 行虚拟表adTable只有一列别名为adid。示例二同一 FROM 子句中多个 Lateral View如果有一张表CREATE TABLE t1(c1 arrayint, c2 arrayint);可以使用多个 Lateral View 子句把列c1和c2同时拆成独立行SELECT myc1, myc2 FROM t1 LATERAL VIEW explode(c1) myTable1 AS myc1 LATERAL VIEW explode(c2) myTable2 AS myc2;这里第二个LATERAL VIEW仍然作用于基表t1的c2列两个虚拟表myTable1、myTable2均可在SELECT中通过各自列别名引用。示例三UDTF 无输出时用 OUTER 保留基表行当 UDTF 不产生任何行时普通的LATERAL VIEW也不会产生行基表行被内连接语义丢弃。使用LATERAL VIEW OUTER则仍会生成行UDTF 对应列以NULL填充SELECT * FROM t1 LATERAL VIEW OUTER explode(array()) C AS a;该查询中explode(array())对每个基表行都返回 0 行由于使用了OUTER基表行仍会出现在结果中列a为NULL。源码印证Lateral 语义在 Flink 中的实现解析器LATERAL 关键字与表函数调用Flink SQL 的解析器模板 Parser.jj 定义了 FROM 子句中 Lateral 表引用table reference的解析分支| LOOKAHEAD(2) [ LATERAL { lateral true; } ] tableRef ParenthesizedExpression(exprContext) // LATERAL (子查询) ... | LOOKAHEAD(2) [ LATERAL ] // LATERAL is implicit with UNNEST, so ignore UNNEST ... // UNNEST(...) ... | [ LATERAL { lateral true; } ] tableRef TableFunctionCall() // LATERAL TABLE(udtf(...)) tableRef addLateral(tableRef, lateral)其中LATERAL关键字本身定义在 Parser.jj 的词法部分。可以看出在 Flink默认方言中与 HiveLATERAL VIEW udtf(...) alias AS col语义等价的写法是LATERAL TABLE(udtf(...)) AS alias(col)UNNEST展开数组/多集列时LATERAL可省略隐式 lateral而显式书写LATERAL时会被包装进SqlStdOperatorTable.LATERAL调用addLateral逻辑规划阶段由 Calcite 的 SqlToRelConverter 完成 lateral join 的转换与谓词下推最终物化为 Correlate相关连接类算子。测试用例Lateral 表函数的运行时验证Flink 的集成测试 CorrelateITCase 覆盖了大量 Lateral 相关查询例如SELECT * FROM LATERAL TABLE(str_split(Jack,John, ,)) as T0(d); SELECT * FROM T1, LATERAL TABLE(str_split(c, ,)) as T2(s);第二例正是基表列驱动 UDTF 逐行展开的标准形态——与 Hive 方言中LATERAL VIEW explode(col)的用法在语义上一一对应。Hive 方言 vs 默认方言写法对照语义Hive 方言table.sql-dialect hiveFlink 默认方言展开数组列为多行LATERAL VIEW explode(col) t AS xLATERAL TABLE(udtf(...)) AS t(x)/UNNESTUDTF 无输出时保留基表行LATERAL VIEW OUTER explode(col) t AS xLEFT JOIN LATERAL /UNNEST的 LEFT JOIN 形式多列级联展开多个LATERAL VIEW后者可引用左侧所有表多个 Lateral 表引用参与 FROM 连接小结LATERAL VIEW [OUTER] udtf(expression) tableAlias AS columnAlias...是 Hive 方言下将 UDTF如explode()与基表逐行连接、把多值列拆成行的标准写法列别名可省略此时继承自 UDTF 返回的 StructObjectInspector 字段名OUTER决定 UDTF 无输出时的行为默认丢弃基表行内连接语义加OUTER则基表行保留且 UDTF 列补NULL左外连接语义一个 FROM 子句支持多个LATERAL VIEW后续子句可引用其左侧任意表的列使用该语法前需按 Hive Dialect 总览 的要求启用 Hive 方言table.sql-dialect hive、当前 catalog 为 HiveCatalog、建议加载 HiveModule并注意该方言主要面向批处理场景从源码看lateral 表引用的解析在 Parser.jj语义校验与计划转换分别在 SqlValidatorImpl 与 SqlToRelConverter 中完成运行时行为可由 CorrelateITCase 中的 Lateral 表函数查询验证。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Apache Spark SQL LATERAL VIEW 子句完全指南语法、OUTER 语义与执行计划剖析Apache Spark SQL LATERAL VIEW 子句完全指南语法、OUTER 语义与执行计划剖析 LATERAL VIEW 是 Apache Sp大数据数据分析批处理流处理机器学习图计算Flink Hive 方言 CREATE 语句完全指南DATABASE / TABLE / VIEW / MACRO / FUNCTION 语法与实现原理Flink Hive 方言 CREATE 语句完全指南DATABASE / TABLE / VIEW / MACRO / FUNCTION 语法与实现原理 F大数据流处理批处理数据工程Flink Hive 方言 LOAD DATA 语句完全指南语法、参数与源码实现剖析Flink Hive 方言 LOAD DATA 语句完全指南语法、参数与源码实现剖析 LOAD DATA 是 Flink Hive 方言中用于将用户指定目录或大数据流处理批处理数据工程上一篇如何快速部署Sheepdog分布式存储系统面向QEMU用户的完整指南下一篇OpenUSD架构解密从MaterialX材质到Hydra渲染的完整技术栈创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考