Reactor 核心 API

共 18 题
📑 题目列表 18 题
#
★★★

1. Flux.fromIterable/fromStream/fromFuture 的工程应用

Flux.fromIterable、Flux.fromStream、Flux.fromFuture 这三个工厂方法分别适用于什么工程场景?使用时有什么注意事项?

  • 三种工厂方法的数据来源与适用场景
  • fromStream 的流不可复用问题
  • fromFuture 与异步任务(CompletableFuture)的桥接

Flux.fromIterable 用于把 Collection/Iterable 集合转换为 Flux,适合批量数据源,数据在订阅时逐条发射,且 Iterable 可被多次订阅复用。Flux.fromStream 用于把 Java Stream 转换为 Flux,但 Stream 只能被消费一次,多次订阅会抛异常,因此通常配合 defer 使用以避免复用问题。Flux.fromFuture 用于把 CompletableFuture 等异步任务的结果桥接为响应式流,返回只发射单个结果的 Flux(单值流;如需 Mono 可用 Mono.fromFuture),适合把基于回调的异步库接入 Reactor 链。工程上:fromIterable 用于内存集合;fromStream 用于已存在的 Stream 数据源(注意一次性);fromFuture 用于与 CompletableFuture 类异步 API 集成。

三者都是"把非响应式数据源拉入响应式世界"的桥梁。区别在于数据源的特性:Iterable 可复用,Stream 一次性,Future 是异步单值。理解这些特性是避免订阅错误的关键。

#
★★★

2. Flux.window/buffer/groupBy 的批处理

Flux.window、Flux.buffer、Flux.groupBy 这三个操作符在批处理上有什么区别?各自适合什么场景?

  • buffer 按数量/时间聚合为 List
  • window 按数量/时间切分为子 Flux
  • groupBy 按键分组为多组子流

buffer 将上游元素按数量(buffer(n))或时间(bufferTimeout)聚合为一个个 List 发射,适合批量处理(如批量写库、批量发送)。window 则把上游元素切分为多个子 Flux(window(n) 每个子 Flux 最多 n 个元素),适合对子流继续做流式操作(如对每个窗口做聚合、与下游组合)。groupBy 根据 key 函数把元素分组,每组产生一个子 Flux,适合按类别或按 key 维度并发处理。三者的选用:需要"定长批次"用 buffer;需要"对流做进一步流式处理"用 window;需要"按 key 分流"用 groupBy。注意 buffer 可能因大列表占用内存,window/groupBy 需要及时消费子流避免背压堆积。

三者都是"聚合/切分"操作符,但产物不同:buffer 产出 List,window 产出子 Flux,groupBy 按 key 产出动态子流。选择依据是下游是以"集合"还是"子流"为单位处理。

#
★★★

3. CombineLatest/merge/concat/zip 的差异

Flux.combineLatest、merge、concat、zip 这几个组合操作符有什么差异?各自适合什么场景?

  • merge 交错合并(并发订阅,按到达顺序)
  • zip 按位置配对(等待对应元素)
  • combineLatest 取最新值组合

这四个操作符用于组合多个 Flux。concat 按顺序拼接:先完成第一个流再订阅第二个,强调顺序性,适合"必须前一阶段完成后才进行下一阶段"。merge 并发订阅所有源,只要有元素就立即发射(交错),不保证顺序,适合"多个独立事件源合并"。zip 按位置/索引配对,会等待每个源的对齐元素同时到达才发射一对,适合"多个源必须一一对应"(如把两个序列对齐)。combineLatest 在任何源发射新值时,取所有源的最新值合并,适合"只要任一源更新就基于最新状态计算"(如表单联动)。选用时:需要严格顺序用 concat,需要最大化吞吐交错用 merge,需要对齐配对用 zip,需要并集最新状态用 combineLatest。

差异核心是"时序与配对的规则"。concat 关心顺序,merge 关心并发交错,zip 关心位置对齐,combineLatest 关心最新状态。理解各自动态行为(如 zip 是否等所有源)是正确选用的关键。

#
★★

4. ConnectableFlux 与 share/refCount 的取舍

ConnectableFlux 与 share()/refCount() 有什么区别?使用时如何取舍?

  • ConnectableFlux 的 publish/connect/autoConnect 机制
  • share() 即 publish().refCount(1) 的语义
  • refCount 的引用计数与断开时机

