Skip to content

Redis Pub/Sub

跨实例消息广播(缓存失效通知、实时推送扇出等)。ace-redis 提供两种消费方式:配置桥接(推荐,声明式)与 RedisSubscriber 直接订阅。发布侧始终是一行:

cangjie
redis.publish("orders", "{\"id\":42}")   // 返回收到消息的订阅者数

方式一:桥接到事件总线(推荐)

配置要桥接的通道,收到的每条消息自动经应用事件总线发布为 RedisMessage 事件:

toml
[redis]
"pubsub.bridgeChannels" = "orders, notifications"   # 逗号分隔
cangjie
@Service
public class OrderListener {
    @EventListener
    public func onMessage(e: RedisMessage): Unit {
        // e.channel 触发通道,e.payload 消息正文
        if (e.channel == "orders") {
            println("order event: ${e.payload}")
        }
    }
}

生命周期完全托管:组件启动建订阅、停机自动关闭(先停订阅读协程再关客户端池)。

方式二:RedisSubscriber 直接订阅

需要动态订阅/退订、glob 模式匹配时直接使用:

cangjie
let sub = RedisSubscriber(RedisConfig.defaults())

sub.subscribe("orders", {channel, payload =>
    println("${channel}: ${payload}")
})
sub.psubscribe("logs.*", {channel, payload =>   // glob 模式,回调收到实际触发通道
    println("${channel}: ${payload}")
})

sub.unsubscribe("orders")
sub.punsubscribe("logs.*")

sub.close()   // 幂等;必须调用,否则读协程与连接不释放

实现特性(使用时需要知道的)

  • 独立专用连接:RESP2 订阅态连接只能收推送、不能复用于普通命令,所以订阅不占客户端连接池;首次 subscribe 才惰性建连并启动唯一读协程。
  • 订阅连接不设读超时:普通命令连接受 timeoutMs 约束,订阅连接特意不设(长时间无消息是常态,设了会被误判断连)。
  • handler 异常不杀协程:回调抛异常被捕获打印,订阅继续存活;但 handler 内不要做长阻塞操作(所有通道共享一个读协程,会阻塞后续消息分发)。
  • 断线不自动重连:连接意外断开时读协程打印日志退出,已订阅关系不会自动恢复。对可靠性有要求的场景配合 pingOnReady/监控告警,或在业务层重建订阅。
  • close 协议close() 置停止标志并关连接,读协程随即静默退出——组件桥接模式下由框架停机流程自动完成。

基于 Apache-2.0 许可证发布