如何配置ActiveMQ topic批量消费 实现Mongo批量持久化
ActiveMQ Topic 批量消费落地MongoDB 实现方案
前置配置(ActiveMQ 消费者侧)
要实现批量消费首先要调整ActiveMQ的基础配置:
- 关闭自动ACK,启用客户端手动ACK模式,创建Session时指定参数即可:
// false表示不开启事务,第二个参数指定手动ACK Session session = connection.createSession(false, Session.CLIENT_ACKNOWLEDGE); - 调整Topic预取策略,将预取数量设置为你需要的批量大小100,有两种配置方式:
- 连接URL直接加参数:
tcp://broker地址:61616?jms.prefetchPolicy.topicPrefetch=100 - 代码中给连接工厂配置:
ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory(brokerUrl); factory.getPrefetchPolicy().setTopicPrefetch(100);
- 连接URL直接加参数:
- (可选)如果要求服务重启后不丢失掉线期间的Topic消息,需要配置持久化订阅:给连接设置唯一的clientId,创建消费者时指定持久化订阅标识即可。
核心逻辑实现
批量消费逻辑
- 本地维护一个长度上限为100的消息缓冲列表,同时设置3-5s的超时等待阈值,避免消息量少的时候长时间无法凑够100条导致写入延迟
- 循环拉取消息加入缓冲列表,直到列表长度达到100或者触发超时阈值,进入批量写入流程
批量写入与回执处理
- 将缓冲列表中的消息统一转换为MongoDB存储的Document对象集合
- 调用MongoDB批量写入接口,建议开启非有序写入配置,避免单条消息写入失败阻塞整批其他消息的写入:
List<Document> docList = messageBuffer.stream() .map(this::convertMsgToDocument) .collect(Collectors.toList()); // ordered设为false,单条失败不影响其他消息写入 InsertManyResult writeResult = mongoCollection.insertMany(docList, new InsertManyOptions().ordered(false)); - 按写入结果分别处理消息回执:
- 写入成功的消息,逐条调用
message.acknowledge()向Broker发送ACK确认 - 写入失败的消息,强转为ActiveMQMessage对象后调用
nack()方法,或者调用session.recover()将失败消息退回Broker等待重新消费
- 写入成功的消息,逐条调用
- 清空缓冲列表,进入下一轮批量消费流程
边界异常处理
- 如果批量写入过程中出现数据库连接断开、超时等全量失败的情况,直接调用
session.recover()将整批消息全部退回Broker,避免消息丢失 - 建议配置消息最大重试次数,多次写入失败的消息转入死信队列存储,避免无限重复消费占用资源
内容的提问来源于stack exchange,提问作者Nabeel Ahmad
相关产品推荐
相关产品推荐

