Flux<String> 的使用

Flux<String> 的使用

lx 1 2026-07-24

Flux 的使用

📚 Flux 核心操作符全景图

数据源 → filter → map → flatMap → doOnNext → 
        ↓
    buffer → collectList → reduce → 
        ↓
    timeout → retry → onErrorResume → 
        ↓
    concatWith → mergeWith → take → 
        ↓
    订阅执行

一、数据转换类(3个)

1️⃣ map() - 一对一转换

作用: 对每个元素进行转换,保持数量不变。

// 场景:Ollama返回的chunk包含完整JSON,提取content字段
Flux<String> ollamaStream = Flux.just(
    "{\"response\":\"你好\",\"done\":false}",
    "{\"response\":\"杭州\",\"done\":false}",
    "{\"response\":\"欢迎你\",\"done\":true}"
);

Flux<String> contentStream = ollamaStream
    .map(chunk -> {
        JsonNode node = MAPPER.readTree(chunk);
        return node.get("response").asText();
    })
    .map(String::trim); // 去除空格

// 输出:你好 → 杭州 → 欢迎你
contentStream.subscribe(System.out::print); // "你好杭州欢迎你"

2️⃣ flatMap() - 一对多转换

作用: 每个元素转换为一个 Flux,然后合并所有 Flux。异步、非阻塞

// 场景:流式翻译 - 每个中文词翻译成英文
Flux<String> chineseStream = Flux.just("你好", "杭州", "欢迎");

Flux<String> translatedStream = chineseStream
    .flatMap(chinese -> {
        // 每个词调用翻译API(异步),返回 Flux<String>
        return translateService.translate(chinese); 
        // translate() 返回 Flux<String>,可能产生多个结果
    });

// 注意:顺序不保证!可能输出 "Hello" "Hangzhou" "Welcome"
// 也可能交错输出

3️⃣ flatMapSequential() - 有序的一对多

// 场景:需要保持顺序的场景
Flux<String> translatedStream = chineseStream
    .flatMapSequential(chinese -> translateService.translate(chinese));
// 保证顺序:你好 → 杭州 → 欢迎

4️⃣ cast() - 类型转换

// 场景:从 Object 转为 String
Flux<Object> objectStream = Flux.just("hello", "world");
Flux<String> stringStream = objectStream.cast(String.class);

二、过滤类(3个)

5️⃣ filter() - 按条件过滤

作用: 保留满足条件的元素,移除不满足的。

// 场景:过滤掉空内容和特殊字符
Flux<String> rawStream = Flux.just("你好", "", "   ", null, "杭州", "\n");

Flux<String> cleanStream = rawStream
    .filter(content -> content != null)           // 移除 null
    .filter(content -> !content.trim().isEmpty()) // 移除空串
    .filter(content -> !content.contains("\n"));  // 移除换行

// 输出:你好 → 杭州

6️⃣ take() - 限制数量

// 场景:只取前100个token,防止输出过长
Flux<String> longStream = ollamaStream.map(this::parseContent);

Flux<String> limitedStream = longStream
    .take(100); // 只取前100个chunk

// 场景:前10秒内的数据
Flux<String> timeLimited = longStream
    .take(Duration.ofSeconds(10)); // 只流10秒

7️⃣ distinct() - 去重

// 场景:避免重复输出相同内容
Flux<String> streamWithDuplicates = Flux.just("A", "A", "B", "A", "C");
Flux<String> uniqueStream = streamWithDuplicates.distinct();
// 输出:A → B → C(首次出现)

三、组合类(3个)

8️⃣ concatWith() - 串联

作用: 先完成当前 Flux,再接下一个 Flux。顺序执行

// 场景:输出正文后,追加"感谢阅读"
Flux<String> contentStream = Flux.just("你好", "杭州");
Flux<String> thankStream = Flux.just("感谢阅读", "再见");

