【Flink】Flink 中的时间和窗口之窗口(Window)

1. 窗口的概念

Flink是一种流式计算引擎,主要是来处理无界数据流,数据流的数据是一直都有的,等待流结束输入数据获取所有的流数据在做聚合计算是不可能的。为了更方便高效的处理无界流,一种方式就是把无限的流数据切割成有限的数据块进行处理,这就是Flink中提到的窗口(Windows)

在Flink中,窗口就是用来处理无界流的核心。我们很容易把窗口想象成一个固定位置的,数据源源不断的流过来,到某个时间点窗口该关闭了,就停止收集数据,触发计算并输出结果。

例如,我们定义了一个时间窗口,每10秒统计一次数据,呢么就相当于把窗口放在那里,从0秒开始收集数据,到10秒时,处理当前窗口内所有的数据,输出一个结果,然后清空窗口继续收集数据;到20秒时,再对窗口内所有数据进行计算处理,输出结果;以此类推:
在这里插入图片描述
这里使用的窗口[0,10)窗口是左闭右开区间,即包含起始时间点,但不包括结束时间点。对于处理实时数据的窗口来说,这种方式存在一定问题。因为基于系统时间进行窗口关闭操作,在某些情况下可能会出现处理结果不准确或丢失部分数据的情况。例如,在一个 0-10 秒的窗口关闭后,如果还有一条时间戳为 9 秒的数据到达,则该数据将无法被正确地处理,并只能进入下一个 10-20 秒的窗口中。

然而如果我们采用事件时间语义,就会有一些费解了。由于乱序数据,我们需要设置一个延迟时间来等所有数据到齐。比如上面的例子,我们可以设置延迟时间为2秒,如下图,这样0-10秒的窗口会在时间戳为12秒的数据到来之后,才真正关闭计算输出结果,这样就可以正常包含迟到的9秒数据了。
在这里插入图片描述
但是这样一来,0-10秒的窗口不光包含了迟到的9秒数据,连11秒和12秒的数据也包含进去了。我们为了正确处理迟到数据,结果把早到的数据划分到了错误的窗口----最终结果也是错的

所以为了解决这个问题,窗口其实并不是一个框,流进来的数据被框住只能进这一个窗口。窗口而是一个桶。在Flink中,窗口可以把流切割成有限大小的多个存储桶;每个数据都会分发到对应的桶中,当到达窗口结束时间时,就对每个桶中收集的数据进行计算处理
在这里插入图片描述
在事件时间语义下,窗口的处理过程:

1. 第一个数据时间戳为2,判断之后创建第一个窗口[0,10),并将2秒数据保存进去;
2. 后续数据依次到来,时间戳均在[0,10)范围内,所以全部保存进第一个窗口
3. 11秒数据到来,判断不属于[0,10)窗口,所以创建第二个窗口[10,20),并将11秒的数据保存进去。由于水位线设置延迟时间为2秒,所以现在的时钟是9秒,第一个窗口也没有到关闭时间;
4. 之后又有9秒数据到来,同样进入[0,10)窗口中;
5. 12秒数据到来,判断属于[10,20)窗口,保存进去。这时产生的水位线推进到了10秒,所以[0,10)窗口应该关闭了。第一个 窗口收集到了所有的7个数据,进行处理计算后输出结果,并将窗口关闭销毁;
6. 同样的,之后的数据依次进入第二个窗口,遇到20秒的数据时会创建第三个窗口[20,30)并将数据保存进去;遇到22秒数据时,水位线到了20秒,第二个窗口触发计算,输出结果并关闭

注意!!! Flink 中窗口并不是静态准备好的,而是动态创建的——当有落在这个窗口区间范围的数据到达时,才创建对应的窗口。另外,这里我们认为到达窗口结束时间时,窗口就触发计算并关闭,事实上触发计算和窗口关闭两个行为也可以分开。

2. 窗口的分类

Flink中有很多种类的窗口,上面说的就是最简单的一种时间窗口

2.1 按照驱动类型分类

