详解Flink的窗口操作

我们经常需要在一个时间窗口维度上对数据进行聚合,窗口是流处理应用中经常需要解决的问题。Flink的窗口算子为我们提供了方便易用的API,我们可以将数据流切分成一个个窗口,对窗口内的数据进行处理,下面为大家详细讲解一下Flink的窗口操作。

一、窗口(window)的类型

对于窗口的操作主要分为两种,分别对于Keyedstream和Datastream。他们的主要区别也仅仅在于建立窗口的时候一个为.window(…),一个为.windowAll(…)。对于Keyedstream的窗口来说,他可以使得多任务并行计算,每一个logical key stream将会被独立的进行处理。

stream
      .keyBy(...)               "assigner"
     [.trigger(...)]            "trigger" (else default trigger)
     [.evictor(...)]            "evictor" (else no evictor)
     [.allowedLateness(...)]    "lateness" (else zero)
     [.sideOutputLateData(...)] "output tag" (else no side output for late data)
      .reduce/aggregate/fold/apply()      "function"
     [.getSideOutput(...)]      "output tag"

按照窗口的Assigner来分,窗口可以分为

Tumbling window, sliding window,session window,global window,custom window

每种窗口又可分别基于processing time和event time,这样的话,窗口的类型严格来说就有很多。

还有一种window叫做count window,依据元素到达的数量进行分配,之后也会提到。

窗口的生命周期开始在第一个属于这个窗口的元素到达的时候,结束于第一个不属于这个窗口的元素到达的时候。

二、窗口的操作

2.1 Tumbling window

固定相同间隔分配窗口,每个窗口之间没有重叠看图一眼明白。

下面的例子定义了每隔3毫秒一个窗口的流:

WindowedStream
  
    Rates = rates    .keyBy(MovieRate::getUserId)    .window(TumblingEventTimeWindows.of(Time.milliseconds(3))); 
  

2.2 Sliding Windows

跟上面一样,固定相同间隔分配窗口,只不过每个窗口之间有重叠。窗口重叠的部分如果比窗口小,窗口将会有多个重叠,即一个元素可能被分配到多个窗口里去。

下面的例子给出窗口大小为10毫秒,重叠为5毫秒的流:

WindowedStream
  
    Rates = rates                .keyBy(MovieRate::getUserId)                .window(SlidingEventTimeWindows.of(Time.milliseconds(10), Time.milliseconds(5))); 
  

2.3 Session window

这种窗口主要是根据活动的事件进行窗口化,他们通常不重叠,也没有一个固定的开始和结束时间。一个session window关闭通常是由于一段时间没有收到元素。在这种用户交互事件流中,我们首先想到的是将事件聚合到会话窗口中(一段用户持续活跃的周期),由非活跃的间隙分隔开。

// 静态间隔时间
WindowedStream
  
    Rates = rates                .keyBy(MovieRate::getUserId)                .window(EventTimeSessionWindows.withGap(Time.milliseconds(10))); // 动态时间 WindowedStream
   
     Rates = rates                .keyBy(MovieRate::getUserId)                .window(EventTimeSessionWindows.withDynamicGap(())); 
   
  

2.4 Global window

将所有相同keyed的元素分配到一个窗口里。好吧,就这样:

WindowedStream
  
    Rates = rates    .keyBy(MovieRate::getUserId)    .window(GlobalWindows.create()); 
  

三、窗口函数

窗口函数就是这四个:ReduceFunction,AggregateFunction,FoldFunction,ProcessWindowFunction。前两个执行得更有效,因为Flink可以增量地聚合每个到达窗口的元素。

Flink必须在调用函数之前在内部缓冲窗口中的所有元素,所以使用ProcessWindowFunction进行操作效率不高。不过ProcessWindowFunction可以跟其他的窗口函数结合使用,其他函数接受增量信息,ProcessWindowFunction接受窗口的元数据。

