文章目录
-
- 一.简介
- 二.窗口Join
-
- 2.1 翻滚窗口(Tumbling Window Join)
- 2.2 滑动窗口Join(Sliding Window Join)
-
- 2.3 会话窗口Join(Session Window Join)
- 2.4.小结
- 三.间隔Join
- 四.示例
-
- 4.1 间隔Join
- 4.2 窗口Join
一.简介
Flink DataStream API中内置有两个可以根据实际条件对数据流进行Join算子:基于间隔的Join和基于窗口的Join。
语义注意事项
- 创建两个流元素的成对组合的行为类似内连接,如果来自一个流的元素与另一个流没有相对应要连接的元素,则不会发出该元素。
- 结合在一起的那些元素将其时间戳设置为位于各自窗口中最大时间戳。例如:以[5,10]为边界的窗口将产生连接的元素的时间戳为9。
二.窗口Join
2.1 翻滚窗口(Tumbling Window Join)
执行滚动窗口连接(Tumbling Window Join)时,具有公共Key和公共Tumbling Window的所有元素都以成对组合形式进行连接,并传递给JoinFunction或FlatJoinFunction。因为这就像一个内连接,在滚动窗口中没有来自另一个流的元素的流的元素不会被输出。

如图所示,我们定义了一个大小为2毫秒的滚动窗口,其结果为[0,1],[2,3], …。该图像显示了每个窗口中所有元素的成对组合,这些元素将传递给JoinFunction。注意,在翻滚窗口[6,7]中没有发出任何内容,因为在绿色流中没有元素与橙色元素⑥、⑦连接。
1import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows; 2import org.apache.flink.streaming.api.windowing.time.Time; 3... 4val orangeStream: DataStream[Integer] = ... 5val greenStream: DataStream[Integer] = ... 6orangeStream.join(greenStream) 7 .where(elem => /* select key */) 8 .equalTo(elem => /* select key */) 9 .window(TumblingEventTimeWindows.of(Time.milliseconds(2))) 10 .apply { 11 12 (e1, e2) => e1 + "," + e2 }
2.2 滑动窗口Join(Sliding Window Join)
在执行滑动窗口连接(Sliding Window Join)时,具有公共Key和公共滑动窗口(Sliding Window )的所有元素都作为成对组合进行连接,并传递给JoinFunction或FlatJoinFunction。当前滑动窗口中没有来自另一个流的元素的流的元素不会被发出。
注意,有些元素可能会在一个滑动窗口中连接,但不会在另一个窗口中连接!

在本例中,我们使用的滑动窗口大小为2毫秒,滑动1毫秒,滑动窗口结果[1,0],[0,1],[1,2],[2、3],… x轴以下是每个滑动窗口的Join结果将被传递给JoinFunction的元素。在这里你还可以看到橙②与绿色③窗口Join(2、3),但不与任何窗口Join[1,2]。
1import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows; 2import org.apache.flink.streaming.api.windowing.time.Time; 3... 4val orangeStream: DataStream[Integer] = ... 5val greenStream: DataStream[Integer] = ... 6orangeStream.join(greenStream) 7 .where(elem => /* select key */) 8 .equalTo(elem => /* select key */) 9 .window(SlidingEventTimeWindows.of(Time.milliseconds(2) /* size */, Time.milliseconds(1) /* slide */)) 10 .apply { 11 12 (e1, e2) => e1 + "," + e2 } 13
2.3 会话窗口Join(Session Window Join)
在执行会话窗口连接时,具有相同键的所有元素(当“组合”时满足会话条件)都以成对的组合进行连接,并传递给JoinFunction或FlatJoinFunction。再次执行内部连接,因此如果会话窗口只包含来自一个流的元素,则不会发出任何输出。

