【问题标题】:flink / spark stream - track user inactivityflink / spark stream - 跟踪用户不活动
【发布时间】:2019-09-21 18:50:36
【问题描述】:

我是 Flink 新手,有一个我不知道如何处理的用例。

我有活动要来

{
"id" : "AAA",
"event" : "someEvent",
"eventTime" : "2019/09/14 14:04:25:235"
}

我想创建一个表(在弹性/oracle 中)来跟踪用户的不活动。

标识 ||最后一个事件 ||上次事件时间 ||不活动时间

我的最终目标是在某些用户组的活动时间超过 X 分钟时发出警报。

此表应每 1 分钟更新一次。

我不知道我所有的身份证。新 ID 随时可能出现。

我想也许只是使用简单的流程函数来发出事件(如果存在)或发出时间戳(这将更新非活动列)。

问题

  1. 关于我的解决方案 - 我仍然需要另一段代码来检查事件是否为空并相应更新。如果为 null --> 更新不活动。否则更新 lastEvent。 这段代码可以/应该在同一个 flink/spark 作业中使用吗?
  2. 如何处理新的 id?

  3. 另外,如何在 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


【解决方案1】:

我会使用KeyedProcessFunction 而不是 Windowing API 来满足这些要求。 [1] 流由 id 键控。

KeyedProcessFunction#process 为流的每条记录调用,您可以保持状态和调度计时器。您可以每分钟安排一个计时器,并为每个 od 存储状态中看到的最后一个事件。当计时器触发时,您可以发出事件并清除状态。

就个人而言,我只会存储在数据库中看到的最后一个事件,并在查询数据库时计算不活动时间。这样,您可以在每次发射后清除状态,并且可能无限的键空间不会导致 Flink 中的每个托管状态都在增长。

希望这会有所帮助。

[1]https://ci.apache.org/projects/flink/flink-docs-stable/dev/stream/operators/process_function.html

【讨论】:

  • 关于 Spark Streaming:我绝不是 Structured Streaming 专家,但据我所知,您希望根据单个事件(更新状态、计时器、发射)采取行动的用例是Spark Structured Streaming 本身难以实现和/或效率低下。
  • 仅在查询很聪明但不适用于我的用例时才计算不活动。因为我想在有人不活动超过 X 分钟时发出警报。另外,您能否用触发器和状态给出您的想法的最小可能示例。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-03-22
  • 2023-03-09
  • 2018-01-31
  • 2011-01-31
  • 2020-01-04
  • 1970-01-01
相关资源
最近更新 更多