You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在Apache Flink中使用initializeState()时为何出现NullPointerException?

Hey there, let's break down why you're hitting that NullPointerException and how to fix it right away.

Looking at your code, the root issue is clear: you're trying to work with Operator State (via CheckpointedFunction) but you're accessing the Keyed State Store instead.

Here's the problematic line:

state = context.getKeyedStateStore().getMapState(descriptor);

The getKeyedStateStore() method returns null when your function isn't running on a KeyedStream (meaning you haven't called keyBy() on your stream before applying this function). Calling getMapState() on that null reference triggers the NullPointerException.

Corrected Code

To fix this, switch to using the Operator State Store instead:

public void initializeState(FunctionInitializationContext context) throws Exception {
    MapStateDescriptor<Long, Long> descriptor = new MapStateDescriptor<>(
        "state",
        TypeInformation.of(new TypeHint<Long>() {}),
        TypeInformation.of(new TypeHint<Long>() {})
    );
    // Use Operator State Store for Operator State usage
    state = context.getOperatorStateStore().getMapState(descriptor);
}

Quick Notes to Avoid This in the Future

  • If you actually intended to use Keyed State instead of Operator State, you need to ensure your stream is keyed first. Add a keyBy() operation on your stream before applying this function, and then getKeyedStateStore() will work as expected.
  • For Operator State, Flink supports different distribution modes (like Union, Broadcast, etc.). Depending on your use case, you might need to use methods like getUnionState() instead of getMapState() if you want each parallel instance to hold a full copy of the state.

内容的提问来源于stack exchange,提问作者Weize Sun

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 10:34:18