在这里,定义一个会话窗口连接,其中每个会话被至少1ms的间隔所分割。有三个会话,在前两个会话中,来自两个流的连接元素被传递给JoinFunction。在第三次会话中绿色流没有元素,所以⑧⑨不会Join。
1import org.apache.flink.streaming.api.windowing.assigners.EventTimeSessionWindows; 2import org.apache.flink.streaming.api.windowing.time.Time; 3 4... 5val orangeStream: DataStream[Integer] = ... 6val greenStream: DataStream[Integer] = ... 7orangeStream.join(greenStream) 8 .where(elem => /* select key */) 9 .equalTo(elem => /* select key */) 10 .window(EventTimeSessionWindows.withGap(Time.milliseconds(1))) 11 .apply { 12 13 (e1, e2) => e1 + "," + e2 }
2.4.小结
除了对窗口中两条流进行Join,你还可以对它们进行Cogroup,只需将算子定义开始位置的Join()改为coGroup()即可,Join和Cogroup的总体逻辑相同。
二者区别:Join会为两侧输入中每个事件对调用JoinFunction;而Cogroup中CoGroupFunction会以两个输入的元素遍历器为参数,只在每个窗口中被调用一次。
三.间隔Join
interval join用一个公共Key连接两个流的元素(将它们称为A & B),其中流B的元素的时间戳具有相对于流A中的元素的时间戳。 这也可以更正式地表示为b.timestamp ∈ [a.timestamp + lowerBound; a.timestamp + upperBound] or a.timestamp + lowerBound <= b.timestamp <= a.timestamp + upperBound
其中a和b是A和B中共享一个公钥的元素。下界和上界都可以是负的或正的,只要下界小于或等于上界。interval连接目前只执行内部连接。
当将一对元素传递给ProcessJoinFunction时,它们将给两个元素分配更大的时间戳(可以通过ProcessJoinFunction.Context访问)。
注意:间隔连接目前只支持事件时间。

