Flink flatMap()方法出现NullPointerException问题排查
问题:Flink 1.17.1去重程序中
seen.value()触发NullPointerException的原因排查 在使用Flink 1.17.1、Java 11开发数据去重程序时,自定义FilterDuplicate(继承RichFlatMapFunction)的flatMap方法中,执行!seen.value()时触发NullPointerException,已确认seen对象本身不为空,但问题仍存在。
相关代码
主程序代码
public class VerifyDuplicate { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); SingleOutputStreamOperator<Tuple3<String, String, Integer>> dataStream = env.fromElements( Tuple3.of("ID_1", "subid_1", 1), Tuple3.of("ID_2", "subid_2", 2), Tuple3.of("ID_3", "subid_3", 3), Tuple3.of("ID_4", "subid_4", 4), Tuple3.of("ID_4", "subid_4", 4), Tuple3.of("ID_6", "subid_6", 6), Tuple3.of("ID_4", "subid_7", 7), Tuple3.of("ID_8", "subid_8", 8), Tuple3.of("ID_9", "subid_9", 9), Tuple3.of("ID_10", "subid_10", 10) ); KeyedStream<Tuple3<String, String, Integer>, String> partitionedStream = dataStream.keyBy(new KeySelector<Tuple3<String, String,Integer>, String>() { @Override public String getKey(Tuple3<String, String, Integer> value) throws Exception { return value.f0; // partition on f0 } }); partitionedStream.keyBy(new KeySelector<Tuple3<String, String,Integer>, String>() { @Override public String getKey(Tuple3<String, String, Integer> value) throws Exception { return value.f1; // subid } }).flatMap(new FilterDuplicate()).print(); env.execute("Test"); } }
FilterDuplicate实现代码
public class FilterDuplicate extends RichFlatMapFunction<Tuple3<String, String, Integer>, Tuple3<String, String, Integer>> { private ValueState<Boolean> seen; @Override public void open(Configuration configuration) { StateTtlConfig ttlConfig = StateTtlConfig .newBuilder(Time.seconds(15)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .cleanupFullSnapshot() .build(); ValueStateDescriptor<Boolean> desc = new ValueStateDescriptor<>("seen", Types.BOOLEAN); desc.enableTimeToLive(ttlConfig); seen = getRuntimeContext().getState(desc); } @Override public void flatMap(Tuple3<String, String, Integer> value, Collector<Tuple3<String, String, Integer>> out) throws Exception { if (!seen.value()) { // nullpointerexception // we haven't seen the element yet out.collect(value); // set operator state to true so that we don't emit elements with this key again seen.update(true); } } }
依赖版本
<flink.version>1.17.1</flink.version> <target.java.version>11</target.java.version>
原因分析
- 核心原因:Flink的
ValueState在未被update过的情况下,value()方法返回的是null,而非布尔类型的默认值false。直接对null执行取非操作(!null)会触发NullPointerException。 - 虽然
seen对象本身已通过getRuntimeContext().getState(desc)完成初始化,但状态容器内的实际值在首次访问时是未设置的,默认即为null。 - 额外说明:代码中连续两次
keyBy的操作不影响本次NPE的触发,但会改变去重的粒度——最终是按subid(f1字段)维度进行去重,若不符合业务预期需调整,但与当前异常无关。
解决方案
提供两种可行的修复方案:
方案1:判断时先检查状态值是否为null
在flatMap方法中先获取状态值,判断其是否为null,再处理未见过的元素:
@Override public void flatMap(Tuple3<String, String, Integer> value, Collector<Tuple3<String, String, Integer>> out) throws Exception { Boolean isSeen = seen.value(); // 首次访问时isSeen为null,视为未见过该元素 if (isSeen == null || !isSeen) { out.collect(value); seen.update(true); } }
方案2:给ValueStateDescriptor设置默认值
在创建ValueStateDescriptor时指定默认值为false,这样首次调用seen.value()时会返回默认值,避免null:
@Override public void open(Configuration configuration) { StateTtlConfig ttlConfig = StateTtlConfig .newBuilder(Time.seconds(15)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .cleanupFullSnapshot() .build(); // 添加默认值false ValueStateDescriptor<Boolean> desc = new ValueStateDescriptor<>("seen", Types.BOOLEAN, false); desc.enableTimeToLive(ttlConfig); seen = getRuntimeContext().getState(desc); }
内容的提问来源于stack exchange,提问作者overexchange
相关产品推荐
相关产品推荐

