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

如何实现天气消息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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 23:20:30