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。