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
异常原因分析
非线程安全集合的并发操作问题
代码中使用ArrayList存储异步任务的CompletableFuture,而ArrayList是非线程安全的集合。在并行流的多线程环境下,多个线程同时调用lstCompletableFutures.add()会导致集合内部结构损坏(比如size计数错误、数组元素引用异常),最终触发NullPointerException。并行流与@Async的冗余叠加
@Async已经将任务提交到自定义的threadImport线程池异步执行,再使用并行流(默认使用ForkJoinPool)相当于两层多线程嵌套,不仅造成线程资源浪费,还额外引入了并发安全风险。空指针的具体触发点
根据错误堆栈,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
相关产品推荐
相关产品推荐

