
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载本指南基于 Flink Table/SQL 模块官方文档 docs/content/docs/dev/table/sql/queries/joins.md 展开系统梳理 Flink SQL 在动态表Dynamic Table上提供的全部 Join 类型Regular Join、Interval Join、Temporal Join事件时间 / 处理时间 / 时态表函数、Lookup Join、数组展开与表函数 Join。读者将掌握每种 Join 的适用语义、SQL 语法、状态开销与性能取舍并能结合仓库源码如 StreamExecTemporalJoin.java理解其底层执行原理在实际流批作业中做出正确的 Join 选型。总览Flink SQL 对动态表的 Join 支持Flink SQL 支持在动态表上执行复杂而灵活的 Join 操作。由于输入数据是持续到达的流Join 的语义远比静态批处理丰富。按语义与状态需求可分为如下几大类Join 类型输入要求状态开销典型场景Regular Join任意更新型Insert/Update/Delete输入两侧状态永久保留可能无限增长一般流关联、双流 JOINInterval Join仅 append-only 且带时间属性按时间约束清理状态有界订单与物流按时间窗口关联Temporal Join事件时间probe 侧任意、build 侧为版本表保留自 watermark 以来的版本按历史时点做汇率换算、CDC 表关联Temporal Join处理时间probe 侧 append-only、build 侧为维度表不存旧行仅存当前版本维表关联、实时补全Lookup Join一侧需 lookup source 连接器无状态外部查询JDBC / HBase 维表补全Array Expansion / Table Function Joinappend-only取决于函数实现数组展开、UDTF 横向扩展一个重要的全局前提默认情况下Flink 不会优化 Join 的顺序表按照FROM子句中指定的顺序进行 Join。为了提升性能建议把更新频率最低的表放在前面、更新频率最高的表放在最后。同时要避免写出产生交叉 Join笛卡尔积的查询——此类查询不受支持会导致作业失败。Regular Join通用双流关联Regular Join 是最通用的 Join 类型任意一侧的新记录或变更都会影响整个 Join 结果的全集。例如当左侧出现一条新记录时它会与右侧所有历史及未来的、满足product id相等的记录进行 Join。SELECT * FROM Orders INNER JOIN Product ON Orders.productId Product.id对于流式查询Regular Join 的语法最为灵活允许任意更新类型insert、update、delete的输入表。但其代价是操作层面的重要限制它要求将 Join 两侧的输入永久保存在 Flink 状态中。因此计算查询结果所需的状态量可能随所有输入表的不同输入行数以及中间 Join 结果的数量无限增长。用 State TTL 控制状态膨胀你可以通过查询配置为状态设置合适的生命周期Time-to-Live, TTL来防止状态过度膨胀。对应配置项定义于 ExecutionConfigOptions.java配置键table.exec.state.ttl类型Duration如1 h、30 min默认值0即默认从不清理状态语义指定空闲状态未再被更新的状态至少保留的最短时间间隔状态在空闲超过最小时间后、于某个后续时刻被清除需要特别注意的是设置 TTL 可能影响查询结果的正确性——过期的状态被清理后后续到达的延迟数据可能无法匹配到应有的 Join 记录。TTL 的完整说明见 查询配置文档table-exec-state-ttl条目。状态清理会引入额外的簿记开销bookkeeping overhead因此应权衡正确性与资源占用后谨慎设置 TTL。INNER Equi-Join返回由 Join 条件约束的笛卡尔积。目前仅支持等值 Joinequi-join即至少包含一个等值谓词作为合取条件之一的 Join任意的 cross join 或 theta join非等值比较 Join不受支持。SELECT * FROM Orders INNER JOIN Product ON Orders.product_id Product.idOUTER Equi-Join返回满足 Join 条件的合格笛卡尔积中的所有行外加外部表中与另一张表任何行都不匹配的每一行的一份副本。Flink 支持 LEFT、RIGHT 和 FULL 三种外连接同样目前仅支持等值 Join。SELECT * FROM Orders LEFT JOIN Product ON Orders.product_id Product.id; SELECT * FROM Orders RIGHT JOIN Product ON Orders.product_id Product.id; SELECT * FROM Orders FULL OUTER JOIN Product ON Orders.product_id Product.id;Interval Join带时间约束的等值 JoinInterval Join 返回由 Join 条件和一个时间约束共同限制的笛卡尔积。它要求至少一个等值 Join 谓词以及一个在两侧约束时间的 Join 条件。合法的时间范围谓词包括、、、的组合、BETWEEN谓词或者一个直接比较两侧同类型时间属性的等值谓词处理时间或事件时间均可。例如下面的查询将订单与其对应的物流记录关联起来前提是订单在收到 4 小时之后才发货SELECT * FROM Orders o, Shipments s WHERE o.id s.order_id AND o.order_time BETWEEN s.ship_time - INTERVAL 4 HOUR AND s.ship_time以下是合法的 Interval Join 条件示例ltime rtimeltime rtime AND ltime rtime INTERVAL 10 MINUTEltime BETWEEN rtime - INTERVAL 10 SECOND AND rtime INTERVAL 5 SECOND对流式查询而言与 Regular Join 相比Interval Join仅支持带时间属性的 append-only 表。由于时间属性是准单调递增的Flink 可以在不影响结果正确性的前提下从状态中移除旧值因此其状态是有界的。时间属性的定义方式可参考 时间属性概念文档其中涵盖了事件时间rowtime需声明 WATERMARK与处理时间proctime两种属性的声明语法与语义。Temporal Join与时序演进表关联时态表Temporal Table是一张随时间演进的表在 Flink 中也称为动态表。时态表中的每一行都关联一个或多个时间周期事实上所有 Flink 表都是时态动态的。时态表包含一个或多个带版本的快照它可以是变更历史表跟踪每一次变更例如数据库 changelog包含所有快照变更维度表只物化最新变更例如数据库表包含最新快照。Temporal Join 接受任意表左侧输入 / probe 侧并将每一行与版本表右侧输入 / build 侧中对应行的相关版本关联。Flink 采用 SQL:2011 标准的FOR SYSTEM_TIME AS OF语法来实现这一操作。语法如下SELECT [column_list] FROM table1 [AS alias1] [LEFT] JOIN table2 FOR SYSTEM_TIME AS OF table1.{ proctime | rowtime } [AS alias2] ON table1.column-name1 table2.column-name1从执行层面看事件时间与处理时间的时态 Join 分别对应两个运行时算子TemporalRowTimeJoinOperator与TemporalProcessTimeJoinOperator位于 flink-table-runtime 的 join/temporal 包并由执行节点 StreamExecTemporalJoin 统一翻译为TwoInputTransformation生成时态表函数 Join 是时态表 Join 的子集二者复用大部分执行逻辑仅在校验上有差别。事件时间 Temporal Join按历史时点关联版本表事件时间 Temporal Join 允许关联版本表即可以用随时间变化的元数据对表进行补全并取回该键在某一特定时间点的值。版本表会保留自最近一次 watermark 以来、以时间标识的所有版本。典型场景是货币换算假设有一张订单表每笔订单以不同币种计价。要统一折算成美元每笔订单需要关联其下单时点对应的汇率。-- 订单表标准的 append-only 动态表 CREATE TABLE orders ( order_id STRING, price DECIMAL(32,2), currency STRING, order_time TIMESTAMP(3), WATERMARK FOR order_time AS order_time - INTERVAL 15 SECOND ) WITH (/* ... */); -- 版本化的汇率表 -- 可来自 Debezium 之类的 CDC、compact 后的 Kafka topic -- 或任何其他能定义版本表的方式。 CREATE TABLE currency_rates ( currency STRING, conversion_rate DECIMAL(32, 2), update_time TIMESTAMP(3) METADATA FROM values.source.timestamp VIRTUAL, WATERMARK FOR update_time AS update_time - INTERVAL 15 SECOND, PRIMARY KEY(currency) NOT ENFORCED ) WITH ( connector kafka, value.format debezium-json, /* ... */ ); SELECT order_id, price, orders.currency, conversion_rate, order_time FROM orders LEFT JOIN currency_rates FOR SYSTEM_TIME AS OF orders.order_time ON orders.currency currency_rates.currency; order_id price currency conversion_rate order_time o_001 11.11 EUR 1.14 12:00:00 o_002 12.51 EUR 1.10 12:06:00关于事件时间时态 Join有两点必须注意触发机制事件时间时态 Join 由左右两侧的 watermark 触发。INTERVAL时间减法即 WATERMARK 声明中的- INTERVAL 15 SECOND用于等待迟到事件以确保 Join 符合预期。请务必确认两侧都已正确设置 watermark。主键约束事件时间时态 Join 要求主键包含在时态 Join 条件的等价条件中。例如上例中表currency_rates的主键currency_rates.currency必须被约束在条件orders.currency currency_rates.currency中。与 Regular Join 相反即使 build 侧发生变更之前的时态表结果也不会受到影响与 Interval Join 相比时态表 Join 不定义记录关联的时间窗口——probe 侧记录总是与时间属性指定的时刻对应的 build 侧版本关联因此 build 侧的行可以任意古老。随着时间推移不再需要的记录版本针对给定主键会从状态中移除。处理时间 Temporal Join关联最新版本处理时间时态表 Join 使用处理时间属性将行关联到外部版本表中某个键的最新版本。从定义上讲处理时间 Join 总是返回给定键的最新值。你可以把 build 侧想象成一个存储全部记录的简单HashMapK, V。它的价值在于当无法将表物化为 Flink 内部动态表时允许 Flink 直接对接外部系统。下面通过一个例子说明。LatestRates是一张维度表例如 HBase 表物化了最新汇率在10:15、10:30、10:52三个时刻其内容如下10:15 SELECT * FROM LatestRates; currency rate US Dollar 102 Euro 114 Yen 1 10:30 SELECT * FROM LatestRates; currency rate US Dollar 102 Euro 114 Yen 1 10:52 SELECT * FROM LatestRates; currency rate US Dollar 102 Euro 116 changed from 114 to 116 Yen 1LatestRates在10:15与10:30的内容相同而 Euro 汇率在10:52从 114 变为 116。Orders是表示支付金额与币种的 append-only 表SELECT * FROM Orders; amount currency 2 Euro arrived at time 10:15 1 US Dollar arrived at time 10:30 2 Euro arrived at time 10:52我们希望把所有订单折算成统一货币amount currency rate amount*rate 2 Euro 114 228 arrived at time 10:15 1 US Dollar 102 102 arrived at time 10:30 2 Euro 116 232 arrived at time 10:52重要限制目前FOR SYSTEM_TIME AS OF语法用于关联任意表/视图的最新版本尚未支持此时应使用如下的时态表函数语法SELECT o_amount, r_rate FROM Orders, LATERAL TABLE (Rates(o_proctime)) WHERE r_currency o_currency不支持的唯一原因是语义考量左侧流的 Join 处理不会等待时态表的完整快照这在生产环境中可能误导用户。处理时间时态表函数 Join 同样存在该语义问题但由于其已存在很久出于兼容性考虑被保留支持。处理时间的结果是不确定的non-deterministic。处理时间时态 Join 最常用于用外部表即维表对数据流进行补全。与 Regular Join 相反build 侧变更不影响之前的时态表结果与 Interval Join 相比时态表 Join 不定义记录关联的时间窗口即旧行不存入状态。时态表函数 Join与时态表函数Temporal Table Function做 Join 的语法与表函数 Join相同。目前仅支持与时态表的内连接和左外连接。时态表函数的完整定义方法见 时态表函数概念文档。假设Rates是一个时态表函数Join 可写作SELECT o_amount, r_rate FROM Orders, LATERAL TABLE (Rates(o_proctime)) WHERE r_currency o_currency时态表 DDL 与时态表函数的主要区别如下时态表 DDL 可以在 SQL 中直接定义而时态表函数不可以两者都支持与时态版本表进行 Temporal Join但只有时态表函数可以关联任意表/视图的最新版本。Lookup Join外部系统维表补全Lookup Join 通常用于从外部系统查询数据来补全enrich表数据。它要求一侧表具有处理时间属性另一侧表由**查找源连接器lookup source connector**支撑。Lookup Join 使用上述处理时间时态 Join的语法其中右侧表由 lookup source 连接器支撑。以下示例展示了 Lookup Join 的语法-- Customers 由 JDBC 连接器支撑可用于 lookup join CREATE TEMPORARY TABLE Customers ( id INT, name STRING, country STRING, zip STRING ) WITH ( connector jdbc, url jdbc:mysql://mysqlhost:3306/customerdb, table-name customers ); -- 为每个订单补全客户信息 SELECT o.order_id, o.total, c.country, c.zip FROM Orders AS o JOIN Customers FOR SYSTEM_TIME AS OF o.proc_time AS c ON o.customer_id c.id;上述示例中Orders表通过 MySQL 数据库中的Customers表进行补全。FOR SYSTEM_TIME AS OF子句后随处理时间属性确保Orders的每一行在 Join 算子处理该行的时间点与满足 Join 谓词的Customers行关联同时它也防止了将来某个关联的Customer行被更新时 Join 结果被二次更新。Lookup Join 还要求一个强制性的等值 Join 谓词即示例中的o.customer_id c.id。从执行角度看Lookup Join 对应的执行节点为 CommonExecLookupJoin在流式执行时可结合异步 I/O 提升外部查询吞吐见 AsyncLookupJoinITCase 等运行时测试用例。Array Expansion数组展开数组展开为给定数组中的每个元素返回一行新记录。目前尚不支持WITH ORDINALITY带序号的展开。SELECT order_id, tag FROM Orders CROSS JOIN UNNEST(tags) AS t (tag)表函数 Table Function Join将一张表与表函数Table Function的返回结果进行 Join左侧外部表的每一行都会与该行对应的表函数调用所产生的所有行进行 Join。用户定义的表函数UDTF在使用前必须完成注册。INNER JOIN如果某行的表函数调用返回空结果则该行被丢弃。SELECT order_id, res FROM Orders, LATERAL TABLE(table_func(order_id)) t(res)LEFT OUTER JOIN如果表函数调用返回空结果对应的外部行被保留并以 null 填充结果。目前对 lateral 表做左外连接时ON子句中需要一个TRUE字面量。SELECT order_id, res FROM Orders LEFT OUTER JOIN LATERAL TABLE(table_func(order_id)) t(res) ON TRUE如何选择合适的 Join决策要点综合以上各类 Join 的语义与状态特性选型时可参考以下决策路径需要保留双侧全部历史、支持更新型输入→ Regular Join必要时配合table.exec.state.ttl控制状态膨胀注意正确性风险两侧都是带时间属性的 append-only 流且只需关联彼此时间接近的记录→ Interval Join状态有界、结果确定需要按事件发生的历史时点关联版本化数据CDC / 变更历史表→ 事件时间 Temporal JoinFOR SYSTEM_TIME AS OF ... rowtime要求主键进入等值条件需要用外部系统MySQL / HBase 等的最新值实时补全流→ Lookup Join 或处理时间时态 Join注意处理时间语义结果不确定需要对数组字段按元素拆行、或用 UDTF 横向扩展行→ Array Expansion 与 Table Function Join内连接 /ON TRUE左外连接。无论选择哪种 Join都应遵循文档开篇的通用建议合理排布FROM子句中表的顺序低频更新在前、高频更新在后并避免产生不受支持的笛卡尔积才能在保证语义正确的前提下获得可预期的性能表现。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink SQL Join 全面指南Regular / Interval / Temporal / Lookup 与表函数连接的实战与原理Flink SQL Join 全面指南Regular / Interval / Temporal / Lookup 与表函数连接的实战与原理 Flink SQ大数据流处理批处理数据工程Flink DataStream 双流 Join 完全指南Window Join 与 Interval Join 实战与源码解析Flink DataStream 双流 Join 完全指南Window Join 与 Interval Join 实战与源码解析 Flink DataStre大数据流处理批处理数据工程Flink DataStream Joining 完全指南Window Join 与 Interval Join 实战与原理Flink DataStream Joining 完全指南Window Join 与 Interval Join 实战与原理 本指南聚焦 Apache Fli大数据流处理批处理数据工程上一篇3步永久保存你的QQ空间记忆GetQzonehistory实战指南下一篇OS.js与Electron终极对比Web桌面平台VS原生桌面开发的完整选择指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考