44、Flink 的默认窗口剔除器 evictor 代码示例

2024-06-14 10:44

本文主要是介绍44、Flink 的默认窗口剔除器 evictor 代码示例,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

1、CountEvictor

仅记录用户指定数量的元素,一旦窗口中的元素超过这个数量,多余的元素会从窗口缓存的开头移除。

2、DeltaEvictor

接收 DeltaFunction 和 threshold 参数,计算最后一个元素与窗口缓存中所有元素的差值,并移除差值大于或等于 threshold 的元素。

3、TimeEvictor

接收 interval 参数,以毫秒表示,它会找到窗口中元素的最大 timestamp max_ts,并移除比 max_ts - interval 小的所有元素。

注意

Flink 不对窗口中元素的顺序做任何保证,即使 evictor 从窗口缓存的开头移除一个元素,这个元素也不一定是最先或者最后到达窗口的;

默认情况下,所有内置的 evictor 逻辑都在调用窗口函数前执行;

4、代码示例

package com.xu.flink.datastream.day09;import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.windowing.WindowFunction;
import org.apache.flink.streaming.api.functions.windowing.delta.DeltaFunction;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.evictors.CountEvictor;
import org.apache.flink.streaming.api.windowing.evictors.DeltaEvictor;
import org.apache.flink.streaming.api.windowing.evictors.TimeEvictor;
import org.apache.flink.streaming.api.windowing.triggers.EventTimeTrigger;
import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
import org.apache.flink.util.Collector;import java.time.Duration;/*** Evictor 可以在 trigger 触发后、调用窗口函数之前或之后从窗口中删除元素* evictBefore() 包含在调用窗口函数前的逻辑,在调用窗口函数之前被移除的元素不会被窗口函数计算* evictAfter() 包含在窗口函数调用之后的逻辑* <p>* -默认情况下,所有内置的 evictor 逻辑都在调用窗口函数前执行-* CountEvictor: 仅记录用户指定数量的元素,一旦窗口中的元素超过这个数量,多余的元素会从窗口缓存的开头移除。* DeltaEvictor: 接收 DeltaFunction 和 threshold 参数,计算最后一个元素与窗口缓存中所有元素的差值,并移除差值大于或等于 threshold 的元素。* TimeEvictor: 接收 interval 参数,以毫秒表示,它会找到窗口中元素的最大 timestamp `max_ts`,并移除比 `max_ts - interval` 小的所有元素。* <p>* 注意:Flink 不对窗口中元素的顺序做任何保证,即使 evictor 从窗口缓存的开头移除一个元素,这个元素也不一定是最先或者最后到达窗口的。*/
public class _12_WindowDefaultEvictors {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();DataStreamSource<String> input = env.socketTextStream("localhost", 8888);// 测试时限制了分区数,生产中需要设置空闲数据源env.setParallelism(2);// 事件时间需要设置水位线策略和时间戳SingleOutputStreamOperator<Tuple2<String, Long>> map = input.map(new MapFunction<String, Tuple2<String, Long>>() {@Overridepublic Tuple2<String, Long> map(String input) throws Exception {String[] fields = input.split(",");return new Tuple2<>(fields[0], Long.parseLong(fields[1]));}});SingleOutputStreamOperator<Tuple2<String, Long>> watermarks = map.assignTimestampsAndWatermarks(WatermarkStrategy.<Tuple2<String, Long>>forBoundedOutOfOrderness(Duration.ofSeconds(0)).withTimestampAssigner(new SerializableTimestampAssigner<Tuple2<String, Long>>() {@Overridepublic long extractTimestamp(Tuple2<String, Long> input, long l) {return input.f1;}}));//1、CountEvictor: 仅记录用户指定数量的元素,一旦窗口中的元素超过这个数量,多余的元素会从窗口缓存的开头移除(默认)。//1.1 doEvictAfter = false//a,1718157600000//b,1718157600000//c,1718157600000////a,1718157602000//b,1718157602000//c,1718157602000////a,1718157604000//b,1718157604000//c,1718157604000////a,1718157605001//b,1718157605001////Window 的开始和结束时间=>1718157600000-1718157605000//2> a//Window 的开始和结束时间=>1718157600000-1718157605000//1> b//Window 的开始和结束时间=>1718157600000-1718157605000//1> c////c,1718157605001////1.2 doEvictAfter = true//a,1718157600000//b,1718157600000//c,1718157600000////a,1718157602000//b,1718157602000//c,1718157602000////a,1718157604000//b,1718157604000//c,1718157604000////a,1718157605001//b,1718157605001////Window 的开始和结束时间=>1718157600000-1718157605000//2> a//2> a//2> a//Window 的开始和结束时间=>1718157600000-1718157605000//1> b//1> b//1> b//Window 的开始和结束时间=>1718157600000-1718157605000//1> c//1> c//1> c////c,1718157605001////2、DeltaEvictor:接收 DeltaFunction 和 threshold 参数,计算最后一个元素与窗口缓存中所有元素的差值,并移除差值大于或等于 threshold 的元素。//a,1718157600000//b,1718157600000//c,1718157600000////a,1718157602000//b,1718157602000//c,1718157602000////a,1718157604000//b,1718157604000//c,1718157604000////a,1718157605001-最后的元素//b,1718157605001-最后的元素////first=>(a,1718157600000),last=>(a,1718157604000)//first=>(a,1718157602000),last=>(a,1718157604000)//first=>(a,1718157604000),last=>(a,1718157604000)//first=>(b,1718157600000),last=>(b,1718157604000)//first=>(b,1718157602000),last=>(b,1718157604000)//first=>(b,1718157604000),last=>(b,1718157604000)//first=>(c,1718157600000),last=>(c,1718157604000)//first=>(c,1718157602000),last=>(c,1718157604000)//first=>(c,1718157604000),last=>(c,1718157604000)////Window 的开始和结束时间=>1718157600000-1718157605000//Window 的开始和结束时间=>1718157600000-1718157605000//Window 的开始和结束时间=>1718157600000-1718157605000//2> a//2> a//2> a//c,1718157605001-最后的元素////3、TimeEvictor:接收 interval 参数,以毫秒表示,它会找到窗口中元素的最大 timestamp `max_ts`,并移除比 `max_ts - interval` 小的所有元素//a,1718157600000//b,1718157600000//c,1718157600000////a,1718157602000//b,1718157602000//c,1718157602000////a,1718157604000//b,1718157604000//c,1718157604000////a,1718157605001//b,1718157605001////Window 的开始和结束时间=>1718157600000-1718157605000//2> a//Window 的开始和结束时间=>1718157600000-1718157605000//1> b//Window 的开始和结束时间=>1718157600000-1718157605000//1> c////c,1718157605001watermarks.keyBy(e -> e.f0).window(TumblingEventTimeWindows.of(Duration.ofSeconds(5))).trigger(EventTimeTrigger.create()).evictor(TimeEvictor.of(Duration.ofSeconds(1), false))
//                .evictor(CountEvictor.of(1, true))
//                .evictor(DeltaEvictor.of(1.0, new DeltaFunction<Tuple2<String, Long>>() {
//                    @Override
//                    public double getDelta(Tuple2<String, Long> first, Tuple2<String, Long> last) {
//                        System.out.println("first=>"+first+",last=>"+last);
//                        if(first.f0.startsWith("a") && last.f0.startsWith("a")){
//                            return 0.0;
//                        }
//                        return 2.0;
//                    }
//                }, false)).apply(new WindowFunction<Tuple2<String, Long>, String, String, TimeWindow>() {@Overridepublic void apply(String s, TimeWindow timeWindow, Iterable<Tuple2<String, Long>> iterable, Collector<String> collector) throws Exception {System.out.println("Window 的开始和结束时间=>" + timeWindow.getStart() + "-" + timeWindow.getEnd());for (Tuple2<String, Long> tuple2 : iterable) {collector.collect(tuple2.f0);}}}).print();env.execute();}
}