窗口本身是截取有界数据的一种方式,所以窗口最重要的信息就是怎样截取数据,以什么标准来开始和结束数据的截取,叫做窗口的驱动类型
在这里插入图片描述

2.1.1 时间窗口(Time Window)

时间窗口(Time Window)就是按照时间段去截取数据,这也是最常见的窗口。时间窗口以时间点来定义窗口的开始(start)和结束(end),所以截取出的就是某一时间段的数据。到达结束时间时,窗口不再收集数据,触发计算输出结果,并将窗口关闭销毁,也可以说基本思路就是定点发车

用结束时间减去开始时间,得到这段时间的长度,就是窗口的大小(windows size)。这里的时间可以是不同的语义,所以我们可以定义处理时间窗口和事件时间窗口。

Flink中有一个专门的类来表示时间窗口,名称叫做TimeWindow。这个类只有两个私有属性startend,这表示窗口的开始和结束的时间戳,单位为毫秒。可以通过公有的get方法调用。另外TImeWindow还提供了一个maxTimestamp()方法,用来获取窗口中能够包含数据的最大时间戳。通过代码可以看出最大时间戳就是end-1,这也代表了时间窗口的时间范围都是左闭右开的区间[start,end)

@PublicEvolving
public class TimeWindow extends Window {
    private final long start;
    private final long end;
    public TimeWindow(long start, long end) {
        this.start = start;
        this.end = end;
    }
    public long getStart() {
        return start;
    }
    public long getEnd() {
        return end;
    }
    @Override
    public long maxTimestamp() {
        return end - 1;
    }
    ....
}

2.1.2 计数窗口(CountWindow)

计数窗口是基于元素个数来截取数据,到达固定的个数时就触发计算并关闭窗口。类似于座位有限,坐满就发车,至于是否发车和时间没有任何关系。每个窗口的截取数据的个数,就是窗口的大小。

计数窗口相比时间窗口就更加简单,我们只需要指定窗口大小,就可以把数据分配到对应的窗口中,在Flink中没有相对应的类表示计数窗口,底层通过全局窗口(Global Window)来实现的。maxTimestamp返回的Long.MAX_VALUE

@PublicEvolving
public class GlobalWindow extends Window {

    private static final GlobalWindow INSTANCE = new GlobalWindow();
    private GlobalWindow() {}

    public static GlobalWindow get() {
        return INSTANCE;
    }
    @Override
    public long maxTimestamp() {
        return Long.MAX_VALUE;
    }
    @Override
    public boolean equals(Object o) {
        return this == o || !(o == null || getClass() != o.getClass());
    }
    @Override
    public int hashCode() {
        return 0;
    }
    @Override
    public String toString() {
        return "GlobalWindow";
    }
    ....
}

2.2 按照窗口分配数据的规则分类

时间窗口和计数窗口只是对窗口的一个大致划分,再具体应用时,还需要定义更加精细的规则,来控制数据应该划分到哪个窗口。不同的分配数据的方式,就可以有不同的功能应用。

根据分配数据的规则,窗口的具体实现划分为4类:滚动窗口(Tumbling Window)、滑动窗口(Sliding Window)、会话窗口(Session Window)、以及全局窗口(Global Window)。

2.2.1 滚动窗口(Tumbling Windows)

滚动窗口有固定的大小,是一种对数据进行的均匀切片的划分方式。窗口之间没有重叠,也不会有间隔,是"均匀切片"的划分你方式。窗口之间没有重叠,也不会相隔,是首尾相接的状态。如果我们把多个窗口的创建,看作一个窗口的运动,就类似于在不停的向前翻滚一样。这是最简单的窗口形式。也因为滚动窗口是无缝衔接,所以每个数据都会被分配到一个窗口上,而且也只属于一个窗口。

滚动窗口可以基于时间定义,也可以基于数据个数定义;需要的参数只有一个:窗口的大小(Windows Size)。窗口的大小可以使一个小时一次,也可以是长度为10的数据个数。
在这里插入图片描述
如上图所示,圆点表示数据流的数据,对数据按照userID做了分区。当固定了窗口大小之后,所有的分区的窗口划分都是一致的;窗口没有重叠,每个数据只属于一个窗口。
滚动窗口应用非常广泛,它可以对每个时间段做聚合统计,很多BI分析指标都可以用它来实现。

