Project Reactor冷热Publisher定义及常见算子属性相关疑问
Reactor冷热Publisher常见疑问解答
核心定义先对齐
通用冷热Publisher判断标准:
- 冷Publisher:直到订阅才触发事件生产,每个订阅者都会获得独立、完整的事件流
- 热Publisher:事件可在订阅前就开始生产,多个订阅者共享同一份事件流
1. just是hot算子吗?
是,属于Reactor里特殊的静态热Publisher。
它是热Publisher的特例:不会动态生成持续事件,只会在装配阶段就把传入的值捕获存好,后续所有订阅者拿到的都是同一份预先捕获的值,不会为每个订阅重新计算传入参数的值,完全符合热Publisher「多订阅共享同一份数据,不会为新订阅重新生产事件」的核心特征。
可以用一段简单代码验证:
AtomicInteger counter = new AtomicInteger(0); // 构造Mono时就会执行counter.incrementAndGet(),counter立刻变为1,无需等待订阅 Mono<Integer> testMono = Mono.just(counter.incrementAndGet()); // 两次订阅拿到的都是同一个值1,不会再触发counter自增 testMono.subscribe(System.out::println); testMono.subscribe(System.out::println);
2. 只要在装配阶段获取值的算子就属于hot类型吗?
装配阶段取值是核心判断依据之一,还需要同时满足第二个条件:所有订阅者共享同一份数据,不会为新订阅重新触发数据生产流程,两个条件都满足才会归为hot类型。
3. Flux.fromIterable(iterable)是hot算子吗?
不是,它属于典型的冷Publisher。
和just不同,fromIterable不会在装配阶段就把Iterable里的元素全部取出缓存,而是等到每次订阅的时候,才会遍历你传入的Iterable对象生成事件。可以用下面的代码验证:
List<Integer> dataList = new ArrayList<>(); dataList.add(1); Flux<Integer> testFlux = Flux.fromIterable(dataList); // 构造完Publisher之后新增元素 dataList.add(2); // 第一次订阅会输出1、2 testFlux.subscribe(System.out::println); // 再次新增元素 dataList.add(3); // 第二次订阅会输出1、2、3,每次订阅都会重新遍历Iterable testFlux.subscribe(System.out::println);
4. 怎么判断哪些Publisher会在装配阶段获取值?
最直接的验证方法是给构造Publisher的传入参数加上可观测的副作用,比如计数器自增、打印日志,观察副作用的执行时机:
- 如果还没调用
subscribe()方法,副作用就已经执行,说明是装配阶段取值 - 如果只有调用
subscribe()之后才会执行副作用,说明是订阅阶段取值
另外也可以通过算子参数类型快速判断:如果算子接收的是普通值参数,一般是装配阶段求值;如果接收的是函数式接口参数(比如Supplier、Function),一般是订阅阶段才会调用函数求值。
5. 常见的cold算子有哪些?
绝大多数Reactor提供的原生算子都是冷Publisher,典型包括:
- 延迟求值类:
Mono.defer、Flux.defer - 函数生成类:
Mono.fromSupplier、Flux.create、Flux.generate - 资源读取类:
Flux.fromStream、Flux.fromInputStream - 定时类:
Flux.interval、Mono.delay - 远程调用类:WebClient、R2DBC等数据库/网络组件返回的Mono/Flux
这些冷算子的共同特征是:每次订阅都会触发一次全新的事件生产流程,不同订阅者拿到的事件流相互独立。
内容的提问来源于stack exchange,提问作者RobertWang
相关产品推荐
相关产品推荐

