You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Spring WebFlux构建反应式内存仓库?实现问题求助

如何用Spring WebFlux构建反应式内存仓库?

先聊聊你遇到的问题和背后的原因,再给你正确的实现方案:

为什么event.subscribe(events::add)没效果?

  1. 异步执行特性导致的感知缺失:subscribe()是异步触发订阅逻辑的,但你的save方法返回void,调用方根本无法感知这个异步操作什么时候完成。往往方法刚返回,订阅逻辑还没来得及执行,自然看不到元素被添加到列表里。
  2. 非线程安全集合的隐患:你用的普通List不是线程安全的,WebFlux默认在多线程调度器下运行,直接操作这种集合会导致不可预期的行为,甚至出现元素丢失的情况。
  3. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.11 08:26:06