Flux<String> finalStream = contentStream
    .concatWith(thankStream); // 先输出"你好杭州",再输出"感谢阅读再见"

// 输出:你好 → 杭州 → 感谢阅读 → 再见

9️⃣ mergeWith() - 合并

作用: 多个 Flux 合并,谁先有数据谁先输出

// 场景:同时从多个AI模型获取回复
Flux<String> ollamaStream = ollamaService.generate(prompt);
Flux<String> gptStream = gptService.generate(prompt);

Flux<String> mergedStream = ollamaStream
    .mergeWith(gptStream); // 谁先返回就输出谁

// 输出可能交错:Ollama的chunk1 → GPT的chunk1 → Ollama的chunk2 ...

🔟 zip() - 合并配对

作用: 将两个 Flux 按位置配对。

// 场景:中英对照输出
Flux<String> chineseStream = Flux.just("你好", "杭州");
Flux<String> englishStream = Flux.just("Hello", "Hangzhou");

Flux<String> pairedStream = Flux.zip(chineseStream, englishStream)
    .map(tuple -> tuple.getT1() + " (" + tuple.getT2() + ")");
// 输出:你好 (Hello) → 杭州 (Hangzhou)

四、聚合类(3个)

1️⃣1️⃣ collectList() - 收集为列表

作用: 等待所有数据完成后,收集为 List。

// 场景:完整对话完成后做分析
Flux<String> stream = Flux.just("今天", "天气", "不错");

Mono<List<String>> listMono = stream.collectList();
// 或阻塞方式(测试代码)
List<String> list = stream.collectList().block();

// 使用场景:流结束后做统计分析
List<String> result = stream.collectList().block();
int totalChars = result.stream().mapToInt(String::length).sum();
System.out.println("总字符数: " + totalChars);

1️⃣2️⃣ reduce() - 累积计算

作用: 将数据流累积为单个结果(类似 SQL 的 SUM)。

// 场景:拼接完整内容
Flux<String> stream = Flux.just("你好", "杭州", "欢迎你");

Mono<String> fullTextMono = stream
    .reduce((acc, next) -> acc + next);
// 或指定初始值
Mono<String> fullText = stream
    .reduce("最终结果:", (acc, next) -> acc + next);

// 输出:你好杭州欢迎你
String result = fullTextMono.block();

五、副作用类(2个)

1️⃣3️⃣ doOnNext() - 观察每个元素

作用: 对每个元素执行副作用操作(日志、打印、统计),不影响数据流

// 场景:实时打印 + 统计
AtomicInteger count = new AtomicInteger(0);
Flux<String> stream = Flux.just("你好", "杭州", "欢迎")
    .doOnNext(chunk -> {
        System.out.print(chunk); // 实时打印到控制台
        count.incrementAndGet(); // 统计chunk数量
    })
    .doOnComplete(() -> {
        System.out.println("\n总共 " + count.get() + " 个片段");
    });

// 输出:你好杭州欢迎
// 完成后:总共 3 个片段

1️⃣4️⃣ doOnError() - 错误监听

// 场景:错误日志记录
Flux<String> stream = riskyStream()
    .doOnError(e -> log.error("流式处理异常", e))
    .doOnError(e -> System.err.println("错误: " + e.getMessage()));

六、错误处理类(3个)

1️⃣5️⃣ onErrorResume() - 错误恢复

作用: 发生错误时,切换到备用的 Flux。

// 场景:Ollama不可用时返回默认回复
Flux<String> safeStream = ollamaStream
    .onErrorResume(e -> {
        log.warn("Ollama失败,使用备用回复", e);
        return Flux.just("系统暂时繁忙", "请稍后重试");
    });

// 场景:降级到备用模型
Flux<String> fallbackStream = ollamaStream
    .onErrorResume(e -> {
        log.info("切换到备用模型GPT");
        return gptService.generate(prompt);
    });