ConnectableFlux 是"可连接"的发布者,通过 publish() 获得,数据不会在订阅时自动开始流动,必须调用 connect() 才启动;autoConnect(n) 在达到 n 个订阅者时自动连接,之后不再断开。share() 是 publish().refCount(1) 的便捷封装,第一个订阅者到来时自动连接,最后一个订阅者取消时自动断开(refCount),从而把一个冷流变成"即时共享"的热流。取舍上:若需要手动控制"何时开始发射"(如等待所有订阅者就绪),用 connect();若希望"有订阅者就自动开始、无订阅者自动停止",用 share()/refCount();若需要"永久连接、不因订阅者数量变化而断开",用 autoConnect(1)。refCount 的额外成本是需跟踪订阅者增减,且可能因订阅者全部取消而重新连接。

核心区别在于"连接生命周期由谁控制"。connect 手动控制、autoConnect 按订阅者数量自动连接且不回头、refCount 按订阅者数量动态连接与断开。选择取决于源是"需要精确控制开始时机"还是"跟随订阅者自动迁移"。

#
★★

5. Flux.cache/replay 的热订阅

Flux.cache 与 Flux.replay 在实现"热订阅"(缓存重放)上有什么区别?各自适合什么场景?

  • cache 缓存所有/最近元素并重放给新订阅者
  • replay 可配置缓存大小与时间窗口
  • 缓存的内存与过期语义

cache 会把流中的元素缓存起来,后续订阅者会立即收到缓存的历史元素(默认缓存全部 or 可按 cache(history) 限定条数),再继续接收新元素,适合"昂贵的计算结果需要被多个订阅者共享复用"的场景。replay 是更灵活的缓存重放:replay(n) 仅缓存最近 n 个元素,replay(Duration) 只在时间窗口内缓存,replay(n, Duration) 组合两者。两者都把"冷流"变成"可重放的历史 + 实时续流"的热发布源。取舍上:需要完整历史用 cache;需要限长或限时窗口用 replay,以控制内存占用。都要注意缓存可能导致内存增长,需权衡缓存条数/时间与内存。

cache 与 replay 的本质都是"缓存重放",区别在于缓存范围:cache 默认全量(可用参数限定),replay 提供更细粒度的条数/时间窗口控制。缓存是"以内存换复用",需平衡内存开销。

#
★★

6. Flux.create/Flux.generate 的差异与取舍

Flux.create 与 Flux.generate 有什么区别?各自适合什么场景?

  • create 支持异步多线程发射、可跨线程
  • generate 同步单线程、逐状态生成
  • 背压(create 需配合溢出策略)与可取消性

Flux.create 允许从任意线程(包括异步回调)多次调用 emitter.next()/complete()/error() 发射元素,适合桥接非响应式事件源(如监听器、回调、异步 API),但需要为溢出指定策略(如 onBackpressureBuffer),且 emitter 可被取消,具备灵活性。Flux.generate 是同步、单线程、逐状态生成器:通过一个可变的内部状态(S)在每次调用时生成一个元素并更新状态,适合有界/深度遍历、递归生成等需要"由状态驱动"的场景,背压天然精确(每次只发射一个)。取舍上:桥接外部异步事件源用 create;从状态生成序列(如会话、递归结构)用 generate。create 更灵活但需注意线程安全与背压,generate 更可控但局限于同步单向生成。

create 是"外部事件拉入",generate 是"内部状态推导"。create 面向异步多源、天然多线程,generate 面向同步单源、状态机驱动。背压处理上 create 靠溢出策略,generate 靠每次一步的同步发射。

#
★★

7. Flux.defer/using/usingWhen 的资源管理

Flux.defer、Flux.using、Flux.usingWhen 在资源管理上有什么区别?各自适合什么场景?

  • defer 延迟创建/求值
  • using 同步资源构造与释放(finally 语义)
  • usingWhen 异步资源及响应式释放

Flux.defer 让资源或数据源的创建延迟到订阅时执行,每次订阅都重新执行 Supplier,避免"创建时固化的状态",本身不管理资源释放。Flux.using 用于管理"同步创建的资源":在订阅时创建资源,在流完成/取消/异常时确保资源被释放(相当于必执行的 finally),适合需要"获取-使用-释放"的资源(如连接、文件句柄)。Flux.usingWhen 是响应式版本的资源管理:资源创建与释放都是响应式的(返回 Publisher),适合资源获取/释放本身是异步操作(如获取数据库连接、异步释放)的场景,能正确传播取消与异常。取舍上:仅需延迟求值用 defer;同步资源绑定生命周期用 using;异步资源生命周期用 usingWhen。

三者是"延迟求值"到"生命周期管理"的递进。defer 只管延迟,using 管同步资源释放,usingWhen 管异步资源释放。真正的资源管理(防止泄漏)依赖 using/usingWhen,且 usingWhen 能正确处理响应式取消。

