【问题标题】:in flink processFunction, all mapstate is empty in onTimer() function在flink processFunction中,所有mapstate在onTimer()函数中都是空的
【发布时间】:2020-12-11 16:29:07
【问题描述】:

我想通过processKeyedFunction来实现aggregationFunction,因为默认aggregationFunction不支持丰富的功能, 另外,我尝试了aggreagationFunction + processWindowFunction(https://ci.apache.org/projects/flink/flink-docs-stable/dev/stream/operators/windows.html),但它也不能满足我的需求,所以我必须使用基本的processKeyedFunction来实现aggregationFunction,我的问题细节如​​下:

在processFunction中,我定义了一个windowState用于stage元素的聚合值,代码如下:

public void open(Configuration parameters) throws Exception {
followCacheMap = FollowSet.getInstance();
windowState = getRuntimeContext().getMapState(windowStateDescriptor);
currentTimer = getRuntimeContext().getState(new ValueStateDescriptor<Long>(
        "timer",
        Long.class
));

在processElement()函数中,我使用windowState(这是一个MapState在open函数中初始化)聚合窗口元素,并注册第一次timeServie清除当前窗口状态,代码如下:

 @Override
public void processElement(FollowData value, Context ctx, Collector<FollowData> out) throws Exception 
{
      if ( (currentTimer==null || (currentTimer.value() ==null) || (long)currentTimer.value()==0 ) && value.getClickTime() != null) {
            currentTimer.update(value.getClickTime() + interval);
            ctx.timerService().registerEventTimeTimer((long)currentTimer.value());
        } 
       windowState = doMyAggregation(value);
}

在onTimer()函数中,首先我在next 1分钟内注册下一次timeService,并清除窗口状态

 @Override
    public void onTimer(long timestamp, OnTimerContext ctx, Collector<FollowData> out) throws Exception {
        currentTimer.update(timestamp + interval);   // interval is 1 minute
        ctx.timerService().registerEventTimeTimer((long)currentTimer.value());
        
        out.collect(windowState);
        windowState.clear();
     }

但是程序运行的时候发现onTimer里面的windowState都是空的,但是processElement()函数里面却不是空的,不知道为什么会这样,可能是执行逻辑不一样,怎么可能我解决这个, 提前致谢!



新增关于 doMyAggregation() 部分的代码


windowState是一个MapState,key是“mykey”,value是一个自定义的Object AggregateFollow

public class AggregateFollow {
    private String clicked;
    private String unionid;
    private ArrayList allFollows;
    private int enterCnt;
    private Long clickTime;

}

和doMyAggregation(value)函数差不多就是这样,doMyAggregation的作用是获取源字段为'follow'的所有值,但如果1分钟内没有字段为'click'的值, 'follow'值应该是过时的,一句话,就像'follow'数据和'click'数据的join操作,

AggregateFollow acc = windowState.get(windowkey);
    String flag = acc.getClicked();
    ArrayList<FollowData> followDataList = acc.getAllFollows();
    if ("0".equals(flag)) {
        if ("follow".equals(value.getSource())) {
            followDataList.add(value);
            acc.setAllFollows(followDataList);
        }
        if ("click".equals(value.getSource())) {
            String unionid = value.getUnionid();
            clickTime = value.getClickTime();
            if (followDataList.size() > 0) {
                ArrayList listNew = new ArrayList();
                for (FollowData followData : followDataList) {
                    followData.setUnionid(unionid);
                    followData.setClickTime(clickTime);
                    followData.setSource("joined_flag");   // 
                }
                acc.setAllFollows(listNew);
            }
            acc.setClicked("1");
            acc.setUnionid(unionid);
            acc.setClickTime(clickTime);
            windowState.put(windowkey, acc);
        }
    } else if ("1".equals(flag)) {
        if ("follow".equals(value.getSource())) {
            value.setUnionid(acc.getUnionid());
            value.setClickTime(acc.getClickTime());
            value.setSource("joined_flag");  
            followDataList.add(value);
            acc.setAllFollows(followDataList);
            windowState.put(windowkey, acc);
        }
    }

由于性能问题,原来的 windowAPI 对我来说不是一个有效的选择,我认为这里唯一的方法是使用 processFunction + ontimer 和 Guava Cache , 非常感谢

【问题讨论】:

    标签: apache-flink


    【解决方案1】:

    如果windowState 是空的,看看doMyAggregation(value) 在做什么会很有帮助。

    如果没有更多上下文,很难对此进行调试或提出好的替代方案,但out.collect(windowState) 不会按预期工作。您可能想要做的是迭代此 MapState 并将其包含的每个键/值对收集到输出中。

    【讨论】:

    • 你说得对,大卫,我在 out.collect(windowState) 处简化了代码,并更新了 doMyAggregation() 操作部分的代码,
    • 不知是不是因为processElement上下文和onTimer上下文访问windowState的顺序,对flink的基本机制不熟悉,不胜感激!
    • 你怎么知道windowState中实际存储了任何东西,你怎么知道onTimer中是空的?
    • 谢谢,大卫,我在 processElement 和 onTimer 中添加了一些日志,所以我可以知道一些值已存储到 windowState 中。
    【解决方案2】:

    我把windowState的类型从MapState改成了ValueState,问题就解决了,可能是bug什么的,谁能解释一下?

    【讨论】:

    • 我不明白你的代码在做什么,所以很难确定问题是什么。然而,MapState 经常被误解。 MapState 为所涉及的流的每个键都提供了一个单独的哈希映射——因此您最终会得到一个从键到映射的分片映射映射。如果只需要将键映射到值,那么 ValueState 更合适。但是你不应该在值是地图的地方使用 ValueState,因为 MapState 已经过优化,可以更有效地处理这种情况。
    • 一开始我把我的聚合值存储在mapState中,但是我没有把它当成一个map,因为key是一个常量字符串,所以它实际上是一个只有一个key的mapState,所以我感到困惑,因为我认为它应该像 valueState 一样工作
    • 你是对的,只有一个键的 MapState 应该像 ValueState 一样工作。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-06-25
    相关资源
    最近更新 更多