标签: Shovel

  • 深入 RabbitMQ 陷阱排查:滥用双向 Shovel 引发的环路风暴与全局水位阻塞实战

    某次核心支付系统的异步回调链路突发大面积超时,API 网关 99 线从 50ms 直接飙升至 30s 并伴随大量 504 Gateway Timeout。排查结论令人啼笑皆非:某位业务开发为了实现所谓的“跨机房双活容灾”,在没有任何路由防环设计的情况下,通过 RabbitMQ 管理控制台手动配置了双向 Shovel 插件。结果导致消息在两个集群间形成无限死循环复制,瞬间产生的消息风暴击穿了节点内存,触发了 Erlang VM 的 vm_memory_high_watermark 告警,底层的 TCP 背压(Backpressure)机制直接将所有 Producer 的 Connection 强行置为 blocking 状态,最终引发了波及全业务线的全局雪崩。

    不要把消息队列当成可以随意拉线的网络集线器,在没有深刻理解 AMQP 路由拓扑和底层流控机制前,任何“高可用”架构的尝试都无异于自掘坟墓。

    案发现场:全线假死与消失的吞吐量

    故障发生时,监控大盘上呈现出极其诡异的景象:

    1. QPS 归零:业务网关请求堆积,RabbitMQ 集群的 Inbound 流量在经历了几秒钟的垂直飙升后,瞬间掉底为 0。

    2. CPU 与 Load 暴增:宿主机 Load Average 飙升至 80+,epmdbeam.smp 进程 CPU 占用率满载。

    3. 海量 Connection 被 Block:应用侧疯狂打印 java.util.concurrent.TimeoutException

    登录 RabbitMQ 节点,敲下排查命令,惨烈的情况一览无余:

    # 查看当前连接状态,发现大量连接处于 blocking 或 blocked 状态
    $ rabbitmqctl list_connections pid name port state | awk '{print $4}' | sort | uniq -c
        152 running
       2048 blocking
        512 blocked
    
    # 查看资源告警状态
    $ rabbitmq-diagnostics alarms
    Alarms on node rabbit@mq-node-01:
    [x] memory alarm: true (Memory high watermark set to 0.4. Current usage: 14.2 GB / 32 GB)
    

    查看核心日志 /var/log/rabbitmq/[email protected],满屏的红色警告:

    202X-XX-XX 14:05:12.123 [warning] <0.1453.0> memory resource limit alarm set on node rabbit@mq-node-01.
    202X-XX-XX 14:05:12.124 [info] <0.1455.0> blocking connection <0.2312.0> (10.0.5.12:45123 -> 10.0.2.10:5672)
    202X-XX-XX 14:05:12.124 [info] <0.1455.0> blocking connection <0.2313.0> (10.0.5.13:42123 -> 10.0.2.10:5672)
    ...
    

    很明显,Erlang VM 的内存使用率超过了设定的阈值(默认 40%),RabbitMQ 启动了极端的自我保护机制:全局内存告警阻塞

    拨开迷雾:愚蠢的“跨机房双活”拓扑

    RabbitMQ 的 vm_memory_high_watermark 触发后,所有发布消息(Publish)的连接都会被底层的 TCP 层面挂起。这不是针对单个 VHost 或 Queue 的限制,而是全局核武级别的熔断,只要连在这个节点上发消息的 Client,全部都要死。

    是什么打爆了内存? 通过 rabbitmqctl list_queues name messages memory 发现,两个机房的核心 Topic Exchange 下绑定的队列消息堆积量在以每秒数十万的速度递增。

    进一步排查拓扑配置,真相大白。业务侧通过 Shovel 插件做了如下配置:

    • 机房 A (Shovel-A): Source: Exchange 'pay.topic' (RoutingKey: '#') -> Dest: URI of DC-B / Exchange 'pay.topic'

    • 机房 B (Shovel-B): Source: Exchange 'pay.topic' (RoutingKey: '#') -> Dest: URI of DC-A / Exchange 'pay.topic'

    这就是典型的“无脑双向复制”引发的广播风暴。

    AMQP 协议中的 Shovel 本质上是一个运行在 Erlang VM 内部的客户端。它在源端执行 basic.consume,在目的端执行 basic.publish。 当一条路由键为 pay.success 的消息在机房 A 产生时:

    1. 机房 A 的 Exchange 将其路由到本地队列,同时 Shovel-A 将其拉取。

    2. Shovel-A 将该消息 basic.publish 到机房 B 的 pay.topic

    3. 机房 B 的 Exchange 接收到消息,不仅路由给 B 的本地队列,同时被 Shovel-B 捕获。

    4. Shovel-B 再次将其发回给机房 A…

    一条消息在毫秒级内变成了几万条,呈指数级放大,瞬间榨干网络带宽并击穿了 14GB 的内存水位。

    为什么说这个错误不可原谅?

    如果是单纯为了做高可用和跨集群复制,官方早就提供了 Federation 插件。为什么 Federation 不会环路而 Shovel 会?这是协议层设计的降维打击。

    Federation 插件在跨节点投递消息时,会在 AMQP Header 中注入 x-received-from 属性。 当机房 B 的 Federation 收到来自机房 A 的消息时,检查 Header 发现这条消息曾经来过,或者达到了配置的 max_hops 阈值,就会直接丢弃,从根源上阻断了环路。

    而该业务团队因为“嫌 Federation 配置策略复杂,Shovel 看起来就像个搬运工比较简单”,直接用了 Shovel。要知道,Shovel 是无状态的,它不管消息从哪里来,只负责傻瓜式地搬运,根本没有防环机制。更要命的是,他们在 Topic 匹配上用了最暴力的 #,将整条业务线推向了深渊。

    破局与防御性架构落地

    应急恢复非常粗暴:

    1. 立刻通过 CLI 强制删除双向的 Shovel 链路:rabbitmqctl clear_parameter -p / shovel shovel-a

    2. 执行 rabbitmqctl purge_queue 清空由于环路产生的海量垃圾消息,让内存水位降至 0.4 以下。

    3. 观察 alarm 解除,TCP 连接恢复 running 状态,业务网关自动重连恢复。

    针对此类惨案,运维和架构层面必须落地以下防御性策略:

    1. 废弃控制台 ClickOps,收归配置权限: 禁止任何人通过 Management UI 手动拉取跨机房链路。所有的 Shovel/Federation Policy、Exchange、Binding 配置,必须通过 Terraform 或 Ansible 以 IaC(基础设施即代码)的形式进入 GitOps 流程,强制进行拓扑评审。

    2. 正确使用高可用组件: 跨集群双活/复制,首选 Federation,并严格配置 max-hops = 1。如果非要用 Shovel,路由键必须加上机房前缀(如 dc-a.pay.#),并且 Shovel 目的端只允许写入带有特定后缀的隔离 Exchange。

    3. 多租户与 VHost 物理隔离: 所有核心业务线必须拆分物理集群,至少也要做到 VHost 级别的隔离,并对每个 VHost 限制 max-lengthmax-length-bytes,防止单一野鸡业务把全局水位打爆。

    排查清单:RabbitMQ 内存阻塞与环路问题速查

    1. 确认全局资源告警阻塞 (TCP Backpressure) rabbitmq-diagnostics alarms 如果存在 memory alarm: truedisk_free alarm: true,说明 Broker 已启动自我保护,所有发布消息的 Connection 已被挂起(State: blocking/blocked)。

    2. 快速定位堆积/异常队列 rabbitmqctl list_queues name messages memory message_bytes | sort -k4 -nr | head -n 10 查出占用内存或消息体总和最大的 Top 10 队列,如果是极短时间内暴增,高度疑似环路风暴。

    3. 排查 Shovel / Federation 配置状态 rabbitmqctl list_parameters -p [vhost] 检查是否存在双向配置的参数。对于 Federation,检查 rabbitmqctl federation_status 的链路是否有报错。

    4. 验证连接状态统计 rabbitmqctl list_connections state | grep -c blocking 当出现大量 blocking 连接时,切勿盲目重启应用,需优先解决 MQ 服务端的资源水位问题,否则应用重启后仍会卡死在建立 AMQP Channel 的握手阶段。

  • 深入 RabbitMQ 跨机房雪崩排查:Shovel 环形路由风暴引发的内存高水位封控与 Paging IO 抖动实战

    某次接手处理一个跨机房双活架构的突发故障,业务端疯狂报错 java.util.concurrent.TimeoutException,所有往 RabbitMQ 集群投递消息的生产者全部卡死。登录管控台一看,双机房的 RabbitMQ 节点内存全部顶到告警线,连接状态齐刷刷显示为 blocked。 最终排查发现,这是一个极其低级的架构配置失误:业务侧通过 HTTP API 动态下发了双向 Shovel 任务进行跨机房消息同步,但既没有规划隔离的 Routing Key,也没有利用 Header 进行防环判断。一条消息在两个机房之间构成了无限死循环(Infinite Routing Loop),引发指数级的消息放大。RabbitMQ 在触发 vm_memory_high_watermark 保护机制后,无差别封杀所有生产者 TCP 连接,随后触发海量内存数据 Paging 刷盘,直接把底层存储 IOPS 打满,导致整个消息总线瘫痪。

    跨机房同步不用自带防环机制的 Federation,反而去手捏底层的 Shovel,捏完还不做防环逻辑。这种把插线板插在自己身上企图获得无限能源的操作,是对分布式系统基本功的严重亵渎。

    案发现场:诡异的 Blocked 连接与暴涨的内存

    监控大屏上的指标非常刺眼:

    1. Message Rate 异常:入队速率(Publish)从平时的 3k/s 瞬间飙升到 80k/s,而出队速率(Deliver/Get)几乎跌零。

    2. 连接状态死锁:执行 rabbitmqctl list_connections pid client_properties state,发现数万个生产者连接的 state 全部处于 blockingblocked 状态。

    3. 节点内存报警:系统内存 32G,RabbitMQ 进程占用飙破 12.8G(默认 40% 阈值)。

    4. 日志报警:核心日志里疯狂刷出 alarm_handler 触发的告警: log [warning] <0.324.0> memory resource limit alarm set on node 'rabbit@node1'. [info] <0.326.0> connection <0.1122.0> (10.x.x.x:54321 -> 10.x.x.y:5672): connection is blocked

    深度剖析:环形风暴与 Erlang VM 内存防御机制

    为什么一条循环消息能让整个 RabbitMQ 集群雪崩?这涉及 AMQP 协议的路由盲区以及 Erlang VM 激进的防御机制。

    1. Shovel 双向死环的形成

    在跨机房同步场景中,RabbitMQ 官方推荐的 Federation 插件会在消息 Header 中隐式追加 x-received-from 标记。当节点发现消息的流转链路中已经包含自己的集群名时,会主动丢弃,从而天然防环。 但排查过程中发现,业务侧为了“灵活控制路由”,选择使用了更底层的 Shovel 插件。Shovel 的本质是一个伪装成客户端的 Erlang 进程,它在一端 Consume,在另一端 Publish。 配置示例还原:

    • 机房 A Shovel:源端 Exchange=order.topic,目标端 机房 B Exchange=order.topic

    • 机房 B Shovel:源端 Exchange=order.topic,目标端 机房 A Exchange=order.topic

    由于两者监听的 Routing Key 均为 # 且目标 Exchange 相同,机房 A 产生的一条真实订单消息,被 Shovel 搬运到机房 B 后,立刻被机房 B 的 Shovel 捕获,再次搬回机房 A。消息在两条千兆专线间以网卡极限速度疯狂打乒乓球。

    2. vm_memory_high_watermark 的“休克疗法”

    RabbitMQ 不是以丢消息为代价来保命的系统。当节点内存达到 vm_memory_high_watermark(默认总内存的 0.4 倍)时,RabbitMQ 会触发一种近乎物理断电的保护机制: 底层 Erlang 会调用 erlang:setopts(Socket, [{active, false}]),直接停止读取所有发布消息的 TCP Socket。 这导致操作系统的 TCP 接收缓冲区迅速填满,TCP 窗口滑动为 0(Zero Window),反压(Backpressure)传导至客户端,最终导致所有的 Spring AMQP / Celery 生产者线程因等不到 ACK 甚至无法建立 Socket 发送而全部 Block 阻塞,业务雪崩。

    3. Paging 刷盘引发的 IO 惨案

    内存触顶后,噩梦才刚刚开始。为了腾出内存,RabbitMQ 会根据 vm_memory_high_watermark_paging_ratio(默认 0.5,即达到内存水位线的 50% 时触发)策略,将内存中的瞬态消息(Transient Messages)和队列索引强行 Page Out 到磁盘的 msg_store_transient 目录。

    # 查看内存破拆情况
    rabbitmq-diagnostics memory_breakdown
    # 输出显示 msg_index 和 queue_procs 占据了绝大部分内存
    

    几十万条循环堆积的消息瞬间引发极高频率的随机写 IO,导致磁盘 %%util 打满 100%,iowait 飙升。此时哪怕你想通过命令行去删除队列,都会因为底层 Mnesia 数据库及 Erlang 进程的 IO 阻塞而超时失败。

    破局与防御性修复

    在 IO 打满、连接全卡死的状态下,常规操作已经失效,必须通过底层干预进行“放水排雷”。

    1. 紧急提水位,恢复管控权 必须先骗过 Erlang VM,让它以为内存还够,从而恢复 TCP 处理和管控台响应:

    # 临时将内存告警阈值从 0.4 提至 0.6,争取操作窗口
    rabbitmqctl set_vm_memory_high_watermark 0.6
    

    2. 斩断死环,清理积压 在争取到的几分钟窗口期内,立刻删掉引发风暴的 Shovel 配置,并暴力清空积压队列:

    # 删除恶意 Shovel (注意:需在目标 VHost 下执行)
    rabbitmqctl clear_parameter -p /my_vhost shovel my_evil_shovel_a2b
    
    # 清洗队列(比从 UI 点 Purge 更稳)
    rabbitmqctl purge_queue -p /my_vhost loop_queue_name
    

    3. 架构级防御加固 恢复后,必须进行彻底的架构重构,杜绝此类问题二次发生:

    • 弃用双向 Shovel,改用 Federation:如果非要用双向同步,强制使用 Federation 插件,利用其内置的 x-received-from Header 实现拓扑防环。

    • 如果是 Shovel 刚需,必须做 Header 路由过滤:在 Shovel 配置中注入特定的 Header(例如 add_forward_headers),并在接收端的 Exchange 之前挂载一个 Headers Exchange 进行逻辑判断,拒收带有该机房标记的消息。

    • 死信与 TTL 兜底:任何跨系统调用的队列,绝对不允许无限期堆积。强制设置 x-message-ttlx-max-length。消息堆满立刻进 DLX(死信交换机),并配合报警,将故障控制在局部。

    总结排查清单

    为了避免后续运维和开发再踩坑,总结同类问题速查清单如下:

    1. 连接 Blocked 速查:遇到大量连接呈 blocking/blocked,第一时间看管控台右上角 Node 状态,如果是红色 Memory,说明已触发内存高水位封控,直接查 vm_memory_high_watermark

    2. 路由死环预警:排查有无异常的高 Message Publish 速率。如果有,且入队等于出队,极大概率是 Dead Letter Exchange (DLX) 配置成了死环,或者是 Shovel/Federation 跨机房配置了镜像拓扑。

    3. Paging 引起的性能雪崩:如果 CPU Load Average 极高,且执行 rabbitmqctl 命令频繁超时,检查磁盘 IO 是否被 RabbitMQ 的 msg_store_transientmsg_store_persistent 目录写满。必要时临时调高内存阈值进行急救。

    4. 生产者防阻塞策略:业务代码严禁对 MQ 同步阻塞等待。必须配置 ConnectionFactory 的超时时间,并在框架层捕获 AmqpException 进行降级,防止 MQ 抖动直接把业务 Tomcat/Netty 线程池拖死。