实战-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)
还没有任何评论哟~