#
★★

8. Flux.interval/Mono.delay 的定时器

Flux.interval 与 Mono.delay 在定时器功能上有什么区别?各自适合什么场景?

  • interval 周期性发射递增数字(0,1,2...)
  • delay 延迟后发射一个信号
  • 定时器的调度器与取消

Flux.interval(period) 会按固定周期周期性发射从 0 递增的 Long 值,适合周期性任务(如心跳、轮询、定时刷新),它默认使用 Schedulers.parallel() 工厂,可通过 interval(duration, scheduler) 指定调度器。Mono.delay(duration) 则延迟指定时间后发射一个信号(0 值),只发射一次,适合"延时后执行单个动作"或作为定时触发。两者都可通过操作符取消(dispose)或配合 take 限制次数。工程上:interval 用于周期性心跳/轮询,delay 用于单次延迟。注意 interval 若 .take(n) 限制次数,否则会无限发射;定时器需在取消时正确释放避免线程泄漏。

区别是"周期重复"与"单次延迟"。interval 是周期生成器,delay 是单次定时器。两者都涉及调度器选择,周期性任务需注意资源释放与取消。

#
★★

9. Flux.map/flatMap/concatMap/switchMap/merge/zip 的语义

Flux.map、flatMap、concatMap、switchMap、merge、zip 这几个操作符的语义分别是什么?如何区分?

  • map 一对一同步转换
  • flatMap 一对多异步展平(交错,不保序)
  • concatMap 一对多但保序(按顺序订阅)

map 对每个元素做同步一对一转换,元素数量不变。flatMap 把一个元素映射为一个或多个内部 Publisher 并展平合并,内部流并发订阅,元素到达顺序由各流完成时间决定(不保证顺序),适合"每个元素触发一个异步子任务"。concatMap 也是展平,但保证内部流按顺序依次订阅处理,结果是保序的,适合需要保持顺序的异步处理。switchMap 在内部流尚未完成时,一旦上游新元素到来就取消当前内部流切换到新的,只保留最新内部流,适合"以最新请求为准"(如搜索联想)。merge 并发合并多个独立流,zip 按位置配对合并。选用时:保序用 concatMap,需并发且不关心顺序用 flatMap,需响应最新用 switchMap。

区分核心是"并发与顺序"。map 同步一对一;flatMap 并发展平不保序;concatMap 串行展平保序;switchMap 只响应最新。它们共同支撑了响应式最核心的"异步组合"能力。

#
★★

10. Flux.retry/retryWhen/repeat 的重试边界

Flux.retry、Flux.retryWhen、Flux.repeat 在重试与重复上有什么区别?重试的边界如何界定?

  • retry 在 onError 后重新订阅(次数的边界)
  • retryWhen 自定义重试策略(条件、退避、截止时间)
  • repeat 在 onComplete 后重新订阅(成功后的重复)

retry(n) 在流发射 onError 后重新订阅上游,最多重试 n 次,针对的是"失败后重试";retryWhen 允许自定义重试策略,用函数对错误信号流做转换(如退避、按错误类型过滤、限制总时间),更精细。repeat(n) 则是在流成功 onComplete 后重新订阅,实现"成功后重复执行",repeat 与 retry 针对的信号不同(complete vs error)。重试边界需要明确:重试次数上限、重试条件(何种错误可重试)、退避策略(固定/指数/抖动)、总截止时间(避免无限重试)。非幂等操作重试需谨慎,可能造成重复执行。

retry 针对 onError,repeat 针对 onComplete,这是最本质的区别。retryWhen 提供重试策略的完全控制,是生产环境处理重试的首选。边界界定(次数、条件、时间)是防止重试雪崩的关键。

#
★★

11. Flux.take/takeWhile/skip/takeUntil 在截断与取消语义上的差异,takeUntil 触发后如何取消上游

Flux.take、takeWhile、skip、takeUntil 在截断与取消语义上有什么差异?takeUntil 触发后如何取消上游?

  • take 取前 n 个后取消上游
  • takeWhile 在条件不满足时截断并取消
  • skip 跳过前 n 个

take(n) 取前 n 个元素后自动请求取消上游(cancel),停止接收。takeWhile(predicate) 持续取元素直到条件不满足,一旦条件为假就截断并取消上游。skip(n) 跳过前 n 个元素,继续接收后续。takeUntil(predicate) 一直取元素直到条件第一次为真,一旦条件满足就截断并取消上游(注意 takeUntil 会包含触发条件为真的那个元素本身,而 takeWhile 不包含第一个不满足的元素)。takeUntil 触发后,Reactor 会通过 cancel() 取消上游订阅,同时若上游是可取消的(如 interval、自定义源),会触发 onCancel 钩子释放资源。这些操作符都通过"向下游发出后取消"实现截断,避免继续接收多余元素。

