如何实现天气消息Flux按Site ID分组缓存最后12条后合并及多订阅?
解决方案
你的问题出在cache(12)的使用场景上——它无法让后续订阅者共享每个站点的历史缓存。当第二个订阅者订阅时,flatMap会重新订阅每个分组流,而cache(12)仅为首次订阅的上下文保留缓存,不会跨订阅者复用历史数据。
正确的做法是为每个站点的分组流配置可重放且多播的缓存,使用replay(12).autoConnect()替代cache(12),代码如下:
weatherMessagesSink .asFlux() .groupBy(weatherMessage -> weatherMessage.getSiteID()) .flatMap(groupedFlux -> groupedFlux.replay(12) // 固定缓存每个站点的最后12条消息 .autoConnect() // 自动连接,允许多订阅者共享缓存内容 )
为什么这个方案有效:
- 每个站点独立缓存:
replay(12)会为每个Site ID对应的分组流单独保留最后12条消息,不受其他站点消息速率或数量变化影响,完全符合你的需求。 - 支持多订阅者:
autoConnect()会让每个分组的重放流在有第一个订阅者时自动订阅源流,后续订阅者订阅时,直接获取该站点已缓存的12条消息,同时接收后续新消息。 - 适配站点动态变化:新增站点时,
groupBy会自动生成新的分组流,replay(12)会为其缓存后续消息;站点停止发送消息后,缓存的12条数据会保留到所有订阅者取消订阅(如果需要自动清理闲置站点的缓存,可以添加超时逻辑,比如replay(12).autoConnect(0, conn -> conn.disposeAfter(Duration.ofMinutes(30))))。
相比你提到的“缓存4*12条”的方案,这个实现更可靠,能精准保证每个站点的缓存数量,完全适配站点新增/移除、消息速率差异的场景。
内容的提问来源于stack exchange,提问作者Dean A
相关产品推荐
相关产品推荐