举一个AggregateFunction的例子吧,下面代码为MovieRate按user分组,且分配5毫秒的Tumbling窗口,返回每个user在窗口内评分的所有分数的平均值。

DataStream
  
   > Rates = rates                .keyBy(MovieRate::getUserId)                .window(TumblingEventTimeWindows.of(Time.milliseconds(5)))                .aggregate(new AggregateFunction
   
    >() {                    @Override                    public AverageAccumulator 
    createAccumulator() {                        
    return new AverageAccumulator();                    }                    @Override                    public AverageAccumulator add(MovieRate movieRate, AverageAccumulator acc) {                        acc.userId = movieRate.userId;                        acc.sum += movieRate.rate;                        acc.count++;                        
    return acc;                    }                    @Override                    public Tuple2
    
      getResult(AverageAccumulator acc) {                        
     return  Tuple2.of(acc.userId, acc.sum/(double)acc.count);                    }                    @Override                    public AverageAccumulator merge(AverageAccumulator acc0, AverageAccumulator acc1) {                        acc0.count += acc1.count;                        acc0.sum += acc1.sum;                        
     return acc0;                    }                }); public static class AverageAccumulator{        int userId;        int count;        double sum;    } 
    
   
  

以下是部分输出:

...
1> (44,3.0)
4> (96,0.5)
2> (51,0.5)
3> (90,2.75)
...

看上面的代码,会发现add()函数特别生硬,因为我们想返回Tuple2 类型,即Integer为key,但AggregateFunction似乎没有提供这个机制可以让AverageAccumulator的构造函数提供参数。所以,这里引入ProcessWindowFunction与AggregateFunction的结合版,AggregateFunction进行增量叠加,当窗口关闭时,ProcessWindowFunction将会被提供AggregateFunction返回的结果,进行Tuple封装:

DataStream
  
   > Rates = rates    .keyBy(MovieRate::getUserId)    .window(TumblingEventTimeWindows.of(Time.milliseconds(5)))    .aggregate(new MyAggregateFunction(), new MyProcessWindowFunction()); public static class MyAggregateFunction implements AggregateFunction
   
     {    @Override    public AverageAccumulator 
    createAccumulator() {        
    return new AverageAccumulator();    }    @Override    public AverageAccumulator add(MovieRate movieRate, AverageAccumulator acc) {        acc.sum += movieRate.rate;        acc.count++;        
    return acc;    }    @Override    public Double getResult(AverageAccumulator acc) {        
    return  acc.sum/(double)acc.count;    }    @Override    public AverageAccumulator merge(AverageAccumulator acc0, AverageAccumulator acc1) {        acc0.count += acc1.count;        acc0.sum += acc1.sum;        
    return acc0;    } } public static class MyProcessWindowFunction extends    ProcessWindowFunction
    
     , Integer, TimeWindow> {    @Override    public void process(Integer key,                        Context context,                        Iterable
     
       results,                        Collector
      
       > out) throws Exception {        Double result = results.iterator().next();        out.collect(new Tuple2(key, result));    } } public static class AverageAccumulator{    int count;    double sum; } 
      
     
    
   
  

可以得到,结果与上面一样,但代码好看了很多。

四、其他操作

4.1 Triggers(触发器)

触发器定义了窗口何时准备好被窗口处理。每个窗口分配器默认都有一个触发器,如果默认的触发器不符合你的要求,就可以使用trigger(…)自定义触发器。

通常来说,默认的触发器适用于多种场景。例如,多有的event-time窗口分配器都有一个EventTimeTrigger作为默认触发器。该触发器在watermark通过窗口末尾时出发。

PS:GlobalWindow默认的触发器时NeverTrigger,该触发器从不出发,所以在使用GlobalWindow时必须自定义触发器。

4.2 Evictors(驱逐器)