2.2.2 滑动窗口(Sliding Windows)

滑动窗口和滚动窗口类似,滑动窗口的大小也是很固定的。区别在于窗口之间并不是首尾相连的,而是错开一定的位置。如果看作一个窗口的运动,呢么就像是向前小步滑动一样,所以滑动窗口的参数就有两个,一个是窗口大小(Windows Size),一个是滑动的步长(Windows slide),它其实就代表了窗口计算的频率。滑动的距离代表了下个窗口开始的时间间隔,而窗口大小是固定的,所以也就是两个窗口结束时间的间隔;窗口在结束时间触发计算输出结果,呢么滑动步长就代表了计算频率。例如:我们定义一个长度为1小时,滑动步长为5分钟的滑动窗口,呢么就会统计1小时内的数据,每5分钟统计一次。同样,滑动窗口也可以基于时间定义,也可以基于数据个数定义。
在这里插入图片描述
当滑动步长小于窗口大小时,滑动窗口就会出现重叠,这时候的部分数据也可能被同时分配到多个窗口中去。而具体的个数,就由窗口大小和滑动步长的比值(size/slide)来决定。如图6-18所示,滑动步长刚好是窗口大小的一半,呢么在windows1和windows2的中间部分,每个数据都会被分配到这2个窗口里。。比如窗口长度定义1个小时,滑动步长为30分钟,呢么对于8.55的数据就分别属于[8,9)和[8.30,9.30]这两个窗口;

所以,滑动窗口是固定大小窗口的更广义的一种形式;换句话说,滚动窗口也是一种特殊的滑动窗口——窗口大小等于滑动步长(size==slide)

2.2.3 会话窗口(Session Windows)

会话窗口是基于会话(session)来对数据进行分组的。这里的会话类似Web的会话session概念,不过并不代表两端的通讯过程,而是借用会话超时失效的机制来描述窗口。简单来说就是当有数据来了就开启一个窗口,如果还有数据到来就一直保持开启状态,如果在等待一段时间后没有收到数据,就认为会话失效窗口自动关闭。

与滑动窗口和滚动窗口不同,会话窗口只能基于时间来定义,而没有"会话计数窗口"的概念。类似于"会话"终止的标志就是"隔一段时间没有数据来",如果不依赖时间而改成个数,就成了"隔几个数据没有来",这是自相矛盾的说法。

会话窗口有两个重要概念,一个是这段时间的长度——Size,它表示会话的超时时间,也就是两个会话窗口之间的最小距离。还有一个是两个数据到来的时间间隔——Gap,如果新的数据到来时间小于指定的大小size,那说明还在保持会话,就属于同一个窗口;但如果gap大于size,呢么新来的数据就应该属于新的会话窗口,前一个窗口就需要关闭了。具体实现上还可以设置静态固定大小Size,也可以通过一个自定义提取器(Gap Extractor)动态提取最小间隔Gap的值

考虑到事件时间语义下的乱序流,这里又会有一些麻烦。相邻两个数据的时间间隔 gap
大于指定的 size,我们认为它们属于两个会话窗口,前一个窗口就关闭;可在数据乱序的情况
下,可能会有迟到数据,它的时间戳刚好是在之前的两个数据之间的。这样一来,之前我们判
断的间隔中就不是“一直没有数据”,而缩小后的间隔有可能会比 size 还要小——这代表三个
数据本来应该属于同一个会话窗口。所以在 Flink 底层,对会话窗口的处理会比较特殊:每来一个新的数据,都会创建一个新的会话窗口;然后判断已有窗口之间的距离,如果小于给定的 size,就对它们进行合并(merge)操作。在 Window 算子中,对会话窗口会有单独的处理逻辑。
在这里插入图片描述
会话窗口和之前两种窗口不同,没有固定长度,起始和结束时间也不确定,各个分区之间窗口也是没有联系的。如图 6-19 所示,会话窗口之间一定是不会重叠的,而且会留有至少为 size 的间隔(session gap)。