截断操作符都伴随"取消上游"语义,区别在于截断条件:按数量(take)、按条件持续(takeWhile)、按条件触发(takeUntil)、跳过(skip)。takeUntil 含触发元素、takeWhile 不含首个不满足元素,这是易混淆点。取消通过 cancel() 信号传播。

#
★★

12. Flux.transform/transformDeferred 的可重用

Flux.transform 与 Flux.transformDeferred 在"可重用操作符链"上有什么区别?

  • transform 立即应用转换函数(订阅前)
  • transformDeferred 延迟到订阅时应用
  • 复用操作符链与上下文敏感性的差异

transform 立即对当前 Flux 应用给定的转换函数(返回新的 Publisher),在装配时执行一次,适合"无状态、可复用"的操作符链封装。transformDeferred 则把转换函数延迟到每次订阅时才应用,适合"转换逻辑依赖订阅时上下文(如当前订阅者的线程、装配参数、或每次订阅需重新创建内部状态)"的场景。两者的区别类似于"立即求值"与"延迟求值":transform 复用的是一个装配好的固定链,transformDeferred 每次订阅都重新执行转换函数。若转换函数内部捕获了订阅时才有的状态(如 Mono.defer 的用途),应使用 transformDeferred 以避免状态固化。

transform 与 transformDeferred 的核心是"求值时机"。transform 立即装配一次,transformDeferred 每次订阅重新求值。选择依据是转换函数是否依赖订阅时的上下文或需要每次新建内部状态。

#
★★

13. Mono.just/Mono.empty/Mono.error 的语义

Mono.just、Mono.empty、Mono.error 三个工厂方法分别表示什么语义?各自适合什么场景?

  • just 发射一个给定值
  • empty 发射空(无值、正常完成)
  • error 发射错误信号

Mono.just(value) 创建一个立即发射给定值的 Mono,表示"确定有一个值"的结果,适合已知的单值装配。Mono.empty() 创建不发射任何元素、正常完成的 Mono,表示"没有结果但不是错误"(如查询无记录),是 0..1 语义中"0"的情况。Mono.error(Throwable) 创建发射错误信号的 Mono,表示"立即失败",用于表示故障条件或测试错误路径。三者覆盖了"成功有值 / 成功无值 / 失败"三种结果形态。工程上,empty 与 error 常用于条件分支:无数据返回 empty,异常返回 error。

三者是 Mono 三种终止状态的构建:just 发射值并完成,empty 只完成不发射,error 只发错误。理解它们有助于正确表达"结果是否存在 / 是否失败"。注意 just 是立即求值,若需延迟求值应使用 defer。

#
★★

14. Mono.then/when 与多源编排

Mono.then、Mono.when 在多源编排上有什么区别?各自适合什么场景?

  • then 忽略前序结果、串联执行
  • when 并行等待多个 Mono 完成(zip 风格)
  • then/thenReturn/thenMany 的变体

Mono.then 忽略前一个 Mono 的结果,等前序完成后执行后续(可用于顺序化多个副作用操作),返回一个 Mono ;thenReturn 在完成后发射固定值;thenMany 在完成后切换为另一个 Flux。Mono.when 则并行订阅多个 Mono/Publisher,等待它们全部完成(结果被忽略,返回 Mono),适合"并行执行多个任务并等待全部完成"的编排,类似 CompletableFuture.allOf。取舍上:需要顺序链式执行(后序依赖前序完成)用 then 系列;需要并行执行多个独立任务并同步等待用 when。若需要保留多个结果可用 zip 或 tuple。

区别是"顺序"与"并行"。then 强调顺序串联(前序是后序的依赖),when 强调并行等待(各任务独立)。then 返回 Void 忽略结果,when 聚合多个源等待完成。选择取决于任务间依赖关系。

#
★★

15. ParallelFlux 的并行处理

ParallelFlux 是什么?它如何实现并行处理?与 reactive 的并发订阅有什么区别?

  • ParallelFlux 通过 runOn 把数据分发到多个线程并行处理
  • 与普通 Flux 的并发订阅(flatMap 并发)的区别
  • 并行度与顺序保证

