RxJava异步调用并行执行遇NullPointerException问题求助
RxJava批量并行处理中的NullPointerException问题解决
问题描述
作为RxJava新手,我实现了基于两次Single异步调用的逻辑:先调用外部服务,再根据其响应调用另一外部服务。需求是对数据批量并行处理(如100条数据分为4批,每批25条并行执行),但当前实现出现了NullPointerException(NPE),同时存在未处理的异常导致程序崩溃。
核心实现代码
第一次API调用方法
private Map<String, List<List<String>>> firstAPICall( List<STRequest> stList, int batchSize, int numberOfBatches, final String mainTraceId ) { log.debug("Main Trace ID: {} Initiating firstAPICall", mainTraceId); CountDownLatch countDownLatch = new CountDownLatch(1); //TODO : update the thread count as per the tps supported by downstream service final List<Integer> batchList = new ArrayList<>(); for (int i = 1; i <= numberOfBatches; i++) { batchList.add(i); } Map<String, List<List<String>>> stMap = new HashMap<>(); int parallelThreadsCount = batchList.size(); Disposable disposable = Flowable.fromIterable(batchList) .parallel(parallelThreadsCount) .runOn(Schedulers.io()) .map(batchIndex -> fetchStoreSubList(storeList, batchSize, batchIndex)) // 注意:参数为storeList,但方法入参是stList,存在笔误 .flatMapIterable(stList -> stList) .flatMap(stRequest -> { String traceId = UUID.randomUUID().toString(); return stDao.getStDetails( stRequest.getStUUID(), stRequest.isDetailsFlag(), stRequest.getCode()) .onErrorReturn(throwable -> { log.error("Trace: {} Error occured {}", traceId, throwable.getMessage()); return new StResponse(); // 返回空对象,后续会触发NPE }) .map(stResponse -> { // 问题点1:stResponse.getSt()可能为null,调用getSts()触发NPE stResponse.getSt().getSts().removeIf(i -> i.getNumber().equalsIgnoreCase(stRequest.getStNumber())); return stResponse; }) .flatMap(stResponse -> Single.just(fetchData(stResponse, traceId))) // 调用第二个API .toFlowable(); }) .collect(Collectors.toList()) .subscribe(responseList -> { responseList.forEach(stMap::putAll); countDownLatch.countDown(); }); // 问题点:缺少错误回调,异常未被处理 try { countDownLatch.await(); } catch (Exception e) { log.error("Trace: {} Error while waiting for count down", mainTraceId); } disposable.dispose(); return stMap; }
第二次API调用方法
private Map<String, List<List<String>>> fetchData (StResponse stResponse, final String traceId) { CountDownLatch countDownLatch = new CountDownLatch(1); log.debug("CRdata Trace id: {} Initiating fetchCRData for", traceId); if (ObjectUtils.isEmpty(stResponse) || ObjectUtils.isEmpty(stResponse.getStore()) || ObjectUtils.isEmpty(stResponse.getStore().getStores())) { log.error("Store Details not available - {}", traceId); throw new StDetailsUnavailableException("St Details not available for st"); } List<St> stList = stResponse.getSt().getSts(); // 此处也可能因stResponse.getSt()为null触发NPE Map<String, List<List<String>>> stDataMap = new HashMap<>(); Disposable disposable = Flowable.fromIterable(storeList) .flatMap(stInfo -> { CRDetail detail = stInfo.getDetail(); String spScheduledCount = "SP"; String stream = "P"; return dataDao.getDept(spScheduledCount, stream, "1111", detail.getData(), DEFAULT, "IE" ).flatMap(resp -> { if(ObjectUtils.isEmpty(resp.getCrDetail()) || ObjectUtils.isEmpty(resp.getCrDetail().getCrData())) { log.error("CR data not available for store - {}", stInfo.getNumber()); throw new CRDepartmentDetailsUnavailableException("CR data not available for store"); } // 问题点2:HashMap.put()新增键值对时返回null,Objects.requireNonNull触发NPE return Single.just(Objects.requireNonNull(stDataMap.put(stInfo.getId(), resp.getDetail().getData()))); }).onErrorReturn(throwable -> { log.error("Exception occurred while fetching data - {} {}",stInfo.getNumber(), throwable.getMessage()); return new ArrayList<>(); }).toFlowable(); }) .collect(Collectors.toList()) .subscribe(responseList -> { countDownLatch.countDown(); }); // 缺少错误回调 try { countDownLatch.await(); } catch (Exception e) { log.error("Trace: {} Error while waiting for count down", traceId); } disposable.dispose(); return stDataMap; }
错误堆栈信息
Caused by: java.lang.NullPointerException: 无法调用"StDetail.getSt()",因为"StDetail.getSt()"的返回值为null at com.tesco.store.counts.scheduler.service.impl.TargetedCountServiceImpl.lambda$fetchStoreDetails$4(MyServiceImpl.java:142) at io.reactivex.rxjava3.internal.operators.single.SingleFlatMap$SingleFlatMapCallback.onSuccess(SingleFlatMap.java:77) 线程"RxCachedThreadScheduler-6"中的异常: io.reactivex.rxjava3.exceptions.OnErrorNotImplementedException: 由于subscribe()方法调用中缺少onError处理程序,异常未被处理。 | java.lang.NullPointerException: 无法调用"StDetail.getSts()",因为"StResponse.getSt()"的返回值为null
问题分析
- 空StResponse导致的NPE:当
stDao.getStDetails()调用失败时,onErrorReturn返回了空的StResponse对象,其getSt()方法大概率返回null,后续调用getSts()直接触发NPE;同时fetchData的校验逻辑未覆盖stResponse.getSt()的空值情况。 - HashMap.put()返回值导致的NPE:
HashMap.put()在新增键值对时会返回null,用Objects.requireNonNull包裹该返回值会直接抛出NPE,对应报错信息"The mapper function returned a null value"。 - 未处理的流异常:两处
subscribe都只定义了成功回调,没有错误处理逻辑,流中出现异常时会抛出OnErrorNotImplementedException导致程序崩溃。 - 变量名笔误:
firstAPICall中调用fetchStoreSubList时传入了未定义的storeList,实际应使用方法入参stList。
解决方案
1. 修复空StResponse的NPE
修改onErrorReturn,返回包含空St对象的StResponse(需根据实际构造方法调整):
.onErrorReturn(throwable -> { log.error("Trace: {} Error occured {}", traceId, throwable.getMessage()); St emptySt = new St(new ArrayList<>()); return new StResponse(emptySt); })
在map操作前增加空值校验:
.map(stResponse -> { if (stResponse.getSt() == null || stResponse.getSt().getSts() == null) { log.warn("Trace: {} Invalid St data in response", traceId); return stResponse; } stResponse.getSt().getSts().removeIf(i -> i.getNumber().equalsIgnoreCase(stRequest.getStNumber())); return stResponse; })
补充fetchData的空值校验:
if (ObjectUtils.isEmpty(stResponse) || ObjectUtils.isEmpty(stResponse.getStore()) || ObjectUtils.isEmpty(stResponse.getStore().getStores()) || ObjectUtils.isEmpty(stResponse.getSt()) || ObjectUtils.isEmpty(stResponse.getSt().getSts())) { log.error("Store Details not available - {}", traceId); throw new StDetailsUnavailableException("St Details not available for st"); }
2. 修复HashMap.put()的NPE
去掉Objects.requireNonNull,直接执行put并返回实际数据:
flatMap(resp -> { if(ObjectUtils.isEmpty(resp.getCrDetail()) || ObjectUtils.isEmpty(resp.getCrDetail().getCrData())) { log.error("CR data not available for store - {}", stInfo.getNumber()); throw new CRDepartmentDetailsUnavailableException("CR data not available for store"); } List<List<String>> data = resp.getDetail().getData(); stDataMap.put(stInfo.getId(), data); return Single.just(data); })
3. 补充错误处理回调
在所有subscribe中添加错误回调,避免未捕获异常:
// firstAPICall中的subscribe .subscribe(responseList -> { responseList.forEach(stMap::putAll); countDownLatch.countDown(); }, throwable -> { log.error("Trace: {} Error in firstAPICall flow", mainTraceId, throwable); countDownLatch.countDown(); // 确保latch正常释放 }) // fetchData中的subscribe .subscribe(responseList -> { countDownLatch.countDown(); }, throwable -> { log.error("Trace: {} Error in fetchData flow", traceId, throwable); countDownLatch.countDown(); })
4. 修正变量名笔误
将fetchStoreSubList(storeList, batchSize, batchIndex)改为fetchStoreSubList(stList, batchSize, batchIndex),确保使用正确的数据源。
5. 优化同步方式(可选)
用RxJava原生的blockingGet()替代CountDownLatch,简化代码:
private Map<String, List<List<String>>> firstAPICall( List<STRequest> stList, int batchSize, int numberOfBatches, final String mainTraceId ) { log.debug("Main Trace ID: {} Initiating firstAPICall", mainTraceId); final List<Integer> batchList = new ArrayList<>(); for (int i = 1; i <= numberOfBatches; i++) { batchList.add(i); } Map<String, List<List<String>>> stMap = new HashMap<>(); int parallelThreadsCount = batchList.size(); try { List<Map<String, List<List<String>>>> responseList = Flowable.fromIterable(batchList) .parallel(parallelThreadsCount) .runOn(Schedulers.io()) .map(batchIndex -> fetchStoreSubList(stList, batchSize, batchIndex)) .flatMapIterable(subList -> subList) .flatMap(stRequest -> { String traceId = UUID.randomUUID().toString(); return stDao.getStDetails( stRequest.getStUUID(), stRequest.isDetailsFlag(), stRequest.getCode()) .onErrorReturn(throwable -> { log.error("Trace: {} Error occured {}", traceId, throwable.getMessage()); St emptySt = new St(new ArrayList<>()); return new StResponse(emptySt); }) .filter(stResponse -> stResponse.getSt() != null && !stResponse.getSt().getSts().isEmpty()) .map(stResponse -> { stResponse.getSt().getSts().removeIf(i -> i.getNumber().equalsIgnoreCase(stRequest.getStNumber())); return stResponse; }) .flatMap(stResponse -> Single.just(fetchData(stResponse, traceId))) .toFlowable(); }) .collect(Collectors.toList()) .blockingGet(); responseList.forEach(stMap::putAll); } catch (Exception e) { log.error("Trace: {} Error in firstAPICall", mainTraceId, e); } return stMap; }
内容的提问来源于stack exchange,提问作者Naveen Yadav
相关产品推荐
相关产品推荐

