栏目分类:
子分类:
返回
名师互学网用户登录
快速导航关闭
当前搜索
当前分类
子分类
实用工具
热门搜索
名师互学网 > IT > 前沿技术 > 云计算 > 云平台

Flink window总结

云平台 更新时间: 发布时间: IT归档 最新发布 模块sitemap 名妆网 法律咨询 聚返吧 英语巴士网 伯小乐 网商动力

Flink window总结

补充一:窗口 1.1 触发器
-- 描述
	用来控制窗口什么时候进行计算,本质上是执行窗口函数,可以计算得到结果并输出过程
1.1.1

定义触发器

stream.keyBy(...)
	  .window(...)
	  .trigger(new MyTrigger()

trigger 为 window的内部属性,每个窗口分配器(WindowAssigner)都会对应一个默认的触发器

flink中所有事件时间的窗口触发器都是EventTimeTrigger,类似的有ProcessTimeTrigger和CountTrigger

#作为window的Trigger是一个抽象类,使用时需要实现以下几种方法

 abstract class Trigger implements 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(Iterable> elements, int size, W window, EvictorContext evictorContext);
//执行窗口之后的移除数据操作
void evictAfter(Iterable> elements, int size, W window, EvictorContext evictorContext);
1.3 允许延迟
-- 描述
	在分布式的情况下,总会存在数据因为某种原因,导致到来的数据延迟。我们根据水位线(WaterMark)表示数据的进展,它并不能保证所有数据会按时到来,也就是说当数据到来之前,水位线超过窗口的关闭时间,窗口关闭,导致此条数据没有参与运算。
	对于迟到的数据通常这几种操作
		1.直接丢弃 (不参与计算,丢失数据)
		2.放到侧输出流中 (不参与计算,丢失数据,但能通过侧输出流看到什么数据迟到了)
		3.设置延迟 (等待数据到来,相应会带来延迟)
		
	为了计算的准确性,我们可以在此处设置一个窗口关闭的延迟时间,也就是开一个小时的窗口,加上10分钟延迟时间,窗口并不会在一个小时后关闭,而是在一小时10分后关闭。
	
-- 使用
	stream.key(...)
		  .window(TumbingEventTimeWindows.of(Time.hours(1)))
		  .allowedLateness(Time.minutes(10))  // 指定延迟时间
		  
-- 注意(1)
	到达窗口关闭时间后,会触发计算一次然后输出到下游,但此时窗口并没有真正关闭,窗口内元素没有清除
	在允许延迟时间内,迟到的数据会在之前数据的基础上做运算
	
	这里就可以看到触发和清除这两个动作是分开的,到达窗口关闭时间,只是触发了计算,并没有把窗口内的元素清除
	
-- 注意(2)
	在延迟时间内到达的元素,来一条就重新触发一次计算,可以通过下述代码验证
        DataStreamSource stream = 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 侧输出流
-- 描述
	使用其他方式,多多少少,总会有几个奇葩的数据不会参与运算,而被程序摒弃。如果不想丢掉任何一个数据,可以试试侧输出流,当数据迟到时,使用侧输出流输出出来
	可以理解为班车开走之后,迟到的人,可以乘上“专门”的班车送出去
	
	专门的班车需要提前定义,然后告诉这里的人,你迟到了可以乘这班车走
OutputTag outputTag = new OutputTag("late") {};  // 定义

stream.keyBy(...)
      .window(TumblingEventTimeWindows.of(Time.hours(1)))
      .sideOutputLateData(outputTag)   // 通知迟到的数据从这儿走
-- 注意(1)
	new OutputTag("late") {}  和 new OutputTag("late") 的区别
	带{}的是子类对象,不带{}的是本类对象
	区别:
	
-- 注意(2)
	侧输出流数据类型应与主流数据类型保持一致
	
-- 注意(3)
	同样可以从侧输出流中获取数据,获取方法如下
	getSideOutput是SingleOutputStreamOperator的方法
DataStream lateStream = steam.getSideOutput(outputTag);
1.5 窗口的生命周期 1.5.1 创建
-- 描述
	窗口的类型和基本信息由窗口分配器(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中,那么我们只需根据流中元素不同的标记,而放到不同的流中
	
-- 使用
	定义:
	OutputTag outputTag = 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);
     	 
-- 注意
	侧输出流可以定义多次,相当于备了好几辆车
	
        DataStreamSource stream = 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
转载请注明:文章转载自 www.mshxw.com
本文地址:https://www.mshxw.com/it/895928.html
我们一直用心在做
关于我们 文章归档 网站地图 联系我们

版权所有 (c)2021-2022 MSHXW.COM

ICP备案号:晋ICP备2021003244-6号