1️⃣6️⃣ onErrorReturn() - 返回默认值

// 场景:出错时返回固定值
Flux<String> safeStream = ollamaStream
    .onErrorReturn("抱歉,我暂时无法回答");

1️⃣7️⃣ retry() - 重试

作用: 出错后自动重试。

// 场景:网络不稳定时自动重试3次
Flux<String> robustStream = ollamaStream
    .retry(3) // 失败后重试3次
    .retryWhen(error -> error.delayElements(Duration.ofSeconds(1))); // 延迟1秒重试

// 场景:指数退避重试
Flux<String> smartRetry = ollamaStream
    .retryWhen(error -> error
        .zipWith(Flux.range(1, 3)) // 最多重试3次
        .flatMap(tuple -> {
            int retryCount = tuple.getT2();
            long delay = (long) Math.pow(2, retryCount) * 1000; // 2s, 4s, 8s
            return Mono.delay(Duration.ofMillis(delay));
        })
    );

七、时间控制类(2个)

1️⃣8️⃣ timeout() - 超时控制

// 场景:AI响应超时处理
Flux<String> stream = ollamaStream
    .timeout(Duration.ofSeconds(30)) // 30秒超时
    .onErrorResume(TimeoutException.class, e -> 
        Flux.just("请求超时,请稍后重试")
    );

1️⃣9️⃣ delayElements() - 延迟发送

// 场景:模拟打字机效果(每100ms输出一个字)
Flux<String> stream = Flux.just("你", "好", "杭", "州")
    .delayElements(Duration.ofMillis(100));
// 每100ms输出一个字

八、实战:完整的 AI 流式对话

@Service
public class AiChatService {
  
    public Flux<String> chatStream(String prompt) {
        // 1. 调用Ollama获取流式数据
        Flux<String> rawStream = webClient.post()
            .uri("/api/generate")
            .bodyValue(Map.of("model", "qwen3.5:4b", "prompt", prompt))
            .retrieve()
            .bodyToFlux(String.class);
    
        // 2. 使用操作符链处理
        return rawStream
            // 解析每个chunk
            .map(this::parseContent)
            // 过滤空内容
            .filter(content -> content != null && !content.isEmpty())
            // 过滤掉标点符号(可选)
            .filter(content -> !content.matches("^[,。、!?]$"))
            // 实时打印日志
            .doOnNext(chunk -> log.debug("输出: {}", chunk))
            // 遇到错误时重试
            .retry(2)
            // 30秒超时
            .timeout(Duration.ofSeconds(30))
            // 错误降级
            .onErrorResume(e -> {
                log.error("流式处理失败", e);
                return Flux.just("抱歉,系统繁忙");
            })
            // 流结束后追加提示
            .concatWith(Flux.just("\n\n—— 回答结束 ——"));
    }
  
    private String parseContent(String raw) {
        try {
            JsonNode node = MAPPER.readTree(raw);
            return node.path("response").asText(null);
        } catch (Exception e) {
            return null;
        }
    }
}

🎯 使用频率统计(实际项目)

操作符 使用频率 场景
map() ⭐⭐⭐⭐⭐ 数据转换(必用)
filter() ⭐⭐⭐⭐⭐ 数据过滤(必用)
doOnNext() ⭐⭐⭐⭐ 日志/统计
onErrorResume() ⭐⭐⭐⭐ 错误处理
retry() ⭐⭐⭐ 网络重试
timeout() ⭐⭐⭐ 超时控制
concatWith() ⭐⭐⭐ 追加额外信息
collectList() ⭐⭐ 测试/批量处理
flatMap() ⭐⭐ 并发调用
distinct() 去重(AI场景少)

📝 记忆口诀

转换用 map,过滤用 filter,
合并看 concat,扁平用 flat,
错误找 resume,重试有 retry,
超时 timeout,延迟 delay,
观察 doOn,收集 collect,
完成 reduce,取用 take。