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

如何在Spring Boot中用Flux实现并行调用优化Token生成?

Spring Boot + Flux 实现Token流并行优化方案

核心实现逻辑

  • 并行触发Server1和Server2的异步调用,避免串行等待的性能损耗
  • 用缓存机制临时存储Server2返回的Token流,直到Server1的结果明确
  • 根据Server1的响应分支处理:
    • 若返回有效整数列表:先输出缓存的Token,再继续推送后续实时Token,最后追加整数列表结果
    • 若返回空/无效列表:立即终止Server2的流,返回无可用信息的提示

关键Flux操作符说明

  • cache():临时缓存Server2的流数据,避免重复发起调用,同时支持后续按需输出缓存内容
  • Mono.zip():并行订阅多个异步任务,确保Server1和Server2同时启动执行
  • concatWith():按顺序拼接Token流和最终列表结果,保证返回顺序符合需求
  • dispose():主动取消Server2的订阅,避免无效流持续占用资源

完整代码实现

1. 模拟远程服务调用(Service层)

import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import org.springframework.stereotype.Service;
import java.time.Duration;
import java.util.List;

@Service
public class TokenService {

    // 模拟调用Server1:返回整数列表,模拟1秒延迟
    public Mono<List<Integer>> callServer1(String query) {
        return Mono.delay(Duration.ofSeconds(1))
                .map(tick -> {
                    // 模拟业务逻辑:根据query返回有效列表或空列表
                    return "valid".equals(query) ? List.of(1, 2, 3) : List.of();
                })
                .timeout(Duration.ofSeconds(5))
                .onErrorResume(e -> Mono.just(List.of())); // 异常时返回空列表
    }

    // 模拟调用Server2:返回Token流,每500ms生成一个Token,模拟持续输出
    public Flux<String> callServer2(String query, List<Integer> tempList) {
        return Flux.interval(Duration.ofMillis(500))
                .map(tick -> {
                    // 模拟先输出临时Token,后续结合列表输出正式Token
                    return tick < 3 ? "Temp-Token-" + tick : "Formal-Token-" + tick + "-" + tempList;
                })
                .take(10) // 模拟最多生成10个Token
                .timeout(Duration.ofSeconds(10))
                .onErrorResume(e -> Flux.empty());
    }
}

2. Controller层并行逻辑实现

import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;

@RestController
public class TokenController {

    private final TokenService tokenService;

    public TokenController(TokenService tokenService) {
        this.tokenService = tokenService;
    }

    @GetMapping("/tokens")
    public Flux<String> getTokens(@RequestParam String query) {
        // 1. 启动Server2的流并缓存所有输出
        Flux<String> server2TokenFlux = tokenService.callServer2(query, List.of())
                .cache(Integer.MAX_VALUE);

        // 2. 并行执行Server1和Server2,等待Server1的结果做分支处理
        return Mono.zip(tokenService.callServer1(query), Mono.just(server2TokenFlux))
                .flatMapMany(tuple -> {
                    List<Integer> server1Result = tuple.getT1();
                    Flux<String> tokenStream = tuple.getT2();

                    if (server1Result != null && !server1Result.isEmpty()) {
                        // 有效列表:输出缓存Token + 后续实时Token + 最终列表信息
                        return tokenStream.concatWith(Mono.just("Final-Integer-List: " + server1Result));
                    } else {
                        // 无效列表:主动取消Server2订阅,返回提示
                        tokenStream.subscribe().dispose();
                        return Mono.just("No valid information available");
                    }
                });
    }
}

关键注意事项

  1. 缓存内存控制:若Token流数据量较大,建议给cache()设置固定缓存大小(如cache(100)),避免内存溢出
  2. 订阅取消机制:当Server1返回无效结果时,必须主动调用dispose()取消Server2的订阅,防止无效流持续占用资源
  3. 超时与异常处理:为所有远程调用添加超时和异常兜底逻辑,保证服务的稳定性和容错性
  4. 参数传递优化:若Server2必须依赖Server1的真实列表才能生成完整Token,可调整为Server2先输出基础Token,待Server1结果返回后再补传列表生成后续Token(需结合业务场景调整)

内容的提问来源于stack exchange,提问作者James K J

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 11:21:27