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

Spring Webflux返回MySQL大数据量时JSON格式异常求助

Spring Webflux流式返回大数量JSON时随机格式损坏问题排查

开发一个API需从MySQL返回50万至700万条数据,采用Spring Webflux流式处理方案,逻辑无错误,但返回的JSON数据会随机出现格式损坏问题。

当前实现逻辑:先获取数据库总条数,内部分页,通过线程异步获取数据,再将线程结果重组拼接至输出流。怀疑流数据拼接环节处理不当导致数据重叠,进而引发格式问题,附上相关源码及异常响应截图。

Service类源码:

@Override
public Flux<ObjectReturn> getAllFromDifferentiel() throws InterruptedException, ExecutionException {
    logger.info("Appel de la methode getAllFromDifferentiel()");
    //Initialisation du flux de sortie
    returnedFlux = Flux.empty();

    //Recuperation du total de ligne
    Mono<Integer> totalRecords = repository.count();
    int total = totalRecords.toFuture().get();
    
    //Debut des calculs
    int relicat = total % limit;
    totalPage = (total / limit) + (relicat > 0 ? 1 : 0);

    logger.info("Total page possible: " + totalPage);
    Thread[] tasks = new Thread[totalPage];

    for (int indexPage = 0; indexPage < totalPage; indexPage++) {
      
        //Lancement en parallèle des requetes de pagination

        int finalIndexPage = indexPage;

        tasks[indexPage] = new Thread(() -> {

            offset = finalIndexPage == 0 ? 0 : limit * finalIndexPage + 1;
            logger.info("Page = " + (finalIndexPage + 1) + "      Offset =" + offset);

            //Regroupement des flux générés apres execution de la requete
            returnedFlux = Flux.concat(returnedFlux , repository.getAllByPage(limit, offset).map(ObjectReturn.INSTANCE::mapToDTO));

        });
        tasks[indexPage].start();

        //Agencement ordonné des tâches
        try {
            tasks[indexPage].join();

        } catch (InterruptedException e) {
            throw new RuntimeException(e);
        }
    }

    return returnedFlux 
            .doOnError(error -> {
                logger.error("--------- Stream error ------------" + error.getMessage(), error);
                throw new IllegalStateException(error);
            })
            .doOnCancel(() -> {
                logger.info("--------- Stream canceled -------------");
            })
            .doOnComplete(() -> {
                logger.info("--------- Stream completed -------------");
            });

}

Controller类源码:

@GetMapping(value = "getData", produces = MediaType.APPLICATION_NDJSON_VALUE)
public Flux<ObjectReturn> getAllFromDifferentiel() throws InterruptedException, ExecutionException {
    return service.getAllFromDifferentiel();
}

API调用结果:

JSON格式损坏截图

VSCode中打开的JSON显示,部分数据格式损坏,该问题在文件多处出现。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 10:25:11