2.2.4 全局窗口(Global Windows)

还有一类比较通用的窗口,就是“全局窗口”。这种窗口全局有效,会把相同 key 的所有数据都分配到同一个窗口中;说直白一点,就跟没分窗口一样。无界流的数据永无止尽,所以这种窗口也没有结束的时候,默认是不会做触发计算的。如果希望它能对数据进行计算处理,还需要自定义“触发器”(Trigger)。
在这里插入图片描述
如图 6-20 所示,可以看到,全局窗口没有结束的时间点,所以一般在希望做更加灵活的窗口处理时自定义使用。Flink 中的计数窗口(Count Window),底层就是用全局窗口实现的。

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:/a/411709.html

如若内容造成侵权/违法违规/事实不符,请联系我们进行投诉反馈qq邮箱809451989@qq.com,一经查实,立即删除!

相关文章

【析】装卸一体化车辆路径问题的自适应并行遗传算法

0 引言 国内外有关 VRPSPD的文献较多,求解目标多以最小化车辆行驶距离为主,但现实中可能存在由租赁费用产生的单次派出成本,需要综合考 虑单次派车成本和配送路径成本。…

消息中间件之RocketMQ源码分析(十八)

Broker CommitLog索引机制中的构建过程 1.创建ConsumeQueue和IndexFile。 ConsumeQueue和IndexFile两个索引都是由ReputMessageService类创建的 RequestMessageService类图 ReputMessageService服务启动后的执行过程。 doReput()方法用于创建索引的入口,通常通过…

Redis 管道详解

Redis 管道 关键词:Pipeline Pipeline 简介 Redis 是一种基于 C/S 模型以及请求/响应协议的 TCP 服务。通常情况下,一个 Redis 命令的请求、响应遵循以下步骤: 客户端向服务端发送一个查询请求,并监听 Socket 返回&#xff08…

3 easy 26. 删除有序数组中的重复项

双指针: //给你一个 非严格递增排列 的数组 nums ,请你 原地 删除重复出现的元素,使每个元素 只出现一次 ,返回删除后数组的新长度。元素的 相对顺序 应该保持 //一致 。然后返回 nums 中唯一元素的个数。 // // 考虑 nums 的唯…

PreMaint CMS系统:数字化驱动起重机远程管理的智能未来

在港口起重机领域,数字化解决方案正成为提高效率、降低停机时间的关键。PreMaint CMS作为起重机核心可视化系统,融合了操作人员、维护团队和管理者的需求,通过智能化的方式,将复杂的数据转化为简洁的信息,实现对起重机…

Node.js中的错误处理和日志记录

在Node.js应用程序中,错误处理和日志记录是非常重要的方面。正确的错误处理和日志记录可以帮助我们更好地跟踪问题、排查 bug,并提供更好的用户体验。在本篇博客中,我将为大家介绍在Node.js中如何进行错误处理和日志记录,以及一些…

BIO实战、NIO编程与直接内存、零拷贝深入辨析

BIO实战、NIO编程与直接内存、零拷贝深入辨析 长连接、短连接 长连接 socket连接后不管是否使用都会保持连接状态多用于操作频繁,点对点的通讯,避免频繁socket创建造成资源浪费,比如TCP 短连接 socket连接后发送完数据后就断开早期的http服…

基于JSP的毕业设计选题系统的设计与实现

基于JSP的毕业设计选题系统的设计与实现 (源代码论文) A. 项目简介 毕业设计选题系统就是能够使学生通过互联网完成毕业设计课题的选定,它采用Web方式,同时适用于局域网和Internet,它要实现审核,权限管理,邮件通知…

Spring6学习技术|IoC|基于注解管理bean

学习材料 尚硅谷Spring零基础入门到进阶,一套搞定spring6全套视频教程(源码级讲解) IoC注解 首先这是最常用的方法。(在学Java基础的时候明明是非常不起眼的知识点啊!!!) 从 Java…

