Loading...

SpringAI(Flux)方法示例大全

一、副作用与监听(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();


0

回到顶部