在Apache Flink中使用initializeState()时为何出现NullPointerException?
Fixing NullPointerException When Initializing MapState with CheckpointedFunction in Flink
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 thengetKeyedStateStore()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 ofgetMapState()if you want each parallel instance to hold a full copy of the state.
内容的提问来源于stack exchange,提问作者Weize Sun
相关产品推荐
相关产品推荐

