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

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

问题分析

  1. 空StResponse导致的NPE:当stDao.getStDetails()调用失败时,onErrorReturn返回了空的StResponse对象,其getSt()方法大概率返回null,后续调用getSts()直接触发NPE;同时fetchData的校验逻辑未覆盖stResponse.getSt()的空值情况。
  2. HashMap.put()返回值导致的NPE:HashMap.put()在新增键值对时会返回null,用Objects.requireNonNull包裹该返回值会直接抛出NPE,对应报错信息"The mapper function returned a null value"。
  3. 未处理的流异常:两处subscribe都只定义了成功回调,没有错误处理逻辑,流中出现异常时会抛出OnErrorNotImplementedException导致程序崩溃。
  4. 变量名笔误: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 02:47:04