ParallelFlux 是通过 Flux.parallel(n) 得到的并行流类型,它把数据流分割为 n 个"轨道"(rails),通过 runOn(scheduler) 将各轨道分发到不同线程并行处理,从而真正利用多核 CPU 实现并行计算。与 flatMap 的并发不同:flatMap 是通过"并发订阅多个内部 Publisher"实现交错,而 ParallelFlux 是"把同一数据流分发给多个轨道并行处理"。ParallelFlux 适用于 CPU 密集型的有界数据处理(如对大批量数据做计算),配合 .sequential() 可转回普通 Flux 继续下游处理。注意:ParallelFlux 不保证顺序,并行度需与 CPU 核数匹配,且数据源需支持背压。

ParallelFlux 的本质是"数据分轨并行"。与 flatMap 的并发订阅不同,它是把单一流按轨道分发到多线程。它适合 CPU 计算密集任务,需用 runOn 绑定调度器,用 sequential 汇合。

#

16. Schedulers.parallel()/single()/immediate() 的工程应用

Schedulers.parallel()、single()、immediate() 这三种调度器分别适合什么工程场景?

  • parallel 固定大小并行线程池,适合 CPU 密集
  • single 单线程,适合时间敏感/低开销任务
  • immediate 当前线程立即执行

Schedulers.parallel() 提供一个固定大小(默认 CPU 核数)的并行线程池,适合 CPU 密集型计算任务,如 ParallelFlux 的 runOn、计算密集的 map/flatMap。Schedulers.single() 提供单个线程,适合时间敏感、低开销、需要串行执行的任务(如定时器、单线程状态机),避免线程切换开销。Schedulers.immediate() 表示在当前线程立即执行,不切线程,适合"无需切换线程、希望直接在下游线程执行"的场景,或用于测试。工程上:CPU 密集用 parallel,轻量定时/串行用 single,无需切换用 immediate。注意阻塞任务不应使用 parallel/single,而应使用 boundedElastic。

三种调度器对应不同的线程模型:parallel 并发、single 串行、immediate 当前线程。选择依据是任务类型与线程切换成本。阻塞任务应避开这些,用 boundedElastic。

#

17. Sinks(Sinks.Many/Sinks.One)作为热发布者

Sinks.Many 与 Sinks.One 是什么?如何用它们构建热发布者(程序化数据源)?

  • Sinks 的用途(程序化地向流中推送数据)
  • Sinks.Many 的多种背压模式(unicast/multicast/replay)
  • Sinks.One 的配对(单值)语义

Sinks 是 Reactor 提供的高层"程序化数据源"API,用于从外部代码向响应式流中推送数据,替代已废弃的 Processor。Sinks.Many 支持多元素推送,有几种模式:unicast 只允许一个订阅者、multicast 允许多个订阅者(无背压协调)、replay 缓存历史重放给新订阅者;通过 tryEmitNext 推送元素并返回成功/失败结果。Sinks.One 是单值配对语义,适合表示"一个异步结果"(类似 CompletableFuture),通过 tryEmitValue 或 tryEmitError 完成。Sinks 作为热发布者,通常暴露为 Flux 给下游订阅,再用 tryEmit 推送数据。使用时需注意 tryEmit 返回的 emitResult 需处理(如缓冲满、失败)。

Sinks 是构建自定义热发布源的标准方式,替代 Processor。区别在于 Sinks.Many 多元素、Sinks.One 单值。tryEmit 返回结果需处理,以应对背压或失败。它是把"事件驱动代码"接入响应式流的桥梁。

#

18. 错误处理(onErrorReturn/onErrorResume/onErrorContinue)

onErrorReturn、onErrorResume、onErrorContinue 这三个错误处理操作符有什么区别?各自适合什么场景?

  • onErrorReturn 失败时返回默认值
  • onErrorResume 失败时切换到备用流
  • onErrorContinue 错误后继续处理后续元素

onErrorReturn 在流发出 onError 时返回一个默认值,让流以正常值结束,适合"可接受用默认值兜底"的场景。onErrorResume 在 onError 时订阅另一个备用 Publisher(可按错误类型选择不同备用流),适合"失败后可切换到备选数据源"的场景(如降级为缓存数据)。onErrorContinue 与其两者不同:它不终止流,而是跳过出错的那个元素,继续处理后续元素,适合"批处理中单个元素失败不应中断整批"的场景(如逐条处理记录时某条失败可跳过)。注意 onErrorContinue 只对某些操作符(如 map、flatMap)的错误生效,且它改变了错误语义,可能掩盖错误,需谨慎使用。三者区别是"终止后兜底 / 终止后切换 / 不终止继续"。

三者针对错误处理的不同策略:onErrorReturn 返回默认值终止、onErrorResume 切换备用流终止、onErrorContinue 跳过错误元素继续。选择依据是"出错后是终止整个流还是跳过单个元素",以及是否可降级。