Advertisement

实战-07flink平台自定义触发器开发实现count和timeout功能

阅读量:

有用的实战功能,搭配了适度源码讲解

    • 背景
      • 代码
      • 番外篇:探讨在eventTime存在延迟数据的情况下,何时应当对窗口内的数据进行清理

背景

假设我们需依据processTime对数据进行处理,并采用时长为5分钟的滑动窗口机制,其伪代码示例如下

复制代码
    window(TumblingProcessingTimeWindows.of(Time.minites(10)))
                .process(new ProcessWindowFunction<Tuple2<String, Long>, String, String, TimeWindow>() {
                    @Override
                    public void process(String s, Context context, Iterable<Tuple2<String, Long>> elements, Collector<String> out) throws Exception {
                        out.collect(elements.toString());
                    }

全部评论 (0)

还没有任何评论哟~