• 大数据-玩转数据-Flink定时器


    一、说明

    基于处理时间或者事件时间处理过一个元素之后, 注册一个定时器, 然后指定的时间执行.
    Context和OnTimerContext所持有的TimerService对象拥有以下方法:
    currentProcessingTime(): Long 返回当前处理时间
    currentWatermark(): Long 返回当前watermark的时间戳
    registerProcessingTimeTimer(timestamp: Long): Unit 会注册当前key的processing time的定时器。当processing time到达定时时间时,触发timer。
    registerEventTimeTimer(timestamp: Long): Unit 会注册当前key的event time 定时器。当水位线大于等于定时器注册的时间时,触发定时器执行回调函数。
    deleteProcessingTimeTimer(timestamp: Long): Unit 删除之前注册处理时间定时器。如果没有这个时间戳的定时器,则不执行。
    deleteEventTimeTimer(timestamp: Long): Unit 删除之前注册的事件时间定时器,如果没有此时间戳的定时器,则不执行。

    二、基于处理时间的定时器

    package com.lyh.flink08;
    
    import com.lyh.bean.WaterSensor;
    import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
    import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
    import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
    import org.apache.flink.util.Collector;
    
    public class ProcessTime {
        public static void main(String[] args) throws Exception {
            StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
            env.setParallelism(1);
            SingleOutputStreamOperator<WaterSensor> stream = env.socketTextStream("hadoop100", 9999)
                    .map(line -> {
                        String[] datas = line.split(",");
                        return new WaterSensor(datas[0],
                                Long.valueOf(datas[1]),
                                Integer.valueOf(datas[2]));
    
                    });
            stream.keyBy(WaterSensor::getId)
                    .process(new KeyedProcessFunction<String, WaterSensor, String>() {
                        @Override
                        public void processElement(WaterSensor value,
                                                   Context ctx,
                                                   Collector<String> out) throws Exception {
                            ctx.timerService().registerProcessingTimeTimer(ctx.timerService().currentProcessingTime() + 5000);
                            out.collect(value.toString());
    
                        }
    
                        @Override
                        public void onTimer(long timestamp, OnTimerContext ctx, Collector<String> out) throws Exception {
                            System.out.println(timestamp);
                            out.collect("wo be chu fa le ");
                        }
                    }).print();
            env.execute();
        }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40

    三、基于事件时间的定时器

    package com.lyh.flink08;
    
    import com.lyh.bean.WaterSensor;
    import org.apache.flink.api.common.eventtime.WatermarkStrategy;
    import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
    import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
    import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
    import org.apache.flink.util.Collector;
    
    import java.time.Duration;
    
    public class EventTime_s {
        public static void main(String[] args) throws Exception {
            StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
            env.setParallelism(1);
            SingleOutputStreamOperator<WaterSensor> stream = env.socketTextStream("hadoop100", 9999)
                    .map(line -> {
                        String[] datas = line.split(",");
                        return new WaterSensor(
                                datas[0],
                                Long.valueOf(datas[1]),
                                Integer.valueOf(datas[2]));
                    });
    
                WatermarkStrategy<WaterSensor> wms = WatermarkStrategy
                .<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(3))
                .withTimestampAssigner((element,recordTimestamp) -> element.getTs() * 1000);
                stream.assignTimestampsAndWatermarks(wms)
                        .keyBy(WaterSensor::getId)
                        .process(new KeyedProcessFunction<String, WaterSensor, String>() {
                            @Override
                            public void processElement(WaterSensor value,
                                                       Context ctx,
                                                       Collector<String> out) throws Exception {
                                System.out.println(ctx.timestamp());
                                ctx.timerService().registerProcessingTimeTimer(ctx.timestamp()+5000);
                                out.collect(value.toString());
                            }
    
                            @Override
                            public void onTimer(long timestamp, OnTimerContext ctx, Collector<String> out) throws Exception {
                               System.out.println("定时器被触发了");
                            }
                        }).print();
                env.execute();
        }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
  • 相关阅读:
    吴恩达机器学习课程笔记1-2
    算法通关村第十五关:青铜-用4KB内存寻找重复元素
    全球电梯空气消毒机行业调研及趋势分析报告
    【c++】cpp类和对象
    数据结构实战开发教程(三)线性表的本质和操作、顺序存储结构的抽象实现、顺序存储线性表的分析、数组类的创建
    2024年06月在线IDE流行度最新排名
    用Abp实现找回密码和密码强制过期策略
    在PostgreSQL中如何实现递归查询,例如使用WITH RECURSIVE构建层次结构数据?
    BOSS直聘新财报:用户、技术两手抓
    公众号查题:一步一步搭建属于自己的查题
  • 原文地址:https://blog.csdn.net/s_unbo/article/details/132628077