ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

ActiveMQ生产环境高级特性指南:异步投递、死信队列与存储选型

ActiveMQ生产环境高级特性指南:异步投递、死信队列与存储选型 很多朋友学了ActiveMQ的基本API、Queue/Topic用法之后都觉得“消息中间件也就这样了”。但真把应用丢到生产环境一套并发压上去问题就全冒出来了消费者明明很空闲消息端到端延迟却高得离谱有人认为丢了几条消息无所谓结果订单状态对不上死信队列天天爆运维天天找你。这篇番外想讲的恰恰是入门教程里不常讲、但线上真正决定生死的ActiveMQ高级特性——异步投递与确认模式、存储引擎选型、死信队列与重投策略、虚拟主题、延迟投递、多协议接入以及Spring Boot整合时那些坑。这里的内容基于ActiveMQ 5.x尤其是5.18的日常使用经验适合已经跑通过Demo、准备把这些知识沉淀到生产项目的读者。1. 异步投递与消息确认高性能与不丢消息之间的权衡1.1 别让发送链路死在“同步确认”上先说我第一次被ActiveMQ性能教育的故事。当时我给一个对账系统发消息每天几百万条单条消息也就几百字节。压测时发现生产者吞吐死活上不去CPU也没跑满链路里全是等待。排查到最后问题出在同步确认上。JMS规范里持久化消息发送后生产者通常要等broker的确认回执ActiveMQ的OpenWire协议默认也是这么干的客户端发送持久化消息broker把日志刷到磁盘、分配完MessageID再回一个Response。这个过程本身没问题但一次磁盘fsync哪怕只花1到5毫秒在单线程高并发场景下吞吐就卡死在这个“每次发送都等一次落盘反馈”的串行流程里。如果业务能接受极端情况下的少量消息丢失可以在连接URL里打开异步发送tcp://127.0.0.1:61616?jms.useAsyncSendtrue打开后客户端不再等待broker的同步确认消息直接扔进网络缓冲区由协议层的心跳和重连机制兜底。实测下来在短小消息、高频率发送的场景里吞吐量提升往往不止一个数量级。代价也很明确如果发送后进程瞬间宕机或者broker在处理这批消息前崩溃这些“已发送但未确认”的消息可能就没了。所以异步发送适合日志、通知、监控采样这类丢一条不致命的场景不适合订单、支付、库存等强一致链路。还有一个细节容易被忽略ActiveMQ对非持久化消息默认就是异步发送对持久化消息默认才是同步。也就是说如果消息本身是NON_PERSISTENT即使不配置useAsyncSend它也是异步路径。很多人排查性能问题时没注意这个点以为“我没开异步消息肯定都同步落盘”结果消息早就走的是非持久化异步通道。提醒一下打开异步发送后消息仍然可以在producer端设置DeliveryMode.PERSISTENT这不冲突。前者管的是客户端与broker之间的交互方式后者管的是broker自己要不要落盘。两个维度不要混。1.2 消费者的ACK模式不是小事消费者端的确认模式是另一个容易写错“看上去没问题”的地方。JMS里有四种Session模式对应关系如下Session模式确认语义典型适用场景AUTO_ACKNOWLEDGE消费方法返回后自动确认通知、日志、不敏感数据CLIENT_ACKNOWLEDGE调用message.acknowledge()手动确认业务处理成功后才确认DUPS_OK_ACKNOWLEDGE允许延迟确认、可重复投递对偶发重复消息不敏感的场景SESSION_TRANSACTED事务提交时统一确认需要与业务操作强绑定先说AUTO_ACK。这个模式下ActiveMQ默认是“消息交给应用后立即确认”看起来最稳。但如果配合了optimizeAcknowledge优化broker会攒一批消息再批量确认这能显著降低确认报文数量代价是如果连接在这期间断开一批已投递的消息会重新排队表现为消费者收到重复消息。所以开了optimizeAcknowledge一定要做好消费幂等。再说CLIENT_ACK。这里有个经常踩到的坑JMS规范里客户端调用acknowledge()时确认的是当前Session里所有已消费未确认的消息而不是只确认当前这一条。在ActiveMQ里这个行为默认也是这样。如果你在循环里消费多条消息只对最后一条调acknowledge()前面那些消息可能全部被确认了一旦后面某条业务处理失败你想“我只让这一条重投”是做不到的。解决办法有几种要么收到消息后立即ack要么在连接URL上加参数让它变成单条确认语义要么干脆用事务会话。事务会话是我自己最常用的模式。消息onMessage里处理业务如果成功就session.commit()失败就session.rollback()回滚后broker会重新投递这条消息并且遵循redeliveryPolicy的重投次数限制。语义简单直接不容易出现“以为自己确认了其实没确认”的混乱。1.3 我在生产里怎么选我的习惯是按业务等级分三类第一类纯通知、日志、实时指标。用AUTO_ACK加optimizeAcknowledge消息不持久化追求吞吐。丢了就丢了业务本来也不依赖它。第二类工单、审批、报表任务这类“丢了会麻烦但能补偿”的场景。消息持久化消费者用CLIENT_ACK或事务业务处理成功后才确认。第三类订单状态、支付结果回调、库存变动。必须持久化加事务同时消费端实现幂等控制因为无论怎么设计分布式环境下都无法百分之百避免重复消息。这里有个特别实在的体会不要指望“确认模式”帮你解决所有可靠性问题。确认模式保住的只是“broker到消费者”这一段业务逻辑里漏了幂等照样会出事。2. 存储引擎KahaDB与JDBC怎么选消息落盘背后发生了什么2.1 KahaDB内部机制ActiveMQ默认的数据存储是KahaDB很多人只知道“默认就是它”不清楚它到底怎么工作。你可以把它想成一个带预写日志的小型数据库broker收到消息后先把消息顺序追加到data log文件里这个过程是纯顺序写速度很快然后内存里的索引和缓存做更新后台通过checkpoint把内存状态定期固化并回收不再需要的日志文件。理解了这个机制很多现象就解释得通了。比如“为什么broker重启后恢复这么慢”很可能是data log文件太多、checkpoint间隔太长启动时要把未固化的日志重新扫描建立索引。再比如“为什么KahaDB目录在不断变大”可能是因为大量小消息一直堆积在未回收的journal里而清理线程还没轮到它。KahaDB常用的几个配置项directory数据存储目录生产环境务必放到独立磁盘别和系统盘、日志盘混用。journalMaxFileLength单个journal文件的最大字节数默认32MB。消息偏小可以调小一点减少空间浪费消息偏大就调大一点减少文件数量。checkpointIntervalcheckpoint间隔默认5秒。太短会频繁刷盘影响性能太长会拖慢异常恢复速度。indexCacheSize索引缓存大小默认10000页。消息种类多、队列多的时候可以适当加大。2.2 什么时候才应该考虑JDBC存储JDBC存储是被问得最多的“高级特性”之一但我必须说实话它不是用来提升性能的而是用来满足“消息数据可以进数据库统一备份、审计、对账”这类需求的。比如公司要求所有业务数据必须落在数据库里或者你要做多broker共享同一套消息状态这时候才考虑配置JDBC Persistence Adapter。配置方式大致是在activemq.xml里把persistenceAdapter换成jdbcPersistenceAdapter并给它配一个DataSourcebroker xmlnshttp://activemq.apache.org/schema/core brokerNamelocalhost useJmxtrue persistenceAdapter jdbcPersistenceAdapter dataSource#mysql-ds createTablesOnStartuptrue / /persistenceAdapter /broker用JDBC存储消息读写要走数据库事务延迟和吞吐完全取决于数据库的IO能力。数据库一旦抖动broker跟着抖动。我在生产里见过最典型的失败案例团队为了“更可靠”把ActiveMQ切到MySQL存储结果数据库主从切换时消息发送直接超时最后又老老实实改回KahaDB。如果你确实必须用JDBC存储建议做好三件事数据库连接池给足、SQL执行慢查询日志打开、broker和数据库放在同机房低延迟网络里。另外不要把生产数据库和ActiveMQ的库放在同一个实例上互相拖累的教训太多了。2.3 消息过期策略与磁盘水位消息设置了TTL过期后不会像你想象的那样“到点就消失”。ActiveMQ有过期消息扫描机制但它是周期性的不是实时的。消息过期后先由后台任务识别并标记然后清除。默认的过期消息检测周期是可以调的在policyEntry里可以设置expireMessagesPeriod之类的参数。如果ExpiryScan执行不及时你会在JMX里看到很多expired消息还占着内存。磁盘水位的处理也类似。KahaDB会周期性地执行数据清理把已经消费掉、不再被引用的事务日志回收掉。但如果你长时间不消费、队列积压又大数据文件不会自动减少磁盘占用会持续上升。这时候优先查消费者而不是疯狂地调GC参数。3. 死信队列与重投策略别等到线上积压才来求医3.1 一条消息什么时候会被丢进DLQ死信队列DLQ在ActiveMQ里叫ActiveMQ.DLQ默认情况下queue和topic类型消息的死信都进这个队列。如果你发现死信队列里的消息堆积如山先别急着清理想想这些消息是怎么进去的消费者处理抛出异常导致消息被回滚或主动拒绝重投次数达到maximumRedeliveries上限。事务会话里rollback次数达到上限。消息超过有效期在特定配置下进入DLQ。消费者端没有正确确认broker反复投递最终被判定为“毒消息”。一条消息反复重投是有代价的。每次重投都占网络、占Broker内存、占日志文件写入如果后台有一个“永远处理不了”的消息它能把整个队列的消费速度拖垮。所以重投策略绝不只是DLQ的事它直接影响正常消息的吞吐。3.2 redeliveryPolicy参数拆解在ActiveMQ里控制重投行为的主要是RedeliveryPolicy。我贴一个实际用过的配置片段policyEntry queue redeliveryPolicy maximumRedeliveries3 initialRedeliveryDelay1000 useExponentialBackOfftrue backOffMultiplier2 useCollisionAvoidancefalse / /policyEntry几个参数的实际含义maximumRedeliveries最大重投次数默认6。不是说重投6次就停而是达到上限后消息转入DLQ。initialRedeliveryDelay第一次重投前的延迟默认1秒左右。如果业务处理失败是由于外部服务暂时不可用一失败就立刻重投往往是最差的选择。useExponentialBackOff与backOffMultiplier是否指数退避以及每次退避的倍率。比如第一次等1秒第二次等2秒第三次等4秒给下游系统留出恢复时间。useCollisionAvoidance碰撞避免给延迟加一点随机扰动防止大量消息同时重投造成“惊群”。这个在纯技术文档里看到频率很高但实际业务场景用得不多除非你确实遇到过所有失败消息同时重试把下游打挂的情况。这些参数可以配置在broker的destinationPolicy里也可以在客户端ConnectionFactory上单独设置。我的建议是全局设置一个保守值具体业务队列用policyEntry单独覆盖不要一锅端。3.3 我处理DLQ的实际套路死信消息不是垃圾它只是“目前处理不了”的消息。我的日常处理分三步第一步先在线排查。用管理控制台或JMX看一下ActiveMQ.DLQ里的消息条数、消息大小、入队时间分布。如果堆积时间集中在某个业务发布窗口多半是发布时的代码问题回滚重启业务即可死信一般不会继续增长。第二步抽样查看死信内容。消息体里的业务字段、报错堆栈、第一次投递时间都是线索。最快的办法是用QueueBrowser浏览消息不用消费掉直接看属性。这里有个小技巧死信消息的原始Queue名可以通过JMSDestination属性看方便你定位“这消息本来该进哪个队列”。第三步决定回投、丢弃还是转人工。回投就是把死信消息重新send回原队列Queue originalQueue (Queue) deadLetterMessage.getJMSDestination(); producer.send(originalQueue, deadLetterMessage);注意回投后消息会变成一条全新的消息重投计数清零业务方会再次收到。如果消息本身业务语义已经过时比如超时订单回投只会让它再次死掉不如走人工补偿通道。这里还想提醒一点如果某个队列的死信消息长期持续产生别每次只清死信一定要查根因。死信队列本质是业务系统bug的“指示灯”你一直把红灯按掉问题不会消失。4. 虚拟主题与延迟投递订阅形态的高级玩法4.1 虚拟主题为什么能解决“订阅者重启丢消息”用Topic做广播时持久订阅者离线期间的消息会累积在broker上等它回来这听着没问题但真实系统里每个消费者实例的消费进度、恢复、集群扩容都不好管。虚拟主题Virtual Topic换了个思路让每个消费者实例背后都有一条物理队列。具体玩法是这样的。生产者发送消息到VirtualTopic.orders消费者端监听的不是一个Topic而是Consumer.组名.VirtualTopic.orders。ActiveMQ会在broker里为这个消费者自动创建一条队列并让虚拟主题上的消息复制一份进去。这样做的好处非常直接消费者实例可以随时重启队列里的消息不会丢因为它是持久队列。新增消费组只需要新增一个Consumer前缀不需要改生产者。某个消费组处理慢不会影响到其他消费组。我在实际项目里用虚拟主题替代了很多原本用Topic加持久订阅的场景。比如订单状态变更需要让订单服务、库存服务、风控服务分别消费互不干扰虚拟主题是最省心的方案。4.2 定时与延迟投递ActiveMQ原生支持延迟与定时投递前提是broker启动时开启schedulerSupportbroker xmlnshttp://activemq.apache.org/schema/core brokerNamelocalhost schedulerSupporttrue发送延迟消息时给消息加上对应的属性message.setLongProperty(ScheduledMessage.AMQ_SCHEDULED_DELAY, 5000); producer.send(message);常用的三个属性AMQ_SCHEDULED_DELAY延迟多少毫秒投递。AMQ_SCHEDULED_PERIOD周期性投递的间隔。AMQ_SCHEDULED_CRONCron表达式按日历时间触发。用这个做订单超时检查、优惠券到期提醒、定时调度任务都很方便。不需要额外部署一套调度中间件。不过要注意定时消息在broker端会有一个专门的调度线程维护消息量大且到期时间集中时会带来一定的CPU峰值。另外如果broker集群是多实例要保证定时消息的调度不会多实例重复执行。4.3 镜像队列看着很美落地时别抱太大期望镜像队列Mirrored Queue可以理解为把写入某个Queue的消息复制一份发布到一个镜像Topic上方便做监控或者实时分析。配置大概长这样destinationInterceptors mirroredQueue copyMessagetrue postfix.mirror / /destinationInterceptors但实际生产里我基本不会依赖它来做核心逻辑。一个是它不能保证复制操作的原子性在极端情况下原队列消息和镜像消息可能不一致另一个是它会对所有队列做隐式复制增加无谓的存储和网络开销。如果你只是想给每个消费组独立进度虚拟主题才是正道只想做审计不如直接在消费端做一份日志转发。5. 多协议接入AMQP、MQTT与JMS的协议枢纽5.1 ActiveMQ是一个协议枢纽而不是只能JMS很多团队的现状是Java服务用JMSPython/Node服务用AMQP物联网设备用MQTT旧系统用STOMP。如果每个协议都各自搭一套消息中间件运维成本和跨系统打通都很痛苦。ActiveMQ的定位恰恰是可以同时开多个transportConnector让不同协议连接到同一套broker、同一套destination上。在conf/activemq.xml里你可以这样开多协议端口transportConnectors transportConnector nameopenwire uritcp://0.0.0.0:61616/ transportConnector nameamqp uriamqp://0.0.0.0:5672/ transportConnector namemqtt urimqtt://0.0.0.0:1883/ transportConnector namestomp uristomp://0.0.0.0:61613/ /transportConnectors开端口很简单麻烦的是协议间的语义差异。5.2 AMQP 1.0与JMS的映射关系先提个醒ActiveMQ的AMQP支持是基于AMQP 1.0的不是RabbitMQ那种AMQP 0-9-1客户端库选型时别搞混。AMQP 1.0的Terminus、Source/Target概念和JMS的Queue/Topic不是一一对应ActiveMQ在内部会做一层映射发送到AMQP address时目标会被解析成queue或topic消息的header、properties也会有一部分落到JMS消息属性里。实际对接时最容易出的问题有两个。一是AMQP客户端里设置的message-format、annotationJMS消费者拿到的message属性可能对不上二是AMQP 1.0事务和JMS本地事务的交互比想象中复杂。我的建议是跨协议互通只传递简单的文本或JSON消息体别依赖复杂的头字段。任何涉及事务、高可靠性保证的链路尽量保持两端协议一致。5.3 MQTT QoS与保留消息物联网设备一般用MQTT接入broker端会把MQTT消息映到内部的Topic上。MQTT的QoS 0/1/2在ActiveMQ里会和JMS的某些机制融合但这个映射不是无损耗的。比如MQTT的retained message保留消息在ActiveMQ里会作为一条消息发布到对应Topic新订阅者能否收到这条保留消息取决于你的topic策略配置和连接时机。如果你打算在物联网场景直接用ActiveMQ做MQTT broker我建议先做小流量的协议语义验证重点测试QoS 1和持久会话下的离线消息恢复行为。ActiveMQ的MQTT支持满足大部分场景但它不是EMQX那样的专业物联网broker连接数上万、或者要求细粒度会话管理时还是换专业产品更稳妥。6. Spring Boot整合事务配置、监听器与5.18落地经验6.1 连接工厂和连接池设置Spring Boot整合ActiveMQ最常见的配置大概是这样spring: activemq: broker-url: tcp://127.0.0.1:61616?jms.useAsyncSendtrue user: admin password: admin pool: enabled: true max-connections: 20 expiry-timeout: 30000生产环境我建议打开连接池否则每次JmsTemplate.send和监听器创建都新建连接连接创建成本高连接数还容易把broker打爆。但连接池不是越大越好连接数上去了消费者端的prefetch和事务并发就更容易互相干扰。有个真实的坑是连接池开了之后很多人忘记设置max-connections默认值其实是很大的压测时broker侧会看到成百上千个连接内存占用直接上了一个台阶。我习惯把连接池压到实际并发的一半左右再配合监控慢慢调。6.2 JmsListener的并发与确认Spring里监听器默认容器工厂的行为和原生JMS不同。JmsListener出现异常时消息会回到broker等待重投如果你希望“失败就死”或者“失败就进DLQ”要在onMessage里自己catch异常并按需确认。并发配置也不能乱填。我看过有人把concurrency配成50-100结果一个队列同时起了100个消费者会话每个会话默认prefetch值又不低消息瞬间被拉走broker内存被打满。合理做法是从1-3起步根据实际单条消息处理耗时长压测。还有一个容易踩的点监听器里通过JmsTemplate再往其他队列发消息如果不做事务绑定业务处理成功但发送失败时两边状态就分叉了。6.3 事务与幂等Spring Boot里让JmsTemplate走会话级事务比较简单jmsTemplate.setSessionTransacted(true);这样发送操作会在一个本地事务里事务提交时消息才真正发送发送失败会抛出异常回滚。但有一点必须讲清楚如果业务操作涉及关系型数据库你希望“数据库更新成功消息才发出去数据库更新失败消息也不发”跨数据库和ActiveMQ两者之间没有全局事务靠本地事务是不够的。JTA/XA理论上可以做到但ActiveMQ的XA支持在5.x里并不让人省心我见过太多项目为了XA引入分布式事务框架最后复杂度远超收益。如果业务需要这个级别的原子性更好的方案是消息本地消息表或者事务性发件箱模式至少可控、可排查。6.4 5.18版本与工程实践的取舍看到热搜里有人问“activemq 5.18下载”说明不少人在关注新版本。5.18整体上还是Classic的延续但如果从老版本直接升上来有两个方面值得留意一是运行环境的要求对JDK版本有基线部署前先确认好二是内置Web控制台相关组件有所调整控制台能看的东西和旧版不完全一样。我的建议是生产环境升级前先在测试环境跑一遍消息收发、虚拟主题、死信重投这些核心路径对比JMX指标而不是只看“版本号变了没有”。另外新装ActiveMQ时默认用户名密码一定要改这属于安全红线但每年都能看到因为没改默认密码导致的消息丢数据事故。改完密码后还要把broker的JMX访问权限也检查一遍。7. 故障排查与调优从指标、日志到消息丢失定位7.1 先看懂Broker的六项核心指标排查ActiveMQ问题第一步永远是看指标乱猜只会浪费时间。我常用的指标就这几个指标含义关注点EnqueueCount入队消息总数生产者是否正常发送DequeueCount出队消息总数消费者是否正常消费InFlightCount已投递但未确认的消息消费者处理是否积压PendingCount队列中等待投递的消息积压水位MemoryPercentUsagebroker内存使用率是否触顶限流StorePercentUsage持久化存储水位磁盘健康度打开ActiveMQ Web控制台点进一个Queue这些指标都在。我最依赖的是PendingCount和InFlightCount的组合Pending大说明生产快消费慢InFlight大说明消息被消费者拉走了但迟迟不确认这是消费端问题。7.2 慢消费者、pending积压与内存限幅当某个消费者处理消息很慢又不确认时broker上的InFlight消息会堆积。ActiveMQ默认开启了flow controlbroker内存达到阈值后会阻塞生产者继续发送让生产端变慢甚至报sendFailIfNoPendingConnection超时。这是保护机制不是bug。排查顺序一般是看队列入队速度是否正常先确认是不是生产者突增。看出队速度如果出队接近0消费者基本是停了。看InFlight数量如果跟Pending差不多大说明消息全被消费者拉走但没确认。看消费者日志、线程池状态、下游数据库耗时定位是处理逻辑卡住还是外部依赖超时。消费者端的prefetch参数也影响这个局面。prefetch越大broker越愿意一次性把大量消息塞给消费者单线程处理能力跟不上时积压越明显。把prefetch调小比如1到10虽然单个消费者吞吐会降但积压更容易被多个消费者实例分摊。调prefetch不是解决慢消费者的根本手段但能有效缓解“消息全堵在某个实例的本地内存里”的尴尬。7.3 消息丢失别急着怪broker先走排查链路“消息不见了”是ActiveMQ群里最高频的求助。我总结了一条排查链路按这个顺序走基本能定位到90%的问题确认生产者发送时有没有收到异常或回调。如果发送端异常被吞了消息根本没出去。确认消息是不是落到了别的队列或Topic。目标打错、虚拟主题前缀写错消息会跑到你看不见的地方。在broker日志里搜这个MessageID看broker是否确实接收、分发。查该队列的DLQ。很多“丢消息”其实是消息进了ActiveMQ.DLQ只是你平时不关注。查消费者的确认逻辑。AUTO_ACK模式下onMessage返回就算确认如果消息在异步处理流程里失败了也“看起来成功”。这套链路走完八成以上都能定位到问题。还剩下一成的“真的没找到日志”十有八九是消息是非持久化的broker重启时丢了。7.4 我积累的调优习惯最后分享几个不太会写进官方文档的调优习惯。首先是不要盲目加参数。ActiveMQ的配置项非常多乱加只会让问题更难排查。我一般只动这几个维度连接数、消费者并发、prefetch、redeliveryPolicy、内存和存储目录。每次只改一个参数观察一整天指标再决定下一步。其次是磁盘和存储水位监控一定要做。炸存储比炸内存更可怕内存炸了大不了重启broker恢复存储炸了连恢复日志的空间都没有。我的告警设置是StorePercentUsage超过70%开始关注超过85%立刻人工介入。最后是关于升级。ActiveMQ的高版本通常带来更好的性能和更活跃的社区支持但升级前必须做全量回归。我经历过的几次生产事故有一种就是升级后配置兼容问题导致的比如某个XML配置在旧版本是warning新版本直接拒绝启动。看release notes别跳过。做消息中间件这行最深的体会是ActiveMQ本身的问题远少于使用方式的问题。把异步发送当成“百分百不丢”把AUTO_ACK当成“百分百可靠”把Topic当成万能广播这些误解才是线上事故的源头。高级特性不是让你一次性全上而是让你在每一个选择里都知道自己放弃了什么、换来了什么。
返回列表