即时看!聊聊Flink的必知必会(三)
时间:2023-06-17 06:17:23来源:博客园

概述

在进行流处理时,很多时候想要对流的有界子集进行聚合分析。例如有如下的需求场景:(1)每分钟的页面浏览(PV)次数。

(2)每用户每周的会话次数。

(3)每分钟每传感器的最高温度。


(资料图片)

(4)当电商发布一个秒杀活动时,想要每隔10min了解流量数据。

对于这些需求的处理,程序需要处理元素组,而不是单个元素,因此,通常使用窗口来限定在数据流上的聚合(如count、sum等)的范围,例如"过去5min内的计数"或"最后100个元素的总和",所以在处理流数据时,通常更有意义的是考虑有限窗口上的聚合,而不是整个流。

在阿里的限流框架Sentinel中,关键的资源数据统计算法也是基于窗口的概念来做的。

窗口(window)是处理无限流的核心,使用窗口计算无界流上的聚合。窗口将流分割为有限大小的组,用户可以对这样的组进行计算。窗口可以是由时间驱动的(例如,每30s),也可以是由数据驱动的(例如,每100个元素)。如下所示

Flink流窗口

通俗点来说,窗口(window)可以将无界流分成有限大小的「桶」,我们基于这个「桶」之上,可以构建各种各样的计算。而无界流的拆分方式可以按时间、或者事件的数量,我们可以根据业务场景来定义窗口的大小。

如何对定义创建流窗口?Flink支持不同类型的窗口,分别介绍如下。

(1)滚动窗口:Tumbling Window,是在流中创建不重叠的相邻窗口。它们是固定长度的窗口,没有重叠。可以根据时间对元素进行分组(例如,从10:00到10:05的所有元素进入一个组),或者根据计数(前50个元素进入一个单独的组)对元素进行分组。例如,可以用它来回答这样的问题:“在不重叠的5min间隔内计算流中元素的数量”。

(2)滑动窗口:Sliding Window,类似于滚动窗口,但是窗口可以重叠。滑动窗口是固定长度的窗口,通过用户给定的窗口滑动参数与前面的窗口重叠。例如,如果需要计算最后5min的指标,但希望每分钟显示一个输出时。

(3)会话窗口:Session Window,当对发生的事件进行分组时,将时间接近的分到一组(一个窗口中)。还可以提供会话间隔的配置参数,该参数指示在关闭会话之前需要等待多长时间。

(4)全局窗口:Global Window,Flink将所有元素放到一个窗口中。通常在这种情况下,每个元素都被分配给一个单一的per-key全局窗口(Global Window)。如果不指定任何触发器,就不会触发任何计算。这只有在定义自定义触发器时才有用,该触发器定义了窗口何时结束。

这几种窗口类型表示,可按如下图表示

窗口分配器

窗口分配器用于定义如何将元素分配给窗口。这是通过在调用window()(针对Keyed Stream)或windowAll()(针对non-keyed stream)时指定所选择的WindowAssigner实现的。WindowAssigner负责将每个传入元素分配给一个或多个窗口。

内置窗口分配器

Flink为最常见的场景(滚动时间窗口、滑动时间窗口、全局窗口和会话窗口)提供了预定义的窗口分配器,它们分别如下。

(1)滚动时间窗口:例如,每分钟PV数据(浏览量),代码如下:

TumblingEventTimeWindows.of(Time.minutes(1))

(2)滑动时间窗口:例如,每10s计算一次每分钟的页面浏览量,代码如下:

SlidingEventTimeWindows.of(Time.minutes(1),Time.seconds(10))

(3)会话窗口:例如,每个会话的PV数据,其中会话定义为会话之间至少30min的间隔,代码如下:

EventTimeSessionWindows.withGap(Time.minutes(30))

所有内置的窗口分配器(全局窗口除外)都根据时间向窗口分配元素。基于时间的窗口分配程序(包括会话窗口)有事件时间和处理时间两种形式。示例如下:

自定义窗口分配器

一个Flink窗口程序的总体结构如下Keyed Stream表示如下,在Keyed Stream的情况下,可以使用传入事件的任何属性作为key。在Keyed Stream的窗口计算由多个任务并行执行,因为每个逻辑Keyed Stream都可以独立于其他流进行处理。所有引用相同key的元素将被发送到相同的并行任务。

