Flink从入门到真香(13、时间语义的定义)

在watermark之前先说下时间的概念,在https://blog.51cto.com/mapengfei/2554577 里面有各种时间窗口,实际生产中那是以哪个时间为准产生的窗口呢? 事件发生的时间? 进入flink程序的时间?还是flink开始处理的时间
Flink提供了一套设计解决方案
设置可以在代码中env直接设置

1 val env = StreamExecutionEnvironment.getExecutionEnvironment 2 3// env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime) //以事件时间作为窗口聚合 4//env.setStreamTimeCharacteristic(TimeCharacteristic.IngestionTime) //以数据进入flink的时间作为窗口时间 5// env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime) //以Flink实际处理时间作为窗口时间

时间语义

Flink从入门到真香(13、时间语义的定义)

只能说不同的场景下,每个时间都有使用场景,具体根据实际情况来实施

在代码中设置

我们可以直接在代码中,对执行环境调用setStreamCharacteristic方法,设置流的时间特性
具体的时间,还需要从数据中提取时间戳(timestamp),
如果要用事件时间,还需要设置具体取的哪个字段和格式,否则flink也不知道你用的哪个字段

val env = StreamExecutionEnvironment.getExecutionEnvironment
//从调用时刻开始给env创建的每个stream追加时间特性
env.setStreamTimeCharcteristic(TimeCharacteristic.EventTime)

点赞
收藏

评论区

加载中...

相关推荐

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以前

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

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

Flink的WaterMark,及demo实例

实际生产中,由于各种原因,导致事件创建时间与处理时间不一致,收集的规定对实时推荐有较大的影响。所以一般情况时选取创建时间,然后事先创建flink的时间窗口。但是问题来了,如何保证这个窗口的时间内所有事件都到齐了?这个时候就可以设置水位线(waterMark)。概念:支持基于时间窗口操作,由于事件的时间来源于源头系统,很多时候由于网络延迟、分布式处理,以

Flink从入门到真香(13、时间语义的定义) - HelloWorld