在上面的示例中,我们将“橙色”和“绿色”两个流连接起来,它们的下界为-2毫秒,上界为+1毫秒。默认情况下,这些是包含边界的,但是可以通过.lowerboundexclusive()和. upperboundexclusive()进行设置。
再用更正式的符号来表示angeElem.ts + lowerBound <= greenElem.ts <= orangeElem.ts + upperBound 如三角形所示。
1import org.apache.flink.streaming.api.functions.co.ProcessJoinFunction; 2import org.apache.flink.streaming.api.windowing.time.Time; 3... 4val orangeStream: DataStream[Integer] = ... 5val greenStream: DataStream[Integer] = ... 6orangeStream 7 .keyBy(elem => /* select key */) 8 .intervalJoin(greenStream.keyBy(elem => /* select key */)) 9 .between(Time.milliseconds(-2), Time.milliseconds(1)) 10 .process(new ProcessJoinFunction[Integer, Integer, String] { 11 12 13 override def processElement(left: Integer, right: Integer, ctx: ProcessJoinFunction[Integer, Integer, String]#Context, out: Collector[String]): Unit = { 14 15 16 out.collect(left + "," + right); 17 } 18 }); 19 });
四.示例
4.1 间隔Join
1package com.lm.flink.datastream.join 2import org.apache.flink.api.common.restartstrategy.RestartStrategies 3import org.apache.flink.streaming.api.CheckpointingMode 4import org.apache.flink.streaming.api.environment.CheckpointConfig 5import org.apache.flink.streaming.api.functions.AssignerWithPeriodicWatermarks 6import org.apache.flink.streaming.api.functions.co.ProcessJoinFunction 7import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment 8import org.apache.flink.streaming.api.watermark.Watermark 9import org.apache.flink.streaming.api.windowing.time.Time 10import org.apache.flink.util.Collector 11/** 12 * @Classname IntervalJoin 13 * @Description TODO 14 * @Date 2020/10/27 20:32 15 * @Created by limeng 16 * 区间关联当前仅支持EventTime 17 * Interval JOIN 相对于UnBounded的双流JOIN来说是Bounded JOIN。就是每条流的每一条数据会与另一条流上的不同时间区域的数据进行JOIN。 18 */ 19object IntervalJoin { 20 21 22 def main(args: Array[String]): Unit = { 23 24 25 //设置至少一次或仅此一次语义 26 val env = StreamExecutionEnvironment.getExecutionEnvironment 27 //设置至少一次或仅此一次语义 28 env.enableCheckpointing(20000,CheckpointingMode.EXACTLY_ONCE) 29 //设置 30 env.getCheckpointConfig 31 .enableExternalizedCheckpoints(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION) 32 //设置重启策略 33 env.setRestartStrategy(RestartStrategies.fixedDelayRestart(5,50000)) 34 env.setParallelism(1) 35 val dataStream1 = env.socketTextStream("localhost",9999) 36 val dataStream2 = env.socketTextStream("localhost",9998) 37 import org.apache.flink.api.scala._ 38 val dataStreamMap1 = dataStream1.map(f=>{ 39 40 41 val tokens = f.split(",") 42 StockTransaction(tokens(0),tokens(1),tokens(2).toDouble) 43 }).assignTimestampsAndWatermarks(new AssignerWithPeriodicWatermarks[StockTransaction]{ 44 45 46 var currentTimestamp = 0L 47 val maxOutOfOrderness = 1000L 48 override def getCurrentWatermark: Watermark = { 49 50 51 val tmpTimestamp = currentTimestamp - maxOutOfOrderness 52 println(s"wall clock is ${System.currentTimeMillis()} new watermark ${tmpTimestamp}") 53 new Watermark(tmpTimestamp) 54 } 55 override def extractTimestamp(element: StockTransaction, previousElementTimestamp: Long): Long = { 56 57 58 val timestamp = element.txTime.toLong 59 currentTimestamp = Math.max(timestamp,currentTimestamp) 60 println(s"get timestamp is $timestamp currentMaxTimestamp $currentTimestamp") 61 currentTimestamp 62 } 63 }) 64 65 val dataStreamMap2 = dataStream2.map(f=>{ 66 67 68 val tokens = f.split(",") 69 StockSnapshot(tokens(0),tokens(1),tokens(2).toDouble) 70 }).assignTimestampsAndWatermarks(new AssignerWithPeriodicWatermarks[StockSnapshot]{ 71 72 73 var currentTimestamp = 0L 74 val maxOutOfOrderness = 1000L 75 override def getCurrentWatermark: Watermark = { 76 77 78 val tmpTimestamp = currentTimestamp - maxOutOfOrderness 79 println(s"wall clock is ${System.currentTimeMillis()} new watermark ${tmpTimestamp}") 80 new Watermark(tmpTimestamp) 81 } 82 override def extractTimestamp(element: StockSnapshot, previousElementTimestamp: Long): Long = { 83 84 85 val timestamp = element.mdTime.toLong 86 currentTimestamp = Math.max(timestamp,currentTimestamp) 87 println(s"get timestamp is $timestamp currentMaxTimestamp $currentTimestamp") 88 currentTimestamp 89 } 90 }) 91 dataStreamMap1.print("dataStreamMap1") 92 dataStreamMap2.print("dataStreamMap2") 93 dataStreamMap1.keyBy(_.txCode) 94 .intervalJoin(dataStreamMap2.keyBy(_.mdCode)) 95 .between(Time.minutes(-10),Time.seconds(0)) 96 .process(new ProcessJoinFunction[StockTransaction,StockSnapshot,String] { 97 98 99 override def processElement(left: StockTransaction, right: StockSnapshot, ctx: ProcessJoinFunction[StockTransaction, StockSnapshot, String]#Context, out: Collector[String]): Unit = { 100 101 102 out.collect(left.toString +" =Interval Join=> "+right.toString) 103 } 104 }).print() 105 106 env.execute("IntervalJoin") 107 } 108 case class StockTransaction(txTime:String,txCode:String,txValue:Double) extends Serializable{ 109 110 111 override def toString: String = txTime +"#"+txCode+"#"+txValue 112 } 113 case class StockSnapshot(mdTime:String,mdCode:String,mdValue:Double) extends Serializable { 114 115 116 override def toString: String = mdTime +"#"+mdCode+"#"+mdValue 117 } 118}
结果
1get timestamp is 1603708942 currentMaxTimestamp 1603708942 2dataStreamMap1> 1603708942#000001#10.4 3get timestamp is 1603708942 currentMaxTimestamp 1603708942 4dataStreamMap2> 1603708942#000001#10.4 51603708942#000001#10.4 =Interval Join=> 1603708942#000001#10.4
4.2 窗口Join
1package com.lm.flink.datastream.join 2import java.lang 3import org.apache.flink.api.common.functions.CoGroupFunction 4import org.apache.flink.streaming.api.TimeCharacteristic 5import org.apache.flink.streaming.api.functions.AssignerWithPeriodicWatermarks 6import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment 7import org.apache.flink.streaming.api.watermark.Watermark 8import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows 9import org.apache.flink.streaming.api.windowing.time.Time 10import org.apache.flink.util.Collector 11/** 12 * @Classname InnerLeftRightJoinTest 13 * @Description TODO 14 * @Date 2020/10/26 17:22 15 * @Created by limeng 16 * window join 17 */ 18object InnerLeftRightJoinTest { 19 20 21 def main(args: Array[String]): Unit = { 22 23 24 val env = StreamExecutionEnvironment.getExecutionEnvironment 25 env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime) 26 //每9秒发出一个watermark 27 env.setParallelism(1) 28 env.getConfig.setAutoWatermarkInterval(9000) 29 30 val dataStream1 = env.socketTextStream("localhost", 9999) 31 val dataStream2 = env.socketTextStream("localhost", 9998) 32 33 /** 34 * operator操作 35 * 数据格式: 36 * tx: 2020/10/26 18:42:22,000002,10.2 37 * md: 2020/10/26 18:42:22,000002,10.2 38 * 39 * 这里由于是测试,固水位线采用升序(即数据的Event Time 本身是升序输入) 40 */ 41 import org.apache.flink.api.scala._ 42 val dataStreamMap1 = dataStream1 43 .map(f => { 44 45 46 val tokens = f.split(",") 47 StockTransaction(tokens(0), tokens(1), tokens(2).toDouble) 48 }).assignTimestampsAndWatermarks(new AssignerWithPeriodicWatermarks[StockTransaction] { 49 50 51 var currentTimestamp = 0L 52 val maxOutOfOrderness = 1000L 53 override def getCurrentWatermark: Watermark = { 54 55 56 val tmpTimestamp = currentTimestamp - maxOutOfOrderness 57 println(s"wall clock is ${System.currentTimeMillis()} new watermark ${tmpTimestamp}") 58 new Watermark(tmpTimestamp) 59 } 60 override def extractTimestamp(element: StockTransaction, previousElementTimestamp: Long): Long = { 61 62 63 val timestamp = element.txTime.toLong 64 currentTimestamp = Math.max(timestamp, currentTimestamp) 65 println(s"get timestamp is $timestamp currentMaxTimestamp $currentTimestamp") 66 currentTimestamp 67 } 68 }) 69 70 val dataStreamMap2 = dataStream2 71 .map(f => { 72 73 74 val tokens = f.split(",") 75 StockSnapshot(tokens(0), tokens(1), tokens(2).toDouble) 76 }).assignTimestampsAndWatermarks(new AssignerWithPeriodicWatermarks[StockSnapshot] { 77 78 79 var currentTimestamp = 0L 80 val maxOutOfOrderness = 1000L 81 override def getCurrentWatermark: Watermark = { 82 83 84 val tmpTimestamp = currentTimestamp - maxOutOfOrderness 85 println(s"wall clock is ${System.currentTimeMillis()} new watermark ${tmpTimestamp}") 86 new Watermark(tmpTimestamp) 87 } 88 override def extractTimestamp(element: StockSnapshot, previousElementTimestamp: Long): Long = { 89 90 91 val timestamp = element.mdTime.toLong 92 currentTimestamp = Math.max(timestamp, currentTimestamp) 93 println(s"get timestamp is $timestamp currentMaxTimestamp $currentTimestamp") 94 currentTimestamp 95 } 96 }) 97 98 dataStreamMap1.print("dataStreamMap1") 99 dataStreamMap2.print("dataStreamMap2") 100 101 /** 102 * Join操作 103 * 限定范围是3秒钟的Event Time窗口 104 */ 105 val joinedStream = dataStreamMap1.coGroup(dataStreamMap2) 106 .where(_.txCode) 107 .equalTo(_.mdCode) 108 .window(TumblingEventTimeWindows.of(Time.seconds(3))) 109 110 val innerJoinedStream = joinedStream.apply(new InnerJoinFunction) 111 val leftJoinedStream = joinedStream.apply(new LeftJoinFunction) 112 val rightJoinedStream = joinedStream.apply(new RightJoinFunction) 113 innerJoinedStream.name("InnerJoinedStream").print() 114 leftJoinedStream.name("LeftJoinedStream").print() 115 rightJoinedStream.name("RightJoinedStream").print() 116 env.execute("InnerLeftRightJoinTest") 117 } 118 119 class InnerJoinFunction extends CoGroupFunction[StockTransaction, StockSnapshot, (String, String, String, Double, Double, String)] { 120 121 122 override def coGroup(first: lang.Iterable[StockTransaction], second: lang.Iterable[StockSnapshot], out: Collector[(String, String, String, Double, Double, String)]): Unit = { 123 124 125 import scala.collection.JavaConverters._ 126 val scalaT1 = first.asScala.toList 127 val scalaT2 = second.asScala.toList 128 129 println(scalaT1.size) 130 println(scalaT2.size) 131 /** 132 * Inner join 要比较的是同一个key下,同一个时间窗口内 133 */ 134 if (scalaT1.nonEmpty && scalaT2.nonEmpty) { 135 136 137 for (transaction <- scalaT1) { 138 139 140 for (snapshot <- scalaT2) { 141 142 143 out.collect(transaction.txCode, transaction.txTime, snapshot.mdTime, transaction.txValue, snapshot.mdValue, "Inner Join Test") 144 } 145 } 146 } 147 } 148 } 149 class LeftJoinFunction extends CoGroupFunction[StockTransaction, StockSnapshot, (String, String, String, Double, Double, String)] { 150 151 152 override def coGroup(T1: java.lang.Iterable[StockTransaction], T2: java.lang.Iterable[StockSnapshot], out: Collector[(String, String, String, Double, Double, String)]): Unit = { 153 154 155 /** 156 * 将Java中的Iterable对象转换为Scala的Iterable 157 * scala的集合操作效率高,简洁 158 */ 159 import scala.collection.JavaConverters._ 160 val scalaT1 = T1.asScala.toList 161 val scalaT2 = T2.asScala.toList 162 /** 163 * Left Join要比较的是同一个key下,同一个时间窗口内的数据 164 */ 165 if (scalaT1.nonEmpty && scalaT2.isEmpty) { 166 167 168 for (transaction <- scalaT1) { 169 170 171 out.collect(transaction.txCode, transaction.txTime, "", transaction.txValue, 0, "Left Join Test") 172 } 173 } 174 } 175 } 176 class RightJoinFunction extends CoGroupFunction[StockTransaction, StockSnapshot, (String, String, String, Double, Double, String)] { 177 178 179 override def coGroup(T1: java.lang.Iterable[StockTransaction], T2: java.lang.Iterable[StockSnapshot], out: Collector[(String, String, String, Double, Double, String)]): Unit = { 180 181 182 /** 183 * 将Java中的Iterable对象转换为Scala的Iterable 184 * scala的集合操作效率高,简洁 185 */ 186 import scala.collection.JavaConverters._ 187 val scalaT1 = T1.asScala.toList 188 val scalaT2 = T2.asScala.toList 189 /** 190 * Right Join要比较的是同一个key下,同一个时间窗口内的数据 191 */ 192 if (scalaT1.isEmpty && scalaT2.nonEmpty) { 193 194 195 for (snapshot <- scalaT2) { 196 197 198 out.collect(snapshot.mdCode, "", snapshot.mdTime, 0, snapshot.mdValue, "Right Join Test") 199 } 200 } 201 } 202 } 203 204 case class StockTransaction(txTime: String, txCode: String, txValue: Double) 205 case class StockSnapshot(mdTime: String, mdCode: String, mdValue: Double) 206}
参考
https://www.jianshu.com/p/ba19e4d1d802
公众号

名称:大数据计算
微信号:bigdata_limeng