IT数码 购物 网址 头条 软件 日历 阅读 图书馆
TxT小说阅读器
↓语音阅读,小说下载,古典文学↓
图片批量下载器
↓批量下载图片,美女图库↓
图片自动播放器
↓图片自动播放器↓
一键清除垃圾
↓轻轻一点,清除系统垃圾↓
开发: C++知识库 Java知识库 JavaScript Python PHP知识库 人工智能 区块链 大数据 移动开发 嵌入式 开发工具 数据结构与算法 开发测试 游戏开发 网络协议 系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程
数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁
 
   -> 游戏开发 -> flink 学习(十一)Watermark -> 正文阅读

[游戏开发]flink 学习(十一)Watermark


前言

一、时间语义

1、Event Time

????????事件时间,是事件发生时的时间,在数据中带有描述时间的字段,由于从事件发生时到数据处理的过程中会经过不同的时间段,事件发生时间则很好的描述了数据的原始时间。相比其他时间语义,Event Time的事件时间是确定的,可以使用数据中的时间,也可以在数据到达flink之后按照一定的规则生成时间。

2、Ingestion Time

????????摄入时间,是事件到达Flink Souce的时间,是数据进入Flink 的时间。

3、Processing Time

????????处理时间,是flink在算子中处理数据的时间,在每个算子中处理事件可能都不相同,具有不确定性。

二、Watermark

????????Ingestion Time 和 Processing Time 不需要设置Watermark。Event Time则需要设置Watermark。Watermark时间戳是一直递增的,同时可以通过设置延迟来控制准确度,保证乱序的晚来的数据也能够被解析到。

????????Event Time存在于每条数据中,需要通过 WatermarkGenerator 来配置 watermark 的生成方式,让Flink 应用程序知道事件时间戳对应的字段,意味着数据流中的每个元素都需要拥有可分配的事件时间戳。其通常通过使用 TimestampAssigner API 从元素中的某个字段去访问/提取时间戳。

三、AscendingTimestampsWatermarks

时间有序递增时,可以使用AscendingTimestampsWatermarks

(1)开启nc

nc -lp 8888

(2)输入数据

0
10
15
19
20
21
22
25
29
30

(3)示例

@Test
    public void forMonotonousTimestampsTest() throws Exception {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setRuntimeMode(RuntimeExecutionMode.STREAMING)
                .setParallelism(1);
        DataStreamSource<String> source = env.socketTextStream("172.16.10.159", 8888);
        source.map(new MapFunction<String, Long>() {
            @Override
            public Long map(String value) throws Exception {
                return Long.parseLong(value);
            }
        })
                //设置时间戳和水印
                .assignTimestampsAndWatermarks(WatermarkStrategy.<Long>forMonotonousTimestamps().withTimestampAssigner((element, recordTimestamp) -> element))
                //基于事件时间的滚动窗口,时间间隔是10毫秒
                .windowAll(TumblingEventTimeWindows.of(Time.milliseconds(10)))
                .process(new ProcessAllWindowFunction<Long, Long, TimeWindow>() {
                    @Override
                    public void process(Context context, Iterable<Long> elements, Collector<Long> out) throws Exception {
                        Iterator<Long> it = elements.iterator();

                        Long last = null;
                        while (it.hasNext()) {
                            Long next = it.next();

                            last = next;
                            System.out.println("元素: " + next);
                        }
                        out.collect(last);
                    }
                })
                .print("forMonotonousTimestamps")
        ;
        env.execute("Watermark");
    }

(4)结果

元素: 0
forMonotonousTimestamps> 0
元素: 10
元素: 15
元素: 19
forMonotonousTimestamps> 19
元素: 20
元素: 21
元素: 22
元素: 25
元素: 29
forMonotonousTimestamps> 29

可以看出,当时间相差10时进行一次输出,
输入10时,10和0相差为10,则输出 forMonotonousTimestamps> 0

四、BoundedOutOfOrdernessWatermarks

允许一定程度的乱序

(1)开启nc

nc -lp 8888

(2)示例

 @Test
    public void forBoundedOutOfOrdernessTest() throws Exception {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setRuntimeMode(RuntimeExecutionMode.STREAMING)
                .setParallelism(1);
        DataStreamSource<String> source = env.socketTextStream("172.16.10.159", 8888);
        source.map(new MapFunction<String, Long>() {
            @Override
            public Long map(String value) throws Exception {
                return Long.parseLong(value);
            }
        })
                //设置时间戳和水印
                .assignTimestampsAndWatermarks(WatermarkStrategy.<Long>forBoundedOutOfOrderness(Duration.ofMillis(2)).withTimestampAssigner((element, recordTimestamp) -> element))
                //基于事件时间的滚动窗口,时间间隔10毫秒,允许乱序时间2毫秒
                .windowAll(TumblingEventTimeWindows.of(Time.milliseconds(10)))
                .process(new ProcessAllWindowFunction<Long, Long, TimeWindow>() {
                    @Override
                    public void process(Context context, Iterable<Long> elements, Collector<Long> out) throws Exception {
                        Iterator<Long> it = elements.iterator();

                        Long last = null;
                        while (it.hasNext()) {
                            Long next = it.next();

                            last = next;
                            System.out.println("元素: " + next);
                        }
                        out.collect(last);
                    }
                })
                .print("forMonotonousTimestamps")
        ;
        env.execute("WindowFunction");
    }

(3)输入数据

0
10
11
12
20
22

当输入到12的时候,12与0的距离达到时间间隔10毫秒加上允许乱序时间2毫秒,即10 + 2,这是会输出0~10之间的数据
在这里插入图片描述

当输入到22的时候,同样触发输出10~20之间的数据

在这里插入图片描述

  游戏开发 最新文章
6、英飞凌-AURIX-TC3XX: PWM实验之使用 GT
泛型自动装箱
CubeMax添加Rtthread操作系统 组件STM32F10
python多线程编程:如何优雅地关闭线程
数据类型隐式转换导致的阻塞
WebAPi实现多文件上传,并附带参数
from origin ‘null‘ has been blocked by
UE4 蓝图调用C++函数(附带项目工程)
Unity学习笔记(一)结构体的简单理解与应用
【Memory As a Programming Concept in C a
上一篇文章      下一篇文章      查看所有文章
加:2022-04-23 11:07:26  更:2022-04-23 11:08:15 
 
开发: C++知识库 Java知识库 JavaScript Python PHP知识库 人工智能 区块链 大数据 移动开发 嵌入式 开发工具 数据结构与算法 开发测试 游戏开发 网络协议 系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程
数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁

360图书馆 购物 三丰科技 阅读网 日历 万年历 2025年1日历 -2025/1/16 21:46:51-

图片自动播放器
↓图片自动播放器↓
TxT小说阅读器
↓语音阅读,小说下载,古典文学↓
一键清除垃圾
↓轻轻一点,清除系统垃圾↓
图片批量下载器
↓批量下载图片,美女图库↓
  网站联系: qq:121756557 email:121756557@qq.com  IT数码