// Keyed Windowsstream    .keyBy(...)    .window()    .reduce/aggregate/apply()

non-keyed-stream表示如下,在Keyed Stream的情况下,可以使用传入事件的任何属性作为key。在Keyed Stream的窗口计算由多个任务并行执行,因为每个逻辑Keyed Stream都可以独立于其他流进行处理。所有引用相同key的元素将被发送到相同的并行任务。

// Keyed Windowsstream    .windowAll()    .reduce/aggregate/apply()

参考《Flink原理深入与编程实战》

Flink的Window

标签:

  • 上一篇文章: 菜鸟驿站上线生鲜包裹到站专属提醒
  • 下一篇文章: 最后一页
  • 最新
  • 即时看!聊聊Flink的必知必会(三)

    概述在进行流处理时,很多时候想要对流的有界子集进行聚合分析。例

  • 菜鸟驿站上线生鲜包裹到站专属提醒

    6月16日消息,全国菜鸟驿站近日上线生鲜包裹到站专属提醒功能。当消费

  • 世界焦点!刚刚,巨无霸IPO,通过!

    刚刚,巨无霸IPO,通过!,a股,农化,先正达,ipo,上交所,巨无霸,公司主营业务

  • 当前关注:《少年歇洛克》展开营救行动

    本报讯(记者王洋)记者自北方演艺集团获悉,大型实景沉浸式儿童互动体

  • 上海今日启动“亮剑浦江·消费领域个人信息权益保护专项执法行动”

    6月16日下午,上海市网信办、上海市市场监管局共同启动“亮剑浦江·消

  • 日本自卫队出兵3000支持乌军?日韩高调援乌,已抓住俄罗斯软肋

    日韩支持乌军,等于分散了俄罗斯的注意力,让俄军投入更多兵力物资打乌

  • 群众举报济宁一煤矿瞒报亡人事故,应急局:经查,非安全生产事故|全球热点评

    群众举报济宁一煤矿瞒报亡人事故,应急局:经查,非安全生产事故

  • 2023年自动衡器概念股名单出炉(6月16日)

    2023年自动衡器概念股名单出炉(6月16日),以下是南方财富网为您整理的

  • 雪橇犬有几种(雪橇犬有几种品种)-世界视讯

    雪橇犬有6种,分别是阿拉斯加雪橇犬、萨摩耶德雪橇犬、西伯利亚雪橇犬

  • 焦点短讯!信用卡逾期被起诉通常会怎么判?信用卡被起诉开庭如何应诉

    信用卡逾期被起诉通常会怎么判法院会根据具体情况进行判决。一般来说,

  • 定本·育儿百科 畅销10年纪念版_关于定本·育儿百科 畅销10年纪念版介绍_世界热讯

    1、《定本·育儿百科(畅销10年纪念版)》是2010年5月华夏出版社出版的

  • 滨医附院(第一临床医学院)成功举办黄河三角洲医学教育研讨会

    通讯员张莹莹徐军记者陈甜田近日,滨州医学院附属医院(第一临床医学院

  • 中国星辰 | 巧妙设计“落客方案” 确保41名“乘客”到站有序“下车”

    中国航天“一箭多星”发射再创新纪录。6月15日,在我国太原卫星发射中

  • 世界实时:推动盲盒加强合规治理 市场监管总局出台盲盒经营行为规范指引

    新华社北京6月15日电 题:推动盲盒加强合规治理市场监管总局出台盲盒

  • 怎么找银行协商停息挂账?信用卡欠多少会坐牢?

    怎么找银行协商停息挂账可以和银行协商向银行提交停息挂账的申请,

  • 上行速率和下行速率是什么意思 1000兆宽带上行和下行是多少?

    上行速率和下行速率是什么意思?上行速率通常用于描述向互联网发送数

  • 旅游
    • 火影忍者惠比寿的忍者级别是多少?卡卡西实力如何?

    • 筹码分布存在什么特征?筹码分布日线准还是周线准?

    • 信用卡逾期六个月被起诉了怎么办?信用卡逾期被起诉了还能协商吗?

    • 浙江一临时工棚爆燃致5人死亡是怎么回事?炉膛爆燃怎么处理?