Flink的WaterMark,及demo实例

实际生产中,由于各种原因,导致事件创建时间与处理时间不一致,收集的规定对实时推荐有较大的影响。所以一般情况时选取创建时间,然后事先创建flink的时间窗口。但是问题来了,如何保证这个窗口的时间内所有事件都到齐了?这个时候就可以设置水位线(waterMark)。

概念:支持基于时间窗口操作,由于事件的时间来源于源头系统,很多时候由于网络延迟、分布式处理,以及源头系统等各种原因导致源头数据的事件时间可能乱序。这时可以设定一个时间阈值,或者说水位线(waterMark),其作用定义一个最大乱序时间,比如某条日志时间为2019-01-01 08:00:10,如果乱序最大允许时间为10s,那么就认为2019-01-01 08:00:00之前产生的所有事件都到齐了,可以进行计算。

时间窗口:指定一个固定时间间隔的窗口

一、滑动窗口

1、SlidingEventTimeWindows.of(Time.second(4), Time.seconds(3)):表示滑动窗口大小为4秒,滑动步长是3 秒,同时,每3秒才滑动一次;

2、每条数据存活的时间为滑动窗口的大小;

3、如果滑动窗口超过之前的窗口,那么后面来的属于前面窗口的数据会丢失;

4、来了一条数据,边移动边计算滑动窗口的数据(一个窗口停留,计算一次,不移动,不计算 ),直至窗口到达指定位置。

计算某位置时间的公式:

1//n:时间戳;size窗口大小;slide:滑动长度 2//根据等差公式推导 3an = a1 + (x-1)*s 4a1 = size - slide -1 5x = [n - (size-slide)]/slide //除数后再乘以slide 6s = slide 7 8//当来了一条时间戳为n的事件,就认为指定位置时间之前的所有事件都到齐了 9指定位置 = (size-slide-1) + [(n-waterMark) - (size-slide)]/slide * slide

二、翻滚窗口

基于时间窗口,对连续数据进行迭代计算时,不会重叠。翻滚窗口是一个特殊的滑动窗口,当窗口的长度等于滑动的长度时,滑动窗口就是翻滚窗口。

计算某位置时间的公式:

指定位置 = -1 + (n-waterMark)/size * size     //除数后再乘以size,size为窗口大小,n为时间戳

三、会话窗口

时间间隔达到一定时间长度时才进行统计计算。

测试代码(需要集群telnet一个producer):

1package com.cjs 2 3import org.apache.flink.streaming.api.TimeCharacteristic 4import org.apache.flink.streaming.api.functions.timestamps.BoundedOutOfOrdernessTimestampExtractor 5import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment 6import org.apache.flink.streaming.api.windowing.time.Time 7import org.apache.flink.api.scala._ 8import org.apache.flink.streaming.api.windowing.assigners.{SlidingEventTimeWindows, TumblingEventTimeWindows} 9 10object WaterMarkTest { 11 12 /** 13 *想使用WaterMark,需要3个步骤: 14 * 1、对数据进行timestamp提取,即调用assignTimestampsAndWatermarks函数, 15 * 实例化BoundedOutOfOrdernessTimestampExtractor,重写extractTimestamp方法 16 * 2、设置使用事件时间,因为WaterMark是基于事件时间 17 * 3、定义时间窗口:翻滚窗口(TumblingEventWindows)、滑动窗口(timeWindow) 18 * 任意一个没有实现,都会报异常:Record has Long.MIN_VALUE timestamp (= no timestamp marker). Is the time characteristic set to 'ProcessingTime', or did you forget to call 'DataStream.assignTimestampsAndWatermarks(...)'? 19 */ 20 def main(args: Array[String]): Unit = { 21 val senv = StreamExecutionEnvironment.getExecutionEnvironment 22 23 val streamAdd = senv.socketTextStream("192.168.112.10",9999) 24 val stream = streamAdd.assignTimestampsAndWatermarks( 25 new BoundedOutOfOrdernessTimestampExtractor[String](Time.seconds(0)) { //WaterMark设置 26 //对数据流进行处理,获取timestamp,对数据流就够不影响 27 override def extractTimestamp(element: String): Long ={ 28 //定义timestamp怎么从数据中抽取出来 29 val eventTime = element.split(" ")(0).toLong 30 print(s"$eventTime \n") 31 eventTime 32 } 33 }) //提取时间戳之后,该数据流是带有时间的,用于事件窗口 34 .map(x=>(x.split(" ")(1),1L)).keyBy(0) 35 36 //设置使用事件时间,因为WaterMark是基于事件时间 37 senv.setStreamTimeCharacteristic(TimeCharacteristic.EventTime) 38 //定义翻滚窗口 39// stream.window(TumblingEventTimeWindows.of(Time.seconds(3))).sum(1).print() 40// stream.sum(1).print() //直接输出,没有用到事件时间窗口,flink默认是累计统计,来一个,统计一个 41 //定义滑动窗口 42 stream.window(SlidingEventTimeWindows.of(Time.seconds(4),Time.seconds(2))).sum(1).print() 43 senv.execute("watermark") 44 } 45 46}
点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

Java日期时间API系列31

  时间戳是指格林威治时间1970年01月01日00时00分00秒起至现在的总毫秒数,是所有时间的基础,其他时间可以通过时间戳转换得到。Java中本来已经有相关获取时间戳的方法,Java8后增加新的类Instant等专用于处理时间戳问题。 1获取时间戳的方法和性能对比1.1获取时间戳方法Java8以前

双十一预售活动分析

2022年双十一促销活动已经开始,大家应该都提前开始关注今年双十一活动的时间表了吧?2022年10月24日晚8:00天猫双11预售时间,第一波销售时间10月31日晚8:0,第二波销售时间11月10日晚8:00;天猫双11的优惠力度是跨店每满30050

​一篇文章总结一下Python库中关于时间的常见操作

前言本次来总结一下关于Python时间的相关操作,有一个有趣的问题。如果你的业务用不到时间相关的操作,你的业务基本上会一直用不到。但是如果你的业务一旦用到了时间操作,你就会发现,淦,到处都是时间操作。。。所以思来想去,还是总结一下吧,本次会采用类型注解方式。time包importtime时间戳从1970年1月1日00:00:00标准时区诞生到现在