如何从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,需要优化你的轮询逻辑:
现有实现的问题
- 无间隔轮询:
repeat()会在每次查询结束后立刻发起下一次请求,会给MongoDB和客户端造成不必要的压力 - 多订阅冲突:
AtomicLong放在方法外部时,多个订阅者会共享同一个状态,容易出现数据重复或遗漏 - 代码错误:第一个代码片段里
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
相关产品推荐
相关产品推荐

