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

在keyBy中使用System.currentTimeMillis()导致KeyedProcessFunction抛出NPE

我正在开发一个处理Row元素并通过KeyedProcessFunction实现延迟逻辑的Flink作业,遇到了KeyBy键生成方式的问题:

  • 当直接使用System.currentTimeMillis()生成键时,执行bufferedElements.add()向ListState添加元素会抛出NullPointerException(NPE),原因是setCurrentKeyGroupIndex计算出的键组索引为78,超出当前任务分配的52-76范围;
  • 但使用Row中包含的System.currentTimeMillis()值字段作为键时,状态管理运行正常。

以下是简化代码:

public static void main(String[] args) {
    // Define row types
    RowTypeInfo rowTypeInfo = new RowTypeInfo(
        new TypeInformation[]{Types.STRING, Types.STRING, Types.LONG, Types.LONG}, 
        new String[]{"src_ip", "port", "time", "time2"});

    // Set up the environment
    StreamExecutionEnvironment environment = StreamExecutionEnvironment.getExecutionEnvironment();
    environment.setParallelism(5);

    // Source function
    environment.addSource(new SourceFunction<Row>() {
        @Override
        public void run(SourceContext<Row> ctx) throws Exception {
            for (int i = 0; i < 1000; i++) {
                ctx.collect(Row.of("first_" + 1, "456", System.currentTimeMillis(), System.currentTimeMillis()));
                Thread.sleep(50);
            }
        }

        @Override
        public void cancel() {
        }
    }).returns(rowTypeInfo)
      .keyBy(new KeySelector<Row, String>() {
          @Override
          public String getKey(Row value) throws Exception {
              // Causes NPE
              return String.valueOf(System.currentTimeMillis());
              
              // Works fine
              // return String.valueOf(value.getField(3));
          }
      })
      .process(new DelayFunction(5000))
      .map(new MapFunction<Row, Row>() {
          @Override
          public Row map(Row value) throws Exception {
              return value;
          }
      });

    try {
        environment.execute();
    } catch (Exception e) {
        throw new RuntimeException(e);
    }
}

// Delay function
public static class DelayFunction extends KeyedProcessFunction<String, Row, Row> {
    private final long delayTime;
    private transient ListState<Row> bufferedElements;

    public DelayFunction(long delayTime) {
        this.delayTime = delayTime;
    }

    @Override
    public void open(Configuration parameters) throws Exception {
        ListStateDescriptor<Row> descriptor = new ListStateDescriptor<>("bufferedElements", Row.class);
        bufferedElements = getRuntimeContext().getListState(descriptor);
    }

    @Override
    public void processElement(Row value, Context ctx, Collector<Row> out) throws Exception {
        long currentTime = ctx.timerService().currentProcessingTime();
        long triggerTime = currentTime + delayTime;
        value.setField(0, new SimpleDateFormat("yyyy-MM-dd HH:mm:ss SSS").format(triggerTime));
        bufferedElements.add(value); // NPE occurs here
        ctx.timerService().registerProcessingTimeTimer(triggerTime);
    }
}

原因分析

1. 键的一致性问题

当你在KeySelector中直接调用System.currentTimeMillis()时,这个值是KeyBy算子执行时实时计算的动态值,而非来自Row元素本身的固定属性。这会导致:

  • 元素在KeyBy阶段被分配任务槽时,基于此时生成的键计算键组索引;
  • 进入ProcessFunction操作状态时,Flink会基于当前元素对应的键(此时KeySelector可能已生成新的时间戳)再次计算键组索引,两次计算的键不一致,最终导致键组索引超出当前任务负责的范围。

此时Flink无法找到对应键组的状态分区,bufferedElements对象未被正确初始化,调用add()方法就会抛出NPE。

而使用Row中的value.getField(3)作为键时,这个值是从元素中提取的固定值,整个处理流程中键始终稳定:元素分配任务槽和操作状态时使用的是同一个键,键组索引始终在当前任务的分配范围内,状态对象能正常绑定并初始化。

KeyBy的核心作用是将具有相同业务属性的元素分组,让相同组的元素被同一个任务处理。而用实时生成的时间戳作为键,相当于每个元素的键都几乎唯一,不仅会导致状态管理异常,还会让每个元素都被分到独立的组,完全丧失并行处理的优势,同时会产生大量状态分区,严重影响作业性能。

内容的提问来源于stack exchange,提问作者jd g

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 12:17:03