二进制部署k8s之网络部分

1 CNI 网络组件 1.1 K8S的三种接口 CRI 容器运行时接口 docker containerd podman cri-o CNI 容器网络接口 flannel calico cilium CSI 容器存储接口 nfs ceph gfs oss s3 minio 1.2 K8S的三种网络 节点网络 nodeIP 物理网卡的IP实现节点间的通信 Pod网络 podIP Pod与Po…

RabbitMQ服务启动失败

报错信息: 在服务中启动RabbitMQ服务显示: RabbitMQ 服务正在启动 . RabbitMQ 服务无法启动。 系统出错。 发生系统错误 1067。 进程意外终止 报错原因: 1.Erlang与RabbitMQ是否匹配 2.Erlang与RabbitMQ安装路径是否存在中文或空格 3.电…

解决谷歌浏览器,每次重启都重置所有设置的问题

一、问题: 修改谷歌浏览器的设置 关闭浏览器再打开设置页面后,会显示(部分设置已重置Chrome检测到您的部分设置被其他程序篡改了,因此已将这些设置重置为原始默认设置。了解详情) 之前的设置被重置了 二、解决 创建—…

VS2022调试技巧(一)

什么是bug? 在1945年,美国科学家Grace Hopper在进行计算机编程时,发现一只小虫子钻进了一个真空管,导致计算机无法正常工作。她取出虫子后,计算机恢复了正常,由此,她首次将“Bug”这个词用来描…

跳房子 Ⅰ(C语言)

题目来自于博主算法大师的专栏:最新华为OD机试C卷AB卷OJ(CJavaJSPy) https://blog.csdn.net/banxia_frontend/category_12225173.html 题目描述 跳房子,也叫跳飞机,是一种世界性的儿童游戏。 游戏参与者需要分多个回…

Go语言必知必会100问题-06 生产者端接口

生产者端接口 Go语言必知必会100问题-05 接口污染中介绍了程序中使用接口是有价值的。在编码的时候,接口应该放在哪里呢?这是Go开发人员经常有误解的一个问题,本文将深入分析该问题。 在深入探讨问题之前,先对提及的术语做一个定…

2024年老薛主机开工大吉活动:云服务器5折起,续费同价!

2024年老薛主机开工大吉活动开始了,香港/美国云服务器季付7折,半年付6折,年付5折,续费同价! 活动地址: 点此直达老薛主机官网 活动详情: 老薛主机2024开年促销活动,香港/美国云服…

全网最全AI绘画工具汇总(二)

一.AI绘画 图像 创造人工智能艺术的方式共有多种方法,包括使用数字模式的程序“基于规则”的图像生成、模拟笔触和其他绘画效果的算法,以及人工智能或深度学习算法等。 最早的重要人工智能艺术系统之一是AARON,由哈罗德科恩于1960年代末开…

Pytorch添加自定义算子之(1)-安装配置Eigen库

一、安装对应的ubuntu环境 推荐使用Docker FROM nvcr.io/nvidia/pytorch:23.01-py3 RUN pip install tensorboardX RUN pip install pyyaml RUN pip install yacs RUN pip install termcolor RUN pip install opencv-python RUN pip install timm0.6.12 WORKDIR /app COPY . …

【LeetCode:2476. 二叉搜索树最近节点查询 + 中序遍历 + 有序表】

🚀 算法题 🚀 🌲 算法刷题专栏 | 面试必备算法 | 面试高频算法 🍀 🌲 越难的东西,越要努力坚持,因为它具有很高的价值,算法就是这样✨ 🌲 作者简介:硕风和炜,…

vDPA测试环境搭建

要求: 运行 Linux 发行版的计算机。本指南使用 CentOS 9-stream,但对于其他 Linux 发行版,特别是对于 Red Hat Enterprise Linux 7,命令不应有重大变化。 具有 sudo 权限的用户 ~ 主目录中有 25 GB 的可用空间 至少 8GB 内存 …