TT Lab
开始
学习 学习路径 课程

构建 EAI 中间层

队列不丢消息,但会投递两次

在 TT Lab 中继续学习

一句话总结

异步对接约定的不是“现在就处理”,而是“不要弄丢,稍后再处理”。队列是履行这一约定的装置,但要让约定成立,发布一方必须做 Broker 已收到的确认(publisher confirms),消费一方必须做只在处理完之后才发送的确认(手动 ack),而作为代价,同一条消息可能会到达两次。

为什么需要它

设想保险理赔受理。客户在 App 上提交理赔单,渠道希望立刻显示“已受理”。然而理赔系统因为要核验资料,每笔需要几秒钟,到了月底还会积压几个小时。如果同步调用,渠道就被理赔系统的速度拖住,理赔系统只要短暂宕机,受理本身就会失败。在中间放一个队列,渠道可以在放入队列的那一刻就答复,理赔系统则按自己的速度取走。一方停了,另一方仍在工作。

作为代价,产生了新的问题。怎么知道放进去了?取走的一方在处理途中死掉,消息去哪里了?无法处理的消息(资料不全)会永远在队列里打转吗?理赔系统慢的时候,消费者要扛下多少笔?本模块将逐一回答这些问题。

工作原理

AMQP 0-9-1 的三个部件。发布者发送的不是队列,而是交换机(exchange)。交换机根据路由键和绑定把消息放入队列。direct 交换机会发往路由键与绑定键完全相同的队列(AMQP 概念)。正因为有这一步让发布者不知道队列名称,之后即使再绑定一个审计用的队列,发布者也无需修改。

durable 与 persistent 不同。durable 队列能让队列定义在 Broker 重启后仍然保留。要让消息保留,发布时必须以 persistent(delivery_mode=2)发送。两者只做一个,重启之后队列虽在,却是空的。

发布确认(publisher confirms)。basic_publish 无错误返回,并不意味着 Broker 已经收到。把通道设为确认模式,Broker 就会对每条消息返回 ack,按文档所述,发往 durable 队列的 persistent 消息,会在写入磁盘之后才确认(Confirms)。无处可路由的消息会被悄悄丢弃,但如果用 mandatory 发送,Broker 会在 ack 之前以 basic.return 退回。pika 的 BlockingChannel 会以 UnroutableError 告知这一点。

手动 ack。在自动 ack 模式下,Broker 一发送消息就算传递完成,所以消费者在处理途中死掉,这条消息就会消失。在手动 ack 下,消费者必须在处理完毕后发送 ack 才算结束。如果没有 ack 就关闭了通道,Broker 会自动把那条消息退回,并在再次投递时加上 redelivered 标记(同一份文档)。这就是至少一次(at-least-once)传递,反过来说,已经处理完毕、却在 ack 之前死掉的消息,会再来一次。队列不会提供幂等。必须由消费者用 message_id 留下处理记录来过滤(与第 8 模块的账本是同样的思路)。

拒绝有两种。资料不全(业务拒绝)即使再给一百次,结果也一样。如果把它退回(requeue),同一条消息就会在队列里无限打转——通常称为毒消息(poison message)。对 basic.reject(或 nack)指定 requeue=False,Broker 就会丢弃消息,或者,如果队列上指定了死信交换机(DLX),就重新发布到那里。用队列参数 x-dead-letter-exchange·x-dead-letter-routing-key 来指定,重新发布的消息会在 x-death 头中留下原因(rejected·expired·maxlen·delivery_limit)(DLX)。而对于临时错误(理赔系统 503),稍后再试就行。本实验中,先退回一次,如果是 redelivered 却又失败,就发往 DLQ。(文档建议在生产环境中不要用参数,而要用策略(policy)来设置 DLX——因为这样无需重新部署就可以修改。)

prefetch 就是背压。消费者通道的 basic.qos(prefetch_count) 是未 ack 状态下最多可以扛下的条数。如果是 0(无限制),Broker 会把队列里的消息全部推给消费者。理赔系统一慢,消费者内存里就会堆积数千条,即使再启动一个消费者,也因为已经被全部拿走而没有可分的了。设了上限,其余的就会留在队列里,由新的消费者分担。

Broker 也会保护自己。内存超过告警阈值时,Broker 会阻塞所有发布的连接,等消费推进、内存下降后再放开(内存告警)。阈值的默认值是相对于检测到的内存的比例(按当前文档为 0.6),而文档警告说,在容器中 Broker 并不总能获知 cgroup 限制,并建议使用绝对值。实际上,在这个过程中测量的结果,在 8GB 机器上的 2Gi 容器里,Ubuntu 软件包的 3.12 Broker 把阈值定在了 3.3GB——意思是告警还没响,容器就已经被 OOM 杀死了。所以本实验的辅助脚本用 vm_memory_high_watermark.absolute = 512MiB 来启动。

在现场相遇的样子

最常见的事故,是用自动 ack 写的消费者在部署中重启,把正在处理的消息丢掉了。日志里什么都没有。第二种是毒消息——一条格式错误的消息被无限重新投递,吃掉消费者的 CPU,排在它后面的正常消息被拖延好几个小时。第三种是抱着“放进 MQ 就安全了”的信念而不做发布确认。Broker 因内存告警正在阻塞发布,而发布者只看超时就重试,这期间放进去了什么,没有人知道。第四种是建了 DLQ 却没有人看——DLQ 必须与监控和重处理流程成对存在。

下一项实验要做什么

在 Pod 内启动 RabbitMQ(辅助脚本),声明交换机、队列和 DLX 拓扑,并制作能收到发布确认的发布者。然后依次扩展消费者——手动 ack、业务拒绝走 DLQ、临时错误重新投递一次、prefetch、用 message_id 过滤重复。评分器为了不碰学员的队列,会创建一个临时 vhost,并在那里运行你的脚本(所以所有脚本都读取 EAI_AMQP_URL)。