消息中间件选型从来不是一件孤立的事。有的团队从早年就基于ActiveMQ搭建了消息总线,交易通知、订单流转都跑在上面;后来为了应对高并发的实时场景,又引入了Redis做缓存和轻量级消息推送。两套系统各自工作没问题,麻烦出现在业务需要跨系统协作的时候:ActiveMQ里的消息希望实时通知到只订阅Redis的客户端,Redis里产生的事件也需要进入ActiveMQ的队列参与后续的可靠处理。这就牵扯到一个不太常被讨论的话题——协议适配。

一、先弄清楚两者的消息模型和协议差异
Redis本身没有严格意义上的消息协议,它的发布订阅是基于RESP协议之上的几个命令实现的。PUBLISH负责向某个channel发布消息,SSUBSCRIBE或者SUBSCRIBE负责订阅。这种模型非常简单直接,消息发出去之后不落盘、不确认、没有消费者组概念(Stream出现之前),一旦没有订阅者在线,消息就丢了。
ActiveMQ则完全不同,它实现了完整的JMS规范,还支持OpenWire、STOMP、MQTT等多种传输协议。JMS模型里有Destination的概念,分为Queue(点对点)和Topic(发布订阅)两种,消息带持久化标志、优先级、过期时间,消费端有CLIENT_ACKNOWLEDGE等确认模式。这些语义层面的差异,是做协议适配时最需要小心的地方。
简单总结一下对照关系:Redis的一个channel大致对应ActiveMQ的一个Topic;Redis Stream的consumer group则更接近Queue的消费语义。适配层要做的核心工作,就是把一边的消息语义翻译成另一边能理解的形式。
二、适配层的整体设计与消息格式映射
推荐的做法是单独部署一个桥接服务,我们姑且叫它bridge。它对内维护两个连接池:一个连Redis,一个连ActiveMQ。方向上有两条通路,一条是Redis到ActiveMQ,一条是反向的。为了两边消息结构统一,先定义一个中立的JSON消息信封格式。
public class BridgeMessage {
private String messageId; // 全局唯一ID,用于去重和追踪
private String topic; // 逻辑主题名
private String payload; // 实际业务内容
private long timestamp; // 消息产生时间戳
private int deliveryMode; // 1 非持久化,2 持久化
}
// 主题映射规则示例:
// redis:order.notify <-> activemq:TOPIC::ORDER.NOTIFY
// redis:stream:task <-> activemq:QUEUE::TASK.QUEUE主题映射不能靠硬编码散落在各处,建议集中放在一个配置类或者配置中心里,格式上约定一套命名规范。比如ActiveMQ的主题名统一大写下划线分隔,Redis的channel用小写点号分隔,bridge启动时加载映射表并校验两侧的Destination是否都存在。
消息体的序列化也要统一。ActiveMQ默认用的是ObjectMessage或者TextMessage,跨语言场景下TextMessage加JSON是最稳妥的选择。Redis的PUBLISH只传字节数组,所以bridge在收到消息后统一反序列化为BridgeMessage,再根据目标端的格式重新组装。这样即使未来要接入第三种中间件,只需要增加新的转换器,不动核心逻辑。
三、两个方向的具体实现
1. Redis到ActiveMQ方向
这个方向用Redis Stream更可靠。如果业务侧只能用Pub/Sub,bridge就要承担消息可能丢失的风险,最好在Redis侧开启AOF并配合适当的刷盘策略来缓解。bridge作为消费者读取Stream,转发到ActiveMQ时一定要等ActiveMQ的send调用成功返回后再ACK这条Stream消息,保证至少一次投递。
@Component
public class RedisToActiveMqBridge {
private final StringRedisTemplate redisTemplate;
private final JmsTemplate jmsTemplate;
public void startForward(String streamKey, String queueName) {
String groupName = "bridge-group";
try {
redisTemplate.opsForStream().createGroup(streamKey, groupName);
} catch (Exception ignore) {
// 组已存在则忽略
}
String consumer = "bridge-" + UUID.randomUUID();
while (!Thread.currentThread().isInterrupted()) {
List<MapRecord<String, Object, Object>> records =
redisTemplate.opsForStream().read(Consumer.from(groupName, consumer),
StreamReadOptions.empty().count(50),
StreamOffset.create(streamKey, ReadOffset.lastConsumed()));
if (records == null || records.isEmpty()) {
continue;
}
for (MapRecord<String, Object, Object> record : records) {
String payload = String.valueOf(record.getValue().get("payload"));
// 先成功发送到ActiveMQ,再确认Redis
jmsTemplate.convertAndSend(queueName, payload, msg -> {
msg.setJMSDeliveryMode(DeliveryMode.PERSISTENT);
return msg;
});
redisTemplate.opsForStream().acknowledge(streamKey, groupName, record.getId());
}
}
}
}注意一个细节:jmsTemplate默认开启了异步发送的话,convertAndSend返回不代表broker已经收到。生产环境建议显式配置为同步发送,或者依赖ActiveMQ的连接工厂设置useAsyncSend为false,否则极端情况下会出现Redis已ACK而ActiveMQ实际没收到消息的空洞。
2. ActiveMQ到Redis方向
反向通路用Spring的JMS监听器来实现。bridge作为ActiveMQ的消费者订阅Topic或Queue,收到消息后写回Redis。如果目标是通知类场景,直接PUBLISH到对应channel即可;如果下游还有消费可靠性的要求,改成写入Redis Stream,让下游用消费组模式读取。
@Component
public class ActiveMqToRedisBridge {
private final StringRedisTemplate redisTemplate;
@JmsListener(destination = "ORDER.NOTIFY", containerFactory = "topicFactory")
public void onMessage(TextMessage message) throws JMSException {
String body = message.getText();
// 方案一:直接发布到Redis channel,低延迟但可能丢
redisTemplate.convertAndSend("order.notify", body);
// 方案二:写入Stream,保证下游可回溯
Map<String, String> fields = new HashMap<>();
fields.put("payload", body);
fields.put("jmsMessageId", message.getJMSMessageID());
redisTemplate.opsForStream().add("order.notify.stream", fields);
}
}两种写入方式可以并存,但要控制好Redis的内存。Stream有长度上限的设置,通过MAXLEN参数裁剪历史消息,比如只保留最近十万条,避免桥接服务长期运行把Redis内存吃满。
四、确认机制对接与重复消费处理
协议适配里最容易翻车的就是确认语义。ActiveMQ的消费者确认是事务性或客户端确认模式,broker在收到ACK之前会保留消息;而Redis这边除了Stream的XACK,基本没有确认概念。所以在bridge的设计里,必须明确一条铁律:先在目标端完成写入,再对源端做确认。上一节的代码已经体现这个顺序,这里再强调一次,因为一旦顺序反了,故障恢复时就可能出现消息既没送达又被确认的丢失窗口。
顺序正确带来的是至少一次投递语义,代价是可能重复。比如bridge发送到ActiveMQ成功,但还没来得及XACK就崩溃了,重启后这条消息会被再次投递。因此下游消费者必须做好幂等处理,常见做法是利用消息信封里的messageId做去重,在数据库或Redis里记录已处理的ID集合,设置一个合理的过期时间即可。
还有一种思路是在桥接服务里引入本地消息表,收到消息先落库标记为待转发,转发成功后更新状态,通过定时任务扫描重发。这种方式适合对可靠性要求极高的场景,代价是吞吐量会下降,链路延迟也会增加几百毫秒,需要根据业务权衡。
五、部署与监控要点
bridge服务本身要支持水平扩展。Redis Stream的消费组天然支持多消费者负载均衡,多个bridge实例分别以不同consumer name读取同一个组即可。ActiveMQ侧的Queue也可以多实例并行消费,但Topic监听会出现每个实例都收到一份消息的情况,转发到Redis时就会重复。解决办法是要么只部署一个Topic转发实例配合选主机制,要么在Redis侧靠messageId去重。
监控指标至少要包括:两侧的连接状态、每秒转发的消息数、转发失败的次数、以及Redis Stream的积压长度(XLEN和XPENDING)。积压长度持续增长往往意味着ActiveMQ那边出问题了,及时告警能避免故障扩大。日志层面建议给每条转发消息打上messageId做全链路追踪,出问题时能快速定位是哪个环节丢的。
最后一点经验之谈:协议适配是解决历史包袱的过渡手段,不是终态架构。如果业务允许,长期方案还是尽量收敛到统一的消息中间件上,bridge的存在是为了给迁移争取时间,而不是让两套系统永久并行。明确这个定位,才能在设计上把握好投入的分寸。