
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载Flink DataStream API 提供了两类用于关联两条数据流的核心算子Window Join窗口连接与Interval Join时间区间连接。本文以 docs/content.zh/docs/dev/datastream/operators/joining.md 为骨架系统讲解这两类 Join 的语义、适用场景、完整 Java/Scala 代码示例并结合当前仓库中JoinedStreams、KeyedStream.IntervalJoined、IntervalJoinOperator等源码实现与集成测试剖析其底层原理与边界行为。读完本文你将能够根据业务场景正确选择 Join 类型写出可直接运行的 DataStream Join 程序并理解其内部实现机制。Window Join按窗口对齐两个流的元素Window Join 的作用对象是两个流中 key 相同且落入相同窗口的元素。这些窗口通过 window assigner窗口分配器定义两个流中的元素都会被用于计算窗口的结果。两个流中的元素组合成对之后会被传递给用户定义的JoinFunction或FlatJoinFunction由用户输出符合 join 要求的结果。典型的 Window Join 代码骨架如下stream.join(otherStream) .where(KeySelector) .equalTo(KeySelector) .window(WindowAssigner) .apply(JoinFunction);语义上有两点值得特别注意从两个流中创建的成对元素与inner join类似即一个流中的元素在与另一个流中对应的元素完成 join 之前不会被输出。完成 join 的元素会将它们的 timestamp 设为对应窗口中允许的最大 timestamp。例如一个边界为[5, 10)的窗口中的元素在 join 之后的 timestamp 为9即10 - 1。从源码结构看这个 DSL 定义在 flink-streaming-java/src/main/java/org/apache/flink/streaming/api/datastream/JoinedStreams.java 中Where类负责通过KeySelector提取第一个流的 keyJoinedStreams.java#L123EqualTo类负责匹配第二个流的 keyJoinedStreams.java#L139、L154随后进入WithWindow阶段除指定WindowAssigner外还支持配置trigger、evictor和allowedLatenessJoinedStreams.java#L288、L310、L332。提示窗口相关的 assigner 类型、时间语义与延迟数据处理方式可进一步参考 Window 算子文档。滚动 Window JoinTumbling Window Join使用滚动 Window Join 时所有 key 相同且共享同一个滚动窗口的元素会被组合成对并传递给JoinFunction或FlatJoinFunction。由于该行为与 inner join 类似一个流中的元素如果没有与另一个流中的元素组合起来就不会被输出。下图演示了大小为 2 毫秒的滚动窗口即形成边界为[0,1], [2,3], ...的窗口实例。每个窗口中的元素会两两配对结果传递给JoinFunction。注意滚动窗口[6,7]将不会产生任何输出因为绿色流中没有数据可以与橙色流的元素 ⑥ 和 ⑦ 配对。Java 示例import org.apache.flink.api.java.functions.KeySelector; import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows; import org.apache.flink.streaming.api.windowing.time.Time; ... DataStreamInteger orangeStream ...; DataStreamInteger greenStream ...; orangeStream.join(greenStream) .where(KeySelector) .equalTo(KeySelector) .window(TumblingEventTimeWindows.of(Time.milliseconds(2))) .apply (new JoinFunctionInteger, Integer, String (){ Override public String join(Integer first, Integer second) { return first , second; } });Scala 示例import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows import org.apache.flink.streaming.api.windowing.time.Time ... val orangeStream: DataStream[Integer] ... val greenStream: DataStream[Integer] ... orangeStream.join(greenStream) .where(elem /* select key */) .equalTo(elem /* select key */) .window(TumblingEventTimeWindows.of(Time.milliseconds(2))) .apply { (e1, e2) e1 , e2 }滑动 Window JoinSliding Window Join当使用滑动 Window Join 时所有 key 相同且处于同一个滑动窗口的元素将被组合成对并传递给JoinFunction或FlatJoinFunction。与前面相同这是 inner join 语义当前滑动窗口内如果一个流中的元素没有与另一个流中的元素组合起来它就不会被输出。注意在某个滑动窗口中被 join 的元素不一定会出现在其他滑动窗口的 join 结果中——因为滑动窗口会按滑动步长产生多个窗口实例元素在不同实例中的配对对象可能不同。下图定义了长度为 2 毫秒、滑动距离为 1 毫秒的滑动窗口生成的窗口实例区间为[-1, 0], [0, 1], [1, 2], [2, 3], ...。X 轴下方是每个滑动窗口中被 join 后传递给JoinFunction的元素。可以看到橙色元素 ② 与绿色元素 ③ 在窗口[2,3]中完成 join但没有与窗口[1,2]中的任何元素 join。Java 示例import org.apache.flink.api.java.functions.KeySelector; import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows; import org.apache.flink.streaming.api.windowing.time.Time; ... DataStreamInteger orangeStream ...; DataStreamInteger greenStream ...; orangeStream.join(greenStream) .where(KeySelector) .equalTo(KeySelector) .window(SlidingEventTimeWindows.of(Time.milliseconds(2) /* size */, Time.milliseconds(1) /* slide */)) .apply (new JoinFunctionInteger, Integer, String (){ Override public String join(Integer first, Integer second) { return first , second; } });Scala 示例import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows import org.apache.flink.streaming.api.windowing.time.Time ... val orangeStream: DataStream[Integer] ... val greenStream: DataStream[Integer] ... orangeStream.join(greenStream) .where(elem /* select key */) .equalTo(elem /* select key */) .window(SlidingEventTimeWindows.of(Time.milliseconds(2) /* size */, Time.milliseconds(1) /* slide */)) .apply { (e1, e2) e1 , e2 }会话 Window JoinSession Window Join使用会话 Window Join 时所有 key 相同且组合后符合会话要求的元素将被组合成对并传递给JoinFunction或FlatJoinFunction。该操作同样是 inner join如果一个会话窗口中只含有某一个流的元素这个窗口将不会产生输出。下图定义了一个间隔gap为至少 1 毫秒的会话窗口。图中总共有三个会话前两个会话中两个流都有元素它们被 join 并传递给JoinFunction而第三个会话中绿流没有任何元素所以元素 ⑧ 和 ⑨ 没有被 join。Java 示例import org.apache.flink.api.java.functions.KeySelector; import org.apache.flink.streaming.api.windowing.assigners.EventTimeSessionWindows; import org.apache.flink.streaming.api.windowing.time.Time; ... DataStreamInteger orangeStream ...; DataStreamInteger greenStream ...; orangeStream.join(greenStream) .where(KeySelector) .equalTo(KeySelector) .window(EventTimeSessionWindows.withGap(Time.milliseconds(1))) .apply (new JoinFunctionInteger, Integer, String (){ Override public String join(Integer first, Integer second) { return first , second; } });Scala 示例import org.apache.flink.streaming.api.windowing.assigners.EventTimeSessionWindows import org.apache.flink.streaming.api.windowing.time.Time ... val orangeStream: DataStream[Integer] ... val greenStream: DataStream[Integer] ... orangeStream.join(greenStream) .where(elem /* select key */) .equalTo(elem /* select key */) .window(EventTimeSessionWindows.withGap(Time.milliseconds(1))) .apply { (e1, e2) e1 , e2 }源码视角Window Join 如何落到 CoGroup 上从实现上看Window Join 的apply并不是独立实现一套窗口逻辑而是委托给 CoGroupedWindowedStream在 JoinedStreams.java 中apply(FlatJoinFunction, ...)通过coGroupedWindowedStream.apply(new JoinCoGroupFunction(function), resultType)完成即用一个JoinCoGroupFunction适配器把 coGroup共组的输出转换成 join 结果。这也是为什么 Window Join 天然继承 coGroup 的窗口机制——同一个窗口内先按 key 分组再对两组元素做笛卡尔积式配对。这一点在单元测试 JoinedStreamsTest.java 中也有体现测试testDelegateToCoGrouped构造了.join(...).where(...).equalTo(...).window(...).allowedLateness(...)链式调用后直接断言内部getCoGroupedWindowedStream()的allowedLateness配置被正确透传印证了 Join 与 CoGroup 的委托关系。Interval Join按时间区间对齐两个流的元素Interval Join 组合元素的条件为两个流这里称为 A 和 B中 key 相同且B 中元素的 timestamp 处于 A 中元素 timestamp 的一定范围内。该条件可以形式化表示为b.timestamp ∈ [a.timestamp lowerBound; a.timestamp upperBound]即a.timestamp lowerBound b.timestamp a.timestamp upperBound其中a和b分别为 A 和 B 中共享相同 key 的元素。上界和下界可正可负只要下界永远小于等于上界即可。目前 Interval Join 仅执行 inner join。当一对元素被传递给ProcessJoinFunction时它们的 timestamp 取两个元素 timestamp 中的最大值该 timestamp 可以通过ProcessJoinFunction.Context访问。重要限制Interval Join 目前仅支持 event time。若在 processing time 下调用between(...)会抛出UnsupportedTimeCharacteristicException见下文源码分析。下图演示了对橙色与绿色两个流执行 Interval Joinjoin 条件为下界 -2 毫秒、上界 1 毫秒。默认情况下上下界都被包含在区间内但可以通过.lowerBoundExclusive()和.upperBoundExclusive()将它们排除在外。图中三角形所表示的条件可以写成orangeElem.ts lowerBound greenElem.ts orangeElem.ts upperBoundJava 示例import org.apache.flink.api.java.functions.KeySelector; import org.apache.flink.streaming.api.functions.co.ProcessJoinFunction; import org.apache.flink.streaming.api.windowing.time.Time; ... DataStreamInteger orangeStream ...; DataStreamInteger greenStream ...; orangeStream .keyBy(KeySelector) .intervalJoin(greenStream.keyBy(KeySelector)) .between(Time.milliseconds(-2), Time.milliseconds(1)) .process (new ProcessJoinFunctionInteger, Integer, String(){ Override public void processElement(Integer left, Integer right, Context ctx, CollectorString out) { out.collect(left , right); } });Scala 示例import org.apache.flink.streaming.api.functions.co.ProcessJoinFunction import org.apache.flink.streaming.api.windowing.time.Time ... val orangeStream: DataStream[Integer] ... val greenStream: DataStream[Integer] ... orangeStream .keyBy(elem /* select key */) .intervalJoin(greenStream.keyBy(elem /* select key */)) .between(Time.milliseconds(-2), Time.milliseconds(1)) .process(new ProcessJoinFunction[Integer, Integer, String] { override def processElement(left: Integer, right: Integer, ctx: ProcessJoinFunction[Integer, Integer, String]#Context, out: Collector[String]): Unit { out.collect(left , right) } })Interval Join 的 API 细节与源码佐证Interval Join 的 DSL 定义在 flink-streaming-java/src/main/java/org/apache/flink/streaming/api/datastream/KeyedStream.java 中KeyedStream.intervalJoin(KeyedStream otherStream)KeyedStream.java#L440要求左右两侧都已经是KeyedStream即先各自keyBy这是它与 Window Join 在调用链上最明显的差异。IntervalJoined.between(Duration lowerBound, Duration upperBound)KeyedStream.java#L524-L535是设置时间边界的方法边界值会被转换成毫秒数默认上下界均包含inclusive。从源码可以看到该方法内部会做时间特性校验若timeBehaviour ! TimeBehaviour.EventTime立即抛出UnsupportedTimeCharacteristicException(Time-bounded stream joins are only supported in event time)。同时上下界参数允许为null表示单侧无限源码中仅做了非空校验后直接使用。边界开闭控制lowerBoundExclusive()KeyedStream.java#L594与upperBoundExclusive()KeyedStream.java#L587分别将下界、上界置为排他exclusive。迟到数据处理sideOutputLeftLateData(OutputTag)与sideOutputRightLateData(OutputTag)KeyedStream.java#L604、L615可将 watermark 之后到达的左右两侧迟到数据分别路由到指定的 side output而不是直接丢弃。结果输出process(ProcessJoinFunction)KeyedStream.java#L630完成整个 join 操作每个配对元素都会调用用户函数且输出类型通过TypeExtractor自动推导也可显式传入TypeInformation。上述能力均有集成测试覆盖见 flink-tests/src/test/java/org/apache/flink/test/streaming/runtime/IntervalJoinITCase.java测试testBoundsAreInclusiveByDefaultIntervalJoinITCase.java#L357-L386验证了默认上下界均包含的行为between(0, 2)时timestamp 为 0 的左元素可与右流中 timestamp 为 0、1、2 的元素全部配对。测试testLowerBoundExclusiveAndUpperBoundExclusiveIntervalJoinITCase.java#L314-L317验证了排他边界between(0, 2)配合.upperBoundExclusive().lowerBoundExclusive()后边界点不再参与配对。测试testExecutionFailsInProcessingTimeIntervalJoinITCase.java#L388-L389验证了 processing time 下调用between会抛出UnsupportedTimeCharacteristicException与源码中的时间特性校验一一对应。单侧边界下界为null或上界为null的使用方式也在测试中有所体现IntervalJoinITCase.java#L281-L296。从算子层面看Interval Join 的底层运行时由 flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/co/IntervalJoinOperator.java 承担。它基于 event time 的 watermark 推进为每个 key 维护左右两侧元素的时间有序缓存当 watermark 越过某个元素的“过期”边界时该元素即被清理同时只有两侧元素的时间差落入[lowerBound, upperBound]按开闭配置调整区间内才会触发ProcessJoinFunction。这一设计保证了区间匹配只与事件时间相关与数据的到达顺序解耦。总结与选型建议对比维度Window JoinInterval Join配对条件两侧 key 相同且落入同一窗口两侧 key 相同且时间差在指定区间内窗口类型Tumbling / Sliding / Session由 WindowAssigner 决定无显式窗口由between(lower, upper)定义时间范围调用链join→where→equalTo→window→applykeyBy→intervalJoin→between→process用户函数JoinFunction/FlatJoinFunctionProcessJoinFunction可访问ContextJoin 语义inner join未配对元素不输出inner join仅 inner join时间特性取决于 assigner 与执行环境时间特性仅支持 event time边界控制窗口边界由 assigner 定义上下界默认包含可配置排他上下界可正可负单侧可为 null选型建议业务上需要按**固定时间段滚动/滑动或活动间隙会话**对齐两个流时选择 Window Join它语义直观、与窗口算子生态trigger、evictor、allowedLateness完全打通。业务上需要刻画“一个事件发生后前后一段时间内的关联事件”如点击后 5 秒内的曝光、下单前后 2 分钟内的浏览这类相对时间关联时选择 Interval Join它的区间表达更贴近业务直觉且天然与 event time watermark 的迟到处理机制配合可将迟到数据定向输出到 side output 做补偿处理。两种 Join 在数据流拓扑上都需要按 key 对两侧数据进行分区对齐Window Join 通过窗口机制Interval Join 通过两侧keyBy因此合理设计 key 的选择是保证 Join 正确性与负载均衡的关键。更多窗口机制与时间语义细节可继续阅读 Window 算子文档 与 Process Function 文档。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink SQL 窗口关联Window Join完整指南语法、语义与源码实现解析Flink SQL 窗口关联Window Join完整指南语法、语义与源码实现解析 窗口关联Window Join是 Apache Flink 在时间大数据流处理批处理数据工程Flink SQL Join 全面指南Regular / Interval / Temporal / Lookup 与表函数连接的实战与原理Flink SQL Join 全面指南Regular / Interval / Temporal / Lookup 与表函数连接的实战与原理 Flink SQ大数据流处理批处理数据工程突破Flink SQL性能瓶颈Broadcast Join与Shuffle Join优化指南突破Flink SQL性能瓶颈Broadcast Join与Shuffle Join优化指南 你是否还在为Flink SQL关联查询的性能问题头疼当数据流不大数据流处理批处理数据工程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考