如何用Spring WebFlux构建反应式内存仓库?实现问题求助
如何用Spring WebFlux构建反应式内存仓库?
先聊聊你遇到的问题和背后的原因,再给你正确的实现方案:
为什么event.subscribe(events::add)没效果?
- 异步执行特性导致的感知缺失:
subscribe()是异步触发订阅逻辑的,但你的save方法返回void,调用方根本无法感知这个异步操作什么时候完成。往往方法刚返回,订阅逻辑还没来得及执行,自然看不到元素被添加到列表里。 - 非线程安全集合的隐患:你用的普通
List不是线程安全的,WebFlux默认在多线程调度器下运行,直接操作这种集合会导致不可预期的行为,甚至出现元素丢失的情况。 block()的“饮鸩止渴”:虽然block()能让代码“看起来生效”,但它会阻塞当前线程,完全违背了WebFlux非阻塞的设计初衷,还可能引发事件循环阻塞的严重问题。
正确的反应式内存仓库实现方式
首先,先修正你的仓库接口——反应式方法必须返回Publisher类型(比如Mono/Flux),而不是void,这是反应式编程的核心原则之一:
public interface EventRepository { // 返回Mono<Void>传递保存操作完成的信号 Mono<Void> save(Mono<Event> event); // 返回所有事件的Flux流 Flux<Event> findAll(); }
然后实现这个接口,用线程安全的集合配合正确的反应式操作:
import org.springframework.stereotype.Repository; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.util.Queue; import java.util.concurrent.ConcurrentLinkedQueue; @Repository public class InMemEventRepository implements EventRepository { // 使用线程安全的ConcurrentLinkedQueue存储数据,适配多线程环境 private final Queue<Event> events = new ConcurrentLinkedQueue<>(); @Override public Mono<Void> save(Mono<Event> event) { // 不要用内部subscribe,而是通过doOnNext在元素流向下游时完成添加操作 // then()返回Mono<Void>,给调用方传递操作完成的信号 return event.doOnNext(events::add) .then(); } @Override public Flux<Event> findAll() { // 从线程安全队列创建Flux流 return Flux.fromIterable(events); } }
进阶:支持实时更新的仓库(可选)
如果你的需求是让findAll()的订阅者能实时获取新添加的事件(类似消息推送),可以用ReplayProcessor实现一个可追加的持久化流:
import org.springframework.stereotype.Repository; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.core.publisher.ReplayProcessor; @Repository public class ReactiveInMemEventRepository implements EventRepository { // ReplayProcessor会保存所有历史元素,新订阅者能拿到全量历史数据,同时接收新元素 private final ReplayProcessor<Event> processor = ReplayProcessor.createUnbounded(); @Override public Mono<Void> save(Mono<Event> event) { // 将事件推送到processor,然后返回完成信号 return event.doOnNext(processor::onNext) .then(); } @Override public Flux<Event> findAll() { return processor; } }
关键注意点
- 绝对不要在WebFlux业务代码中随意用
block():除非是初始化等特殊场景,否则会彻底破坏非阻塞模型。 - 必须使用线程安全集合:WebFlux是多线程环境,普通集合会引发并发问题。
- 反应式方法务必返回
Mono/Flux:让调用方通过订阅来控制操作的执行和完成,而不是在内部偷偷调用subscribe。
内容的提问来源于stack exchange,提问作者user14900768
相关产品推荐
相关产品推荐