这篇关于44、Flink 的默认窗口剔除器 evictor 代码示例的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1060179

相关文章

HTML5实现的移动端购物车自动结算功能示例代码

《HTML5实现的移动端购物车自动结算功能示例代码》本文介绍HTML5实现移动端购物车自动结算,通过WebStorage、事件监听、DOM操作等技术,确保实时更新与数据同步,优化性能及无障碍性,提升用... 目录1. 移动端购物车自动结算概述2. 数据存储与状态保存机制2.1 浏览器端的数据存储方式2.1.

基于 HTML5 Canvas 实现图片旋转与下载功能(完整代码展示)

《基于HTML5Canvas实现图片旋转与下载功能(完整代码展示)》本文将深入剖析一段基于HTML5Canvas的代码,该代码实现了图片的旋转(90度和180度)以及旋转后图片的下载... 目录一、引言二、html 结构分析三、css 样式分析四、JavaScript 功能实现一、引言在 Web 开发中,

Python如何去除图片干扰代码示例

《Python如何去除图片干扰代码示例》图片降噪是一个广泛应用于图像处理的技术,可以提高图像质量和相关应用的效果,:本文主要介绍Python如何去除图片干扰的相关资料,文中通过代码介绍的非常详细,... 目录一、噪声去除1. 高斯噪声(像素值正态分布扰动)2. 椒盐噪声(随机黑白像素点)3. 复杂噪声(如伪

Java Spring ApplicationEvent 代码示例解析

《JavaSpringApplicationEvent代码示例解析》本文解析了Spring事件机制,涵盖核心概念(发布-订阅/观察者模式)、代码实现(事件定义、发布、监听)及高级应用(异步处理、... 目录一、Spring 事件机制核心概念1. 事件驱动架构模型2. 核心组件二、代码示例解析1. 事件定义

python使用库爬取m3u8文件的示例

《python使用库爬取m3u8文件的示例》本文主要介绍了python使用库爬取m3u8文件的示例,可以使用requests、m3u8、ffmpeg等库,实现获取、解析、下载视频片段并合并等步骤,具有... 目录一、准备工作二、获取m3u8文件内容三、解析m3u8文件四、下载视频片段五、合并视频片段六、错误

HTML5 getUserMedia API网页录音实现指南示例小结

《HTML5getUserMediaAPI网页录音实现指南示例小结》本教程将指导你如何利用这一API,结合WebAudioAPI,实现网页录音功能,从获取音频流到处理和保存录音,整个过程将逐步... 目录1. html5 getUserMedia API简介1.1 API概念与历史1.2 功能与优势1.3

spring中的ImportSelector接口示例详解

《spring中的ImportSelector接口示例详解》Spring的ImportSelector接口用于动态选择配置类,实现条件化和模块化配置,关键方法selectImports根据注解信息返回... 目录一、核心作用二、关键方法三、扩展功能四、使用示例五、工作原理六、应用场景七、自定义实现Impor

nginx启动命令和默认配置文件的使用

《nginx启动命令和默认配置文件的使用》:本文主要介绍nginx启动命令和默认配置文件的使用,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录常见命令nginx.conf配置文件location匹配规则图片服务器总结常见命令# 默认配置文件启动./nginx

mysql中insert into的基本用法和一些示例

《mysql中insertinto的基本用法和一些示例》INSERTINTO用于向MySQL表插入新行,支持单行/多行及部分列插入,下面给大家介绍mysql中insertinto的基本用法和一些示例... 目录基本语法插入单行数据插入多行数据插入部分列的数据插入默认值注意事项在mysql中,INSERT I

PostgreSQL 默认隔离级别的设置

《PostgreSQL默认隔离级别的设置》PostgreSQL的默认事务隔离级别是读已提交,这是其事务处理系统的基础行为模式,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价... 目录一 默认隔离级别概述1.1 默认设置1.2 各版本一致性二 读已提交的特性2.1 行为特征2.2