一、副作用与监听(doOnXxx 系列)
1. doOnNext —— 每个元素到达时触发
Flux<String> stream = chatClient.prompt("讲个笑话").stream().content();
stream
.doOnNext(chunk -> System.out.println("收到片段: " + chunk))
.subscribe();
2. doOnComplete —— 流正常完成时触发
stream
.doOnComplete(() -> System.out.println("✅ 笑话讲完了!"))
.subscribe();
3. doOnError —— 流发生错误时触发
stream
.doOnError(e -> System.err.println("❌ 出错了: " + e.getMessage()))
.subscribe();
4. doFinally —— 无论何种结局,最后都会执行
stream
.doFinally(signal -> System.out.println("最终信号: " + signal)) // onComplete / onError / cancel
.subscribe();
5. doOnCancel —— 被取消时触发
Disposable disposable = stream
.doOnCancel(() -> System.out.println("用户取消了"))
.subscribe();
// 模拟取消
disposable.dispose();
6. doFirst —— 订阅之前执行(最早)
stream
.doFirst(() -> System.out.println("1. 最先执行"))
.doOnSubscribe(sub -> System.out.println("2. 然后订阅"))
.doOnNext(chunk -> System.out.println("3. 然后收数据"))
.subscribe();
7. doOnSubscribe —— 订阅发生时触发
stream
.doOnSubscribe(sub -> System.out.println("订阅已建立,开始请求"))
.subscribe();
8. doOnRequest —— 每次请求数据时触发
stream
.doOnRequest(n -> System.out.println("下游请求了 " + n + " 个元素"))
.subscribe();
9. doOnTerminate —— 终止时触发(完成或错误,不含取消)
stream
.doOnTerminate(() -> System.out.println("流终止了(完成或出错)"))
.subscribe();
10. doAfterTerminate —— 终止之后触发
stream
.doAfterTerminate(() -> System.out.println("终止信号已传播给下游后执行"))
.subscribe();
二、转换与过滤
11. map —— 一对一同步转换
stream
.map(chunk -> chunk + " 🐋") // 每个字后面加鲸鱼
.map(String::toUpperCase) // 转大写
.subscribe();
12. mapNotNull —— 转换,但 null 值会被过滤掉
stream
.mapNotNull(chunk -> {
if (chunk.equals("敏感词")) return null; // 直接丢弃
return "【" + chunk + "】";
})
.subscribe();
13. filter —— 过滤不符合条件的元素
stream
.filter(chunk -> chunk != null && !chunk.trim().isEmpty())
.filter(chunk -> !chunk.equals("[DONE]"))
.subscribe();
14. cast —— 类型强制转换
Flux<Object> objectStream = chatClient.prompt("hi").stream().content().map(Object.class::cast);
objectStream
.cast(String.class) // 转回 String
.subscribe();
15. ofType —— 只保留指定类型
Flux<Object> mixed = Flux.just("hello", 123, "world", 456);
mixed
.ofType(String.class) // 只保留 String
.subscribe(s -> System.out.println(s)); // hello, world
16. flatMap —— 每个元素转成 Publisher,然后合并(可交叉)
stream
.flatMap(chunk -> {
// 每个片段去查增强信息(异步)
return enrichService.enrich(chunk)
.onErrorReturn("【查不到】");
})
.subscribe();
17. concatMap —— 同 flatMap,但保证顺序(一个接一个)
stream
.concatMap(chunk -> {
// 按顺序处理,前一个完成才处理下一个
return processService.process(chunk);
})
.subscribe();
18. flatMapSequential —— 并发订阅,但按原顺序输出
stream
.flatMapSequential(chunk -> enrichService.enrich(chunk), 4) // 并发数4
.subscribe();
19. flatMapIterable —— 每个元素转成 Iterable 然后扁平化
stream
.flatMapIterable(chunk -> List.of(chunk.split(""))) // 把每个词拆成单字
.subscribe();
20. concatMapIterable —— 同 flatMapIterable,顺序保证
stream
.concatMapIterable(chunk -> List.of(chunk, chunk + "!"))
.subscribe();
21. handle —— 灵活处理,可同时完成过滤+转换+终止
stream
.handle((chunk, sink) -> {
if (chunk.contains("秘密")) {
sink.error(new RuntimeException("发现秘密!"));
} else if (chunk.trim().isEmpty()) {
// 跳过,不调用 sink.next
} else {
sink.next("【" + chunk + "】");
}
})
.subscribe();
三、错误处理
22. onErrorResume —— 发生错误时切换备用流
stream
.onErrorResume(e -> {
System.err.println("备用方案启动");
return Flux.just("(服务器正在休息,请稍后再来)");
})
.subscribe();
23. onErrorReturn —— 发生错误时返回一个固定值
stream
.onErrorReturn("(出错了,这是默认回复)")
.subscribe();
24. onErrorMap —— 将错误转换成另一种错误
stream
.onErrorMap(e -> new BusinessException("AI 服务异常: " + e.getMessage(), e))
.subscribe();
25. onErrorComplete —— 发生错误时直接完成(不传播错误)
stream
.onErrorComplete() // 有错就悄悄结束,当作什么都没发生
.doOnComplete(() -> System.out.println("结束了(可能是有错但被吞了)"))
.subscribe();
26. retry —— 发生错误时重试(固定次数)
stream
.retry(3) // 最多重试3次
.subscribe();
27. retryWhen —— 灵活重试策略
stream
.retryWhen(Retry.backoff(3, Duration.ofSeconds(1))
.doBeforeRetry(rs -> System.out.println("第" + (rs.totalRetries()+1) + "次重试")))
.subscribe();
28. timeout —— 设置超时
stream
.timeout(Duration.ofSeconds(5))
.onErrorReturn(TimeoutException.class, "(超时了,AI 想太久啦)")
.subscribe();
29. onErrorContinue —— 出错的元素被跳过,继续后面的
stream
.onErrorContinue((e, val) -> {
System.err.println("元素 '" + val + "' 处理失败,跳过: " + e);
})
.map(chunk -> {
if (chunk.contains("坏词")) throw new RuntimeException("bad");
return chunk;
})
.subscribe();
30. onErrorStop —— 取消上游的 onErrorContinue 策略
stream
.onErrorContinue((e, v) -> {}) // 上游启用跳过
.map(chunk -> {
if (chunk.contains("危险")) throw new RuntimeException("危险");
return chunk;
})
.onErrorStop() // 这里开始,错误不再被跳过,会正常传播
.subscribe();
四、流控制与截取
31. take —— 取前 N 个
stream
.take(10) // 只要前10个片段
.subscribe();
32. takeLast —— 取最后 N 个(流结束后才发出)
stream
.takeLast(5) // 等流结束,只发最后5个片段
.subscribe();
33. takeUntil —— 取到某个条件满足为止(包含触发元素)
stream
.takeUntil(chunk -> chunk.contains("。"))) // 取到第一个句号为止
.subscribe();
34. takeWhile —— 条件满足时一直取(不包含第一个不满足的)
stream
.takeWhile(chunk -> !chunk.contains("【结束】")) // 遇到【结束】就停,不包含它
.subscribe();
35. takeUntilOther —— 取到另一个 Publisher 发出信号为止
Flux<String> stopSignal = Flux.never().delaySubscription(Duration.ofSeconds(3));
stream
.takeUntilOther(stopSignal) // 3秒后停止
.subscribe();
36. skip —— 跳过前 N 个
stream
.skip(5) // 忽略前5个片段
.subscribe();
37. skipLast —— 跳过后 N 个
stream
.skipLast(3) // 忽略最后3个片段
.subscribe();
38. skipUntil —— 跳过直到条件满足(不包含触发元素)
stream
.skipUntil(chunk -> chunk.contains("正文")) // 跳过"正文"之前的所有
.subscribe();
39. skipWhile —— 条件满足时跳过
stream
.skipWhile(chunk -> chunk.length() < 2) // 跳过单字,直到遇到双字及以上
.subscribe();
40. skipUntilOther —— 跳过直到另一个 Publisher 发出信号
Flux<String> startSignal = Flux.just("开始").delaySubscription(Duration.ofSeconds(2));
stream
.skipUntilOther(startSignal) // 前2秒的内容全部跳过
.subscribe();
41. limitRate —— 控制请求速率(背压优化)
stream
.limitRate(10) // 每次最多请求10个,75%时补充
// 或者定制高低水位
.limitRate(20, 5) // 高水位20,低水位5
.subscribe();
42. sample —— 定期采样最新值
stream
.sample(Duration.ofMillis(200)) // 每200ms取最新一个片段
.subscribe();
43. sampleFirst —— 取每个窗口的第一个值
stream
.sampleFirst(Duration.ofMillis(500)) // 每500ms取第一个,然后跳过这500ms内的其他
.subscribe();
44. sampleTimeout —— 每个元素触发一个窗口,窗口内无新值则发出
stream
.sampleTimeout(chunk -> Mono.delay(Duration.ofMillis(300))) // 每个片段后等300ms,若期间有新片段则替换
.subscribe();
五、分组与聚合
45. window —— 拆分成多个子流(你问的)
stream
.window(5) // 每5个片段一组,每个组是一个新的 Flux<String>
.flatMap(window -> window.collectList()) // 把每组收集成 List
.subscribe(list -> System.out.println("一组: " + list));
46. windowUntil —— 按条件拆分子流(包含触发元素)
stream
.windowUntil(chunk -> chunk.contains("。")) // 遇到句号就关窗,句号在旧窗里
.flatMap(window -> window.collectList())
.subscribe();
47. windowWhile —— 条件满足时保持窗口,不满足时关闭(不包含触发元素)
stream
.windowWhile(chunk -> !chunk.contains("。")) // 没遇到句号就一直在一个窗里
.flatMap(window -> window.collectList())
.subscribe();
48. windowTimeout —— 按数量或时间拆分
stream
.windowTimeout(10, Duration.ofSeconds(1)) // 每10个或每1秒关窗
.flatMap(window -> window.collectList())
.subscribe();
49. windowWhen —— 由外部 Publisher 控制窗口开关
Flux<Long> opens = Flux.interval(Duration.ofSeconds(2));
stream
.windowWhen(opens, open -> Mono.delay(Duration.ofSeconds(1))) // 每2秒开窗,1秒后关窗
.flatMap(window -> window.collectList())
.subscribe();
50. windowUntilChanged —— 元素变化时才关窗
Flux.just("A","A","B","B","A")
.windowUntilChanged()
.flatMap(w -> w.collectList())
.subscribe(list -> System.out.println(list)); // [A,A] [B,B] [A]
51. buffer —— 同 window,但输出 List
stream
.buffer(5) // 每5个片段一个 List
.subscribe(list -> System.out.println("一批: " + list));
52. bufferUntil —— 同 windowUntil,但输出 List
stream
.bufferUntil(chunk -> chunk.contains("。"))
.subscribe(list -> System.out.println(list));
53. bufferWhile —— 同 windowWhile,但输出 List
stream
.bufferWhile(chunk -> !chunk.contains("。"))
.subscribe();
54. bufferTimeout —— 按数量或时间收集 List
stream
.bufferTimeout(10, Duration.ofSeconds(1))
.subscribe(list -> System.out.println("批量: " + list));
55. bufferWhen —— 由外部 Publisher 控制
Flux<Long> opens = Flux.interval(Duration.ofSeconds(2));
stream
.bufferWhen(opens, open -> Mono.delay(Duration.ofSeconds(1)))
.subscribe();
56. collectList —— 收集全部到一个 List(流结束时发出)
stream
.collectList()
.subscribe(list -> System.out.println("完整内容: " + String.join("", list)));
57. collectMap —— 收集成 Map
stream
.collectMap(chunk -> chunk.length(), chunk -> chunk) // key=长度,value=片段
.subscribe(map -> System.out.println(map));
58. collectMultimap —— 收集成 Map<K, Collection<V>>
stream
.collectMultimap(chunk -> chunk.length()) // 相同长度的放一起
.subscribe(map -> System.out.println(map));
59. collect —— 自定义收集器
stream
.collect(
StringBuilder::new,
(sb, chunk) -> sb.append(chunk)
)
.subscribe(sb -> System.out.println("完整: " + sb.toString()));
60. reduce —— 累积归并(同类型)
stream
.reduce((a, b) -> a + b) // 拼接所有片段
.subscribe(optional -> optional.ifPresent(System.out::println));
61. reduce(带初始值)—— 累积归并(可不同类型)
stream
.reduce(0, (count, chunk) -> count + chunk.length()) // 统计总字数
.subscribe(total -> System.out.println("总字数: " + total));
62. reduceWith —— 延迟提供初始值
stream
.reduceWith(() -> 0, (count, chunk) -> count + chunk.length())
.subscribe();
63. scan —— 同 reduce,但每次中间结果都发出
stream
.scan((acc, chunk) -> acc + chunk) // 边拼边发,每次都是"已拼的完整前缀"
.subscribe(partial -> System.out.println("目前: " + partial));
64. scan(带初始值)
stream
.scan("【开始】", (acc, chunk) -> acc + chunk)
.subscribe();
65. scanWith —— 延迟提供初始值
stream
.scanWith(() -> "", (acc, chunk) -> acc + chunk)
.subscribe();
六、组合与合并
66. concatWith —— 拼接另一个流在后面
Flux<String> extra = Flux.just("【完】", "🐋");
stream
.concatWith(extra) // AI回答后面追加结束语
.subscribe();
67. mergeWith —— 合并另一个流(交叉发出)
Flux<String> other = Flux.just("A", "B", "C").delayElements(Duration.ofMillis(100));
Flux<String> merged = stream.mergeWith(other);
merged.subscribe();
68. zipWith —— 一对一配对另一个流
Flux<String> emotions = Flux.just("😊", "😂", "😅");
stream
.zipWith(emotions) // (片段, 表情) 配对
.subscribe(tuple -> System.out.println(tuple.getT1() + tuple.getT2()));
69. zipWithIterable —— 与 Iterable 配对
List<String> tags = List.of("tag1", "tag2", "tag3");
stream
.zipWithIterable(tags)
.subscribe(tuple -> System.out.println(tuple));
70. startWith —— 在前面加元素
stream
.startWith("【开始】", "🐋") // 先发这两个,再发 AI 内容
.subscribe();
71. startWith(Iterable)
List<String> prefix = List.of("【注意】", "以下内容纯属虚构");
stream
.startWith(prefix)
.subscribe();
七、条件与分支
72. defaultIfEmpty —— 如果空流则发默认值
Flux<String> empty = Flux.empty();
empty
.defaultIfEmpty("(AI 什么都没说)")
.subscribe();
73. switchIfEmpty —— 如果空流则切换到另一个流
Flux<String> fallback = Flux.just("(AI 卡住了,这是备用回答)");
stream
.filter(chunk -> false) // 故意让所有元素都被过滤掉,变成空流
.switchIfEmpty(fallback)
.subscribe();
八、其他实用方法
74. distinct —— 去重(全局)
Flux.just("A","B","A","C","B")
.distinct()
.subscribe(); // A, B, C
75. distinctUntilChanged —— 去重(相邻)
Flux.just("A","A","B","B","A")
.distinctUntilChanged()
.subscribe(); // A, B, A
76. count —— 计数
stream
.count()
.subscribe(cnt -> System.out.println("总片段数: " + cnt));
77. hasElements —— 是否有元素
stream
.hasElements()
.subscribe(has -> System.out.println("有内容吗? " + has));
78. hasElement —— 是否包含某个值
stream
.hasElement("🐋")
.subscribe();
79. all —— 所有元素都满足条件?
stream
.all(chunk -> chunk.length() > 0)
.subscribe();
80. any —— 有任何一个元素满足条件?
stream
.any(chunk -> chunk.contains("鲸鱼"))
.subscribe();
81. elementAt —— 取第 N 个元素
stream
.elementAt(5) // 取第6个片段(0-based)
.subscribe();
82. elementAt(带默认值)
stream
.elementAt(100, "(没有第100个)")
.subscribe();
83. index —— 给每个元素加索引
stream
.index() // 变成 Tuple2<Long, String>
.subscribe(tuple -> System.out.println(tuple.getT1() + ": " + tuple.getT2()));
84. timestamp —— 给每个元素加时间戳
stream
.timestamp() // Tuple2<Long, String>
.subscribe(tuple -> System.out.println(tuple.getT1() + "ms: " + tuple.getT2()));
85. elapsed —— 每个元素与前一个元素的时间差
stream
.elapsed() // Tuple2<Long, String>
.subscribe(tuple -> System.out.println("间隔 " + tuple.getT1() + "ms: " + tuple.getT2()));
86. timed —— 更丰富的时间信息
stream
.timed() // Flux<Timed<String>>
.subscribe(timed -> {
System.out.println("值: " + timed.get());
System.out.println("耗时: " + timed.elapsed().toMillis() + "ms");
System.out.println("时间戳: " + timed.timestamp());
});
87. repeat —— 重复订阅
Flux.just("Hello")
.repeat(3) // 总共发出4次(原1次+重播3次)
.subscribe();
88. repeatWhen —— 由条件控制重复
Flux.just("Hello")
.repeatWhen(flux -> flux.take(3)) // 重复3次
.subscribe();
89. cache —— 缓存结果供后续订阅
Flux<String> cached = stream.cache(); // 第一次订阅会真正执行,后续订阅直接发缓存数据
cached.subscribe(System.out::println);
cached.subscribe(System.out::println); // 第二次立即拿到,不会重新请求
90. cache(带大小)
stream.cache(10); // 只缓存最近10个
91. cache(带 TTL)
stream.cache(Duration.ofMinutes(1)); // 缓存1分钟
92. share —— 多订阅共享同一个流(不重复执行)
Flux<String> shared = stream.share();
shared.subscribe(s -> System.out.println("A: " + s));
shared.subscribe(s -> System.out.println("B: " + s)); // A和B收到同样的数据
93. publish —— 更灵活的共享控制
ConnectableFlux<String> connectable = stream.publish();
connectable.subscribe(s -> System.out.println("A: " + s));
connectable.subscribe(s -> System.out.println("B: " + s));
connectable.connect(); // 手动启动
94. replay —— 历史重播
ConnectableFlux<String> replayable = stream.replay(5); // 缓存5个,新订阅也能收到历史
replayable.subscribe(s -> System.out.println("A: " + s));
// 等一会儿...
replayable.subscribe(s -> System.out.println("B: " + s)); // B能收到之前缓存的5个
replayable.connect();
95. expand / expandDeep —— 递归展开
Flux.just(1)
.expand(v -> Flux.just(v + 1).take(3)) // 广度优先:1,2,3,4
.subscribe();
96. groupBy —— 按键分组
Flux.just("cat","dog","cow","rat")
.groupBy(s -> s.length()) // 按长度分组
.flatMap(group -> group.collectList())
.subscribe();
97. log —— 打印日志(调试神器)
stream
.log("AI-Stream") // 打印所有信号
.subscribe();
98. checkpoint —— 在错误栈中标记位置
stream
.map(chunk -> chunk.toUpperCase())
.checkpoint("🚀 转大写之后") // 如果出错,栈里会有这个标记
.subscribe();
99. publishOn —— 切换后续操作执行的线程
stream
.publishOn(Schedulers.boundedElastic()) // 后面的操作在弹性线程池执行
.map(chunk -> heavyProcessing(chunk))
.subscribe();
100. subscribeOn —— 切换订阅和源头执行的线程
stream
.subscribeOn(Schedulers.boundedElastic()) // 整个链从源头就在弹性线程执行
.subscribe();
101. contextWrite —— 写入上下文(传递数据)
stream
.contextWrite(ctx -> ctx.put("userId", "123"))
.map(chunk -> {
String userId = (String) Context.currentView().get("userId");
return "[" + userId + "] " + chunk;
})
.subscribe();
102. transform —— 组装时转换(一次性)
Function<Flux<String>, Flux<String>> addPrefix = flux -> flux.map(s -> ">> " + s);
stream
.transform(addPrefix)
.subscribe();
103. transformDeferred —— 订阅时转换(每次订阅独立)
stream
.transformDeferred(flux -> flux.map(s -> "⏰ " + s))
.subscribe();
104. transformDeferredContextual —— 订阅时转换,可访问上下文
stream
.transformDeferredContextual((flux, ctx) -> {
String user = ctx.getOrDefault("user", "匿名");
return flux.map(s -> "[" + user + "] " + s);
})
.contextWrite(ctx -> ctx.put("user", "鲸鱼娘"))
.subscribe();
105. as —— 转成任意类型(工具方法)
List<String> result = stream
.as(flux -> flux.collectList().block()); // 转成 List
System.out.println(result);
106. hide —— 隐藏实际类型,防止优化
stream
.hide() // 外部看到的是 Flux,不是具体的实现类
.subscribe();