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

如何从Reactive MongoDB源获取无限Flux流?

实现Mongo文档的实时持续查询(响应式)

优先方案:用MongoDB Change Streams实现主动推送

MongoDB原生支持Change Streams,能监听集合的变更(插入、更新、删除等),这是官方推荐的实时数据推送方案——无需客户端轮询,MongoDB会主动把符合条件的新文档推送给你,效率和实时性都远高于轮询。

用ReactiveMongoTemplate实现的示例代码:

public Flux<Document> tail(String conversation) {
    // 构建聚合规则:只监听符合conversation条件的插入操作
    Aggregation aggregation = Aggregation.newAggregation(
        // 匹配目标对话的新插入文档
        Aggregation.match(Criteria.where("fullDocument.conversation").is(conversation)),
        // 只保留插入类型的变更事件
        Aggregation.match(Criteria.where("operationType").is("insert"))
    );

    return reactiveMongoTemplate.changeStream("你的集合名称", aggregation, Document.class)
        .map(ChangeStreamEvent::getFullDocument); // 提取出实际的文档内容
}

这个Flux会一直保持活跃,直到你主动取消订阅,有新匹配文档入库时会自动推送。


轮询方案的优化(针对你的现有实现)

如果因为环境限制没法用Change Streams,需要优化你的轮询逻辑:

现有实现的问题

  1. 无间隔轮询:repeat()会在每次查询结束后立刻发起下一次请求,会给MongoDB和客户端造成不必要的压力
  2. 多订阅冲突:AtomicLong放在方法外部时,多个订阅者会共享同一个状态,容易出现数据重复或遗漏
  3. 代码错误:第一个代码片段里startFrom.set(item.getSequence())的item未定义,应该是d.get("sequence")

优化后的轮询代码

public Flux<Document> tail(String conversation) {
    return Flux.defer(() -> {
        // 每个订阅者拥有独立的sequence追踪状态,避免多订阅冲突
        AtomicLong startFrom = new AtomicLong(0L);
        // 每隔2秒轮询一次,间隔可根据业务需求调整
        return Flux.interval(Duration.ofSeconds(2))
            .flatMap(tick -> reactiveMongoTemplate.find(
                Query.query(
                    Criteria.where("conversation").is(conversation)
                        .and("sequence").gt(startFrom.get())
                ), Document.class
            ))
            .doOnNext(document -> {
                long currentSeq = document.getLong("sequence");
                if (currentSeq > startFrom.get()) {
                    startFrom.set(currentSeq);
                }
            })
            .distinct(); // 避免轮询间隔内出现重复查询结果
    });
}

递归concat方案的问题

递归concat的逻辑会在每次查询完成后立刻发起下一次请求,同样存在无间隔轮询的性能问题;而且长期运行可能触发栈溢出风险(虽然Reactor是异步的,但递归深度累积仍有隐患),不推荐使用。


内容的提问来源于stack exchange,提问作者Paul

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 21:23:17