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

Java并行结合异步处理Kafka消息时出现空指针异常的原因

Java并行流结合@Async处理Kafka消息时出现NullPointerException的原因分析

问题背景

在Java中同时使用并行流和Spring的@Async异步处理Kafka消息时,触发了NullPointerException。尽管异步方法已通过全局try-catch捕获所有异常,但仍出现该错误,需分析异常产生原因。

代码实现

线程池配置

@Bean("threadImport")
public TaskExecutor threadImport() {
    log.info("*** init threadImport with corePoolSize [{}] maxPoolSize [{}]",
            threadImportCorePoolSize, threadImportMaxPoolSize);
    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
    executor.setCorePoolSize(20); /* 20 poolSize */
    executor.setMaxPoolSize(100); /* 100 max poolSize */
    executor.setQueueCapacity(1024000); //1024000 default queue capacity
    executor.setThreadNamePrefix("threadImport-");
    return executor;
}

Kafka消息消费方法

public void consumeT24CustCorpVip(List<String> messages) throws ExecutionException, InterruptedException {
    List<CompletableFuture<BaseResponse>> lstCompletableFutures = new ArrayList<>();
    if (!CollectionUtils.isEmpty(messages)) {
        Gson gson = new GsonBuilder().registerTypeAdapter(LocalDate.class, new LocalDateAdapter()).create(); //create bean Gson custom
        messages.stream().parallel().forEach(message -> {
            lstCompletableFutures.add(executedThreadT24CorpVip.executeT24CustCorpVip(message, gson));
        }); //Parallel processing of messages received when consuming kafka, use
    }
}

异步处理方法

@Async("threadImport") //Use config bean threadImport
public CompletableFuture<BaseResponse> executeT24CustCorpVip(String message, Gson gson) {
    BaseResponse baseResponse = new BaseResponse();
    try {
        T24VpbCustCorpVip t24VpbCustCorpVip = gson.fromJson(message, T24VpbCustCorpVip.class); //convert message to entity
        t24CustCorpVipRepository.save(t24VpbCustCorpVip); //save database
        baseResponse.setStatus(Status.SUCCESS);
    } catch (Exception e) {
        baseResponse.setStatus(Status.ERROR);
    }
    return CompletableFuture.completedFuture(baseResponse); //response CompletableFuture
}

错误堆栈信息

org.springframework.kafka.listener.ListenerExecutionFailedException: Listener method 'public void vn.com.vpbank.corp.api.AirportLoungeController.consumeT24CustCorpVip(java.util.List<java.lang.String>) throws java.util.concurrent.ExecutionException,java.lang.InterruptedException' threw exception; nested exception is java.lang.NullPointerException; nested exception is java.lang.NullPointerException
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.decorateException(KafkaMessageListenerContainer.java:1641)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeBatchErrorHandler(KafkaMessageListenerContainer.java:1387)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeBatchListener(KafkaMessageListenerContainer.java:1274)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeBatchListener(KafkaMessageListenerContainer.java:1179)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeListener(KafkaMessageListenerContainer.java:1162)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:949)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:884)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.run(FutureTask.java:266)
    at java.lang.Thread.run(Thread.java:750)
Caused by: java.lang.NullPointerException: null
    at vn.com.vpbank.corp.service.airportLounge.impl.AirportLoungeServiceImpl.consumeT24CustCorpVip(AirportLoungeServiceImpl.java:352)
    at vn.com.vpbank.corp.service.airportLounge.impl.AirportLoungeServiceImpl$$FastClassBySpringCGLIB$$b3d2718d.invoke(<generated>)
    at org.springframework.cglib.proxy.MethodProxy.invoke(MethodProxy.java:218)
    at org.springframework.aop.framework.CglibAopProxy$DynamicAdvisedInterceptor.intercept(CglibAopProxy.java:685)
    at vn.com.vpbank.corp.service.airportLounge.impl.AirportLoungeServiceImpl$$EnhancerBySpringCGLIB$$7bfb2704.consumeT24CustCorpVip(<generated>)
    at vn.com.vpbank.corp.api.AirportLoungeController.consumeT24CustCorpVip(AirportLoungeController.java:104)
    at sun.reflect.GeneratedMethodAccessor428.invoke(Unknown Source)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:498)
    at org.springframework.messaging.handler.invocation.InvocableHandlerMethod.doInvoke(InvocableHandlerMethod.java:171)
    at org.springframework.messaging.handler.invocation.InvocableHandlerMethod.invoke(InvocableHandlerMethod.java:120)
    at org.springframework.kafka.listener.adapter.HandlerAdapter.invoke(HandlerAdapter.java:48)
    at org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter.invokeHandler(MessagingMessageListenerAdapter.java:301)
    at org.springframework.kafka.listener.adapter.BatchMessagingMessageListenerAdapter.invoke(BatchMessagingMessageListenerAdapter.java:141)
    at org.springframework.kafka.listener.adapter.BatchMessagingMessageListenerAdapter.onMessage(BatchMessagingMessageListenerAdapter.java:133)
    at org.springframework.kafka.listener.adapter.BatchMessagingMessageListenerAdapter.onMessage(BatchMessagingMessageListenerAdapter.java:58)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeBatchOnMessage(KafkaMessageListenerContainer.java:1359)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeBatchOnMessage(KafkaMessageListenerContainer.java:1322)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeBatchListener(KafkaMessageListenerContainer.java:1262)
    ... 7 common frames omitted

异常原因分析

  1. 非线程安全集合的并发操作问题
    代码中使用ArrayList存储异步任务的CompletableFuture,而ArrayList是非线程安全的集合。在并行流的多线程环境下,多个线程同时调用lstCompletableFutures.add()会导致集合内部结构损坏(比如size计数错误、数组元素引用异常),最终触发NullPointerException。

  2. 并行流与@Async的冗余叠加
    @Async已经将任务提交到自定义的threadImport线程池异步执行,再使用并行流(默认使用ForkJoinPool)相当于两层多线程嵌套,不仅造成线程资源浪费,还额外引入了并发安全风险。

  3. 空指针的具体触发点
    根据错误堆栈,NPE出现在consumeT24CustCorpVip方法的第352行,也就是操作lstCompletableFutures的过程中。由于ArrayList在并发add时内部状态不一致,后续对集合的操作(即使代码中暂时没有后续操作,并发add本身也可能触发内部NPE)会抛出空指针。

修复建议

  • 替换为线程安全集合:将ArrayList替换为CopyOnWriteArrayList,或者使用并行流的collect操作安全收集结果:
    List<CompletableFuture<BaseResponse>> lstCompletableFutures = messages.stream()
            .parallel()
            .map(message -> executedThreadT24CorpVip.executeT24CustCorpVip(message, gson))
            .collect(Collectors.toList());
    
  • 移除并行流的冗余:由于@Async已经实现异步处理,直接遍历提交任务即可,无需再使用并行流:
    messages.forEach(message -> {
        lstCompletableFutures.add(executedThreadT24CorpVip.executeT24CustCorpVip(message, gson));
    });
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 04:39:52