【发布时间】:2019-09-21 18:50:36
【问题描述】:
我是 Flink 新手,有一个我不知道如何处理的用例。
我有活动要来
{
"id" : "AAA",
"event" : "someEvent",
"eventTime" : "2019/09/14 14:04:25:235"
}
我想创建一个表(在弹性/oracle 中)来跟踪用户的不活动。
标识 ||最后一个事件 ||上次事件时间 ||不活动时间
我的最终目标是在某些用户组的活动时间超过 X 分钟时发出警报。
此表应每 1 分钟更新一次。
我不知道我所有的身份证。新 ID 随时可能出现。
我想也许只是使用简单的流程函数来发出事件(如果存在)或发出时间戳(这将更新非活动列)。
问题
- 关于我的解决方案 - 我仍然需要另一段代码来检查事件是否为空并相应更新。如果为 null --> 更新不活动。否则更新 lastEvent。 这段代码可以/应该在同一个 flink/spark 作业中使用吗?
如何处理新的 id?
-
另外,如何在 spark 结构化流中处理这个用例?
input .keyBy("id") .window(TumblingEventTimeWindows.of(Time.minutes(1))) .process(new MyProcessWindowFunction()); public class MyProcessWindowFunction extends ProcessWindowFunction<Tuple2<String, Long>, Tuple2<Long, Object>> { @Override public void process(String key, Context context, Iterable<Tuple2<String, Long>> input, Collector<Tuple2<Long, Object>> out) { Object obj = null; while(input.iterator().hasNext()){ obj = input.iterator().next(); } if (obj!=null){ out.collect(Tuple2.of(context.timestamp(), obj)); } else { out.collect(Tuple2.of(context.timestamp(), null)); } }
【问题讨论】:
-
至于 “另外,如何在 spark 结构化流中处理这个用例?” 我会问一个单独的问题。
标签: apache-flink