Evictors可以在触发器触发之后以及窗口函数被应用之前和/或之后可选择的移除元素。使用Evictor可以防止预聚合,因为窗口的所有元素都必须在应用计算逻辑之前先传给Evictor进行处理

4.3 Allowed Lateness

当使用event-time窗口时,元素可能会晚到,例如Flink用于跟踪event-time进度的watermark已经超过了窗口的结束时间戳。

默认来说,当watermark超过窗口的末尾时,晚到的元素会被丢弃。但是flink也允许为窗口operator指定最大的allowed lateness,以至于可以容忍在彻底删除元素之前依然接收晚到的元素,其默认值是0。

为了支持该功能,Flink会保持窗口的状态,知道allowed lateness到期。一旦到期,flink会删除窗口并删除其状态。

把晚到的元素当作side output。

SingleOutputStreamOperator
  
    result = input    .keyBy(
   
    )    .window(
    
     )    .allowedLateness(
     )    .sideOutputLateData(lateOutputTag)    .
      
       (
       
        function>); 
       
      
    
   
  

文章来源网络,作者:管理,如若转载,请注明出处:https://shuyeidc.com/wp/220481.html<

(0)
管理的头像管理
上一篇2025-04-14 14:26
下一篇 2025-04-14 14:28

相关推荐

  • 服务器域名解析失败怎么排查,是什么原因造成的?

    服务器域名解析失败的根本原因在于DNS系统无法将域名正确转换为IP地址,排查应遵循从客户端到服务端的顺序:先检查本地网络和DNS缓存,再验证域名解析记录和权威服务器状态,第一步:检查本地网络与DNS设置测试网络连通性先确认你的设备是否正常联网,打开命令提示符或终端,输入ping 8.8.8.8,如果返回回复数据……

    2026-07-27
    0
  • 站群服务器怎么设置不同环境配置,有哪些注意事项?

    站群服务器设置不同的环境配置,核心在于通过虚拟化或容器化技术实现站点隔离,再结合Web服务器配置为每个站点分配独立的PHP版本、数据库及运行参数,从而满足多样化需求,为什么站群服务器需要环境隔离?不同CMS依赖的PHP版本差异明显,例如WordPress推荐PHP 7.4以上,而Drupal 7仍基于PHP 5……

    2026-07-27
    0
  • 站群服务器如何批量管理更高效,有哪些管理技巧?

    站群服务器批量管理想提效,自动化是唯一出路,通过统一配置管理工具与面板系统,结合服务商提供的底层基础设施支持,能将运维效率提升数倍,批量管理的核心痛点与解决思路多台站群服务器分散管理,最常见的问题就是重复劳动,每次软件更新、配置修改、安全加固,都需要逐台登录操作,不仅耗时,还容易漏掉某台机器,更头疼的是,一旦某……

    2026-07-27
    0
  • 服务器磁盘IO过高如何优化?,磁盘IO过高的原因有哪些?

    服务器磁盘IO过高,核心优化路径是“先定位、再分流、后升级”,你需要通过系统工具精确判断究竟是应用程序、日志策略还是硬件瓶颈导致,然后针对性地从代码、缓存、存储架构和硬件选型四个层面下手,其中选择持有持牌自营机房和增值电信业务经营许可证的服务商,能从根本上保障底层IO稳定性,定位IO瓶颈:动手优化的第一步盲目优……

    2026-07-27
    0
  • 跨境网站访问延迟高怎么解决,网站访问慢的原因是什么?

    跨境网站访问延迟高的核心解决思路在于多维度优化网络路径,包括使用全球CDN加速、选择靠近目标区域的优质IDC机房、调整传输协议以及精简应用层资源,其中服务商的基础设施质量直接决定优化上限,为什么跨境访问延迟高?三大核心因素物理距离与光速限制数据包在海底光缆中的传输速度受限于介质,从中国到美国西海岸的物理往返时间……

    2026-07-27
    0

发表回复

您的邮箱地址不会被公开。必填项已用 * 标注