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

Flink flatMap()方法出现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 19:05:53