-- 描述 用来控制窗口什么时候进行计算,本质上是执行窗口函数,可以计算得到结果并输出过程1.1.1
定义触发器
stream.keyBy(...) .window(...) .trigger(new MyTrigger()
trigger 为 window的内部属性,每个窗口分配器(WindowAssigner)都会对应一个默认的触发器
flink中所有事件时间的窗口触发器都是EventTimeTrigger,类似的有ProcessTimeTrigger和CountTrigger
#作为window的Trigger是一个抽象类,使用时需要实现以下几种方法
abstract class Triggerimplements Serializable onElement():窗口每到来一个元素,都会调用这个方法 onEventTime():当注册的事件事件定时触发时,将调用这个方法 onProcessingTime():当注册的处理时间定时触发时,将调用这个方法 clear():当窗口销毁时,将调用这个方法
TriggerResult onElement(T element, long timestamp, W window, TriggerContext ctx)
此处的TriggerContext对象可以用来注册回调(callback),也就是注册定时器(Timer)
当时间进展到设定的值后,就会执行定义好的操作
TriggerResult onProcessingTime(long time, W window, TriggerContext ctx)
TriggerResult onEventTime(long time, W window, TriggerContext ctx)
说明:
可以看到上面三个方法的返回值都是TriggerResult
他的结构如下 该类型为枚举类型
CONTINUE(false, false) 继续:不对窗口做任何操作
FIRE_AND_PURGE(true, true) 触发清除:触发计算,并清除窗口中数据
FIRE(true, false) 触发:触发计算,保留所有元素
PURGE(false, true) 清除:清除窗口内元素,丢弃该窗口,不做计算
默认情况下:
我们认为窗口的关闭和计算应该是同时发生的,看到上述代码,两者可以根据返回的枚举值不同,影响窗口的行为
1.2 移除器(Evictor)
-- 描述 用来定义移除某些移除数据的逻辑 -- 使用 stream .key(...) .window(...) .evictor(new MyEvictor) 默认情况下会在窗口之前移除数据
//执行窗口之前的移除数据操作 void evictBefore(Iterable1.3 允许延迟> elements, int size, W window, EvictorContext evictorContext); //执行窗口之后的移除数据操作 void evictAfter(Iterable > elements, int size, W window, EvictorContext evictorContext);
-- 描述 在分布式的情况下,总会存在数据因为某种原因,导致到来的数据延迟。我们根据水位线(WaterMark)表示数据的进展,它并不能保证所有数据会按时到来,也就是说当数据到来之前,水位线超过窗口的关闭时间,窗口关闭,导致此条数据没有参与运算。 对于迟到的数据通常这几种操作 1.直接丢弃 (不参与计算,丢失数据) 2.放到侧输出流中 (不参与计算,丢失数据,但能通过侧输出流看到什么数据迟到了) 3.设置延迟 (等待数据到来,相应会带来延迟) 为了计算的准确性,我们可以在此处设置一个窗口关闭的延迟时间,也就是开一个小时的窗口,加上10分钟延迟时间,窗口并不会在一个小时后关闭,而是在一小时10分后关闭。 -- 使用 stream.key(...) .window(TumbingEventTimeWindows.of(Time.hours(1))) .allowedLateness(Time.minutes(10)) // 指定延迟时间 -- 注意(1) 到达窗口关闭时间后,会触发计算一次然后输出到下游,但此时窗口并没有真正关闭,窗口内元素没有清除 在允许延迟时间内,迟到的数据会在之前数据的基础上做运算 这里就可以看到触发和清除这两个动作是分开的,到达窗口关闭时间,只是触发了计算,并没有把窗口内的元素清除 -- 注意(2) 在延迟时间内到达的元素,来一条就重新触发一次计算,可以通过下述代码验证
DataStreamSourcestream = env.addSource(new SourceFunction () { @Override public void run(SourceContext ctx) throws Exception { ctx.collectWithTimestamp(1, 1000L); Thread.sleep(1000L); ctx.collectWithTimestamp(2, 2000L); Thread.sleep(1000L); ctx.emitWatermark(new Watermark(5000L)); Thread.sleep(1000L); ctx.collectWithTimestamp(3, 3000L); ctx.collectWithTimestamp(4, 3000L); Thread.sleep(1000L); } @Override public void cancel() { } }); stream.keyBy(r -> 1) .window(TumblingEventTimeWindows.of(Time.seconds(5))) .allowedLateness(Time.seconds(2)) .process(new ProcessWindowFunction () { @Override public void process(Integer integer, ProcessWindowFunction .Context context, Iterable elements, Collector out) throws Exception { Iterator iterator = elements.iterator(); while (iterator.hasNext()) { Integer next = iterator.next(); out.collect(next); } } }) .print();
-- 注意 窗口允许延迟时间和水位线延迟时间 窗口允许延迟时间 窗口到达关闭时间之后,触发一次计算,但是不销毁窗口和内部的状态,而是等到窗口关闭时间 + 允许延迟时间之后关闭 水位线延迟时间 预料到会有数据发生延迟,将水位线的生成根据元素到来的时间戳延后一段时间1.4 侧输出流
-- 描述 使用其他方式,多多少少,总会有几个奇葩的数据不会参与运算,而被程序摒弃。如果不想丢掉任何一个数据,可以试试侧输出流,当数据迟到时,使用侧输出流输出出来 可以理解为班车开走之后,迟到的人,可以乘上“专门”的班车送出去 专门的班车需要提前定义,然后告诉这里的人,你迟到了可以乘这班车走
OutputTagoutputTag = new OutputTag ("late") {}; // 定义 stream.keyBy(...) .window(TumblingEventTimeWindows.of(Time.hours(1))) .sideOutputLateData(outputTag) // 通知迟到的数据从这儿走
-- 注意(1) new OutputTag("late") {} 和 new OutputTag ("late") 的区别 带{}的是子类对象,不带{}的是本类对象 区别: -- 注意(2) 侧输出流数据类型应与主流数据类型保持一致 -- 注意(3) 同样可以从侧输出流中获取数据,获取方法如下 getSideOutput是SingleOutputStreamOperator的方法
DataStream1.5 窗口的生命周期 1.5.1 创建lateStream = steam.getSideOutput(outputTag);
-- 描述
窗口的类型和基本信息由窗口分配器(window assigners)指定,但是窗口并不是预先建立好的,而是由数据驱动创建。当第一个属于这个窗口的数据元素到达时,创建对应的窗口
-- 元素属于那个窗口
窗口大小为size,开始时间start,结束时间end,元素所属时间t
所属窗口: ceil(t/size)
start = t - t%size
end = start + size
-- 窗口特性
左开又闭 [0,5) 包括0 不包括5 实际上是 [0,4999]
当5秒的数据来了之后,触发第一个窗口的计算
1.5.2 触发
-- 描述 每个窗口都有自己的窗口函数和触发器(Trigger) 窗口函数有两种,一种是增量聚合函数,一种是全量聚合函数 触发器: 1.注册事件时间或者处理时间 2.根据窗口中元素的多少触发1.5.3 销毁
-- 描述 一般情况下,当时间达到了结束点,就会触发计算输出结果,进而清除状态销毁窗口,这里可以认定窗口的销毁和触发计算是同一时刻,flink只针对时间窗口(TimeWindow)有销毁机制 CountWindow和GlobalWindow则没有这种清除状态的机制。 -- 注意 在设置窗口允许延迟时间之后,窗口销毁的时机应该为窗口关闭时间 + 允许延迟时间1.5.4 窗口API调用
-- 描述 两类窗口 1.key之后的窗口 调用.window() 2.没有key的窗口 调用.windowALL() -- 注意 除了定义窗口和窗口的计算方法,其他都是可选的,一般情况下都不需要实现
[外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传(img-ch8Tuz0h-1652605121753)(img/image-20220515151249103.png)]
补充二:迟到数据处理-- 描述 注意只有在事件时间的情况下,才去讨论迟到的数据2.1 设置水位线的延迟时间
-- 描述
水位线是衡量事件进展而产生的,它是我们整个应用的【全局逻辑时钟】
水位线不只针对于窗口的关闭,同样用于Trigger的执行
-- 逻辑时钟
现实世界时钟:我们能看到墙上钟表的一直在走动,不受其他的影响
逻辑时钟:有数据驱动时间发生变化
来了一条数据 1 逻辑时钟 1
来了一条数据 2 逻辑时钟 2
来了一条数据 3 逻辑时钟 3
来了一条数据 5 逻辑时钟 5
来了一条数据 4 逻辑时钟 5
看数据 1,2,3,5 可以发现逻辑时钟是递增的,实际上程序里的逻辑时钟,是算子自己维护的一个变量,这个变量记录着检测到数据的最大值,实际上就是去所有数据的最大值,来一条数据比较一次,有没有这个值大,没有不替换,如果来的值比变量大,则将自己赋予这个变量。
这里可以看到,逻辑时钟的变化仅受数据的影响,不受时间的影响,因为它压根没跟时间交互,你数据来的慢,逻辑时钟的进展就慢,你数据来的快,我逻辑时钟进展的就快
-- 注意
这个逻辑时钟是个全局变量
水位线也是数据流中的元素,这点可以看水位线源码可以看到
水位线到达算子后,水位线带的时间戳变量,和窗口的关闭时间一比较,大于窗口的关闭时间,那就触发运算
2.2 允许窗口处理迟到元素
-- 描述 由于水位线设置的延迟比较小,无法满足后续还有迟到的数据过来。所以使用窗口延迟时间去等待迟到的数据 大部分情况下,设置水位线迟到时间后,迟到的数据不会太多,这样,在水位线到达,触发窗口的关闭,我们先近似正确的结果。然后保持窗口等待延迟数据,每来一条数据窗口就会再次计算,并将更新后的结果输出,这样慢慢修正,就可以得到准确的值2.3 侧输出流
-- 描述 上面我们讲到迟到的数据可以通过侧输出流输出 不止迟到的数据可以这样做,在我们进行业务处理时,如果想让不同的业务分流,如实时数据放到kafka中,维度数据放到hbase中,那么我们只需根据流中元素不同的标记,而放到不同的流中 -- 使用 定义: OutputTagoutputTag = new OutputTag ("late"){}; 解释: 流中的数据类型,应该和主流中的数据类型保持一致 ("late")侧输出流的id,此id不能为空 {} 使用: (1)将迟到的数据放入侧输出流 stream.keyBy(...) .window(...) .sideOutputLateData(outputTag) (2)数据分流 midStream = stream.keyBy(...) .process(...) 在ProcessFunction中调用ctx.output(outputTag,数据) 将输出传入到侧输出流 获取侧输出流: DataStream sideOutput = midStream.getSideOutput(outputTag); -- 注意 侧输出流可以定义多次,相当于备了好几辆车
DataStreamSourcestream = env.fromElements(1, 2, 3); OutputTag outputTag = new OutputTag ("演示2") { }; OutputTag outputTag1 = new OutputTag ("演示3") { }; SingleOutputStreamOperator process = stream.process(new ProcessFunction () { @Override public void processElement(Integer value, Context ctx, Collector out) throws Exception { if (value == 1) { ctx.output(outputTag, value); } else if (value == 2) { ctx.output(outputTag1, value); } else { out.collect(value); } } }); process.getSideOutput(outputTag).print("这里输出1"); process.getSideOutput(outputTag1).print("这里输出2"); process.print("这里输出3");
输出结果
这里输出1> 1 这里输出2> 2 这里输出3> 3
-- 注意 获取从流要从主流获取,而不是一开始的流获取。
(The Dataflow Model: A Practical Approach to Balancing Correctness, Latency, and Cost in Massive-Scale, Unbounded, Out-of-Order Data Processing)
窗口分配器WindowAssigner
出1");
process.getSideOutput(outputTag1).print(“这里输出2”);
process.print(“这里输出3”);
输出结果
这里输出1> 1
这里输出2> 2
这里输出3> 3
```sql -- 注意 获取从流要从主流获取,而不是一开始的流获取。
(The Dataflow Model: A Practical Approach to Balancing Correctness, Latency, and Cost in Massive-Scale, Unbounded, Out-of-Order Data Processing)
窗口分配器WindowAssigner


