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

构建 EAI 中间层

理赔走队列 — 不丢失,不重复处理

在 TT Lab 中继续学习

目标

用 RabbitMQ 异步传递保险理赔。声明拓扑(交换机、队列、DLX),并制作能收到发布确认的发布者,以及具备手动 ack、DLQ、重新投递、prefetch 和去重的消费者。

为什么重要

队列承诺“不丢失”,但这个承诺只有在做了发布确认和手动 ack 时才成立,而作为代价,同一条消息会到达两次。如果把该拒绝的消息退回,就会变成毒消息,没有 prefetch,消费者就会在缓慢的目标前全部扛下。异步对接的事故,大部分都出在这五点上。

步骤

  1. 用 bash /opt/lab/fixtures/eaimw/mq/mq-up.sh 启动 Broker,并用 rabbitmqctl -n rabbit@localhost status > /root/eaimw/mq/status.txt 留下状态。确认内存告警阈值是绝对值 512MiB(0.5369 gb)。
  2. /root/eaimw/mq/topology.py:交换机 eai.claim(direct,durable)、eai.dlx(direct,durable),队列 claim.in(durable,参数 x-dead-letter-exchange=eai.dlx、x-dead-letter-routing-key=claim.dead)、claim.dead(durable),绑定 eai.claim→claim.in(键 claim)、eai.dlx→claim.dead(键 claim.dead)。运行两次也不能出错。
  3. /root/eaimw/mq/publish.py <청구JSON파일>(占位符为理赔 JSON 文件):把 JSON 原样作为正文,以路由键 claim 发布到 eai.claim。persistent(delivery_mode 2),message_id 取 JSON 中的 guid,content_type 为 application/json。开启发布确认并以 mandatory 发送,无法路由或被拒绝时,以非 0 的退出码结束。
  4. /root/eaimw/mq/consumer.py --claim <청구시스템URL> --max-seconds <초>(占位符依次为理赔系统 URL、秒数):消费 claim.in,把正文通过 POST /v1/claims 传出,如果是 201,这时才 ack(手动 ack)。其他结果先退回(nack,requeue)。超过 --max-seconds 就结束。
  5. 业务拒绝(422,资料不全)不要退回,而是用 reject(requeue=False) 发往 DLX。
  6. 临时错误(503、连接失败、超时)退回一次,如果被重新投递的(redelivered)消息又失败,就发往 DLX。
  7. 用 --prefetch(默认 5)设置 basic_qos,使即使在缓慢的理赔系统前,未 ack 的消息也不超过 5 条。
  8. 把处理记录留在 SQLite(EAI_MQ_DB,默认 /root/eaimw/mq/processed.db)中,对于已经处理过的 message_id,不调用理赔系统,只做 ack(即使重新启动消费者也能记住)。

参考

在 Pod 内启动 Broker

用 mq-up.sh 启动 RabbitMQ,并把 rabbitmqctl status 的输出保存到 /root/eaimw/mq/status.txt。

辅助脚本把节点名称设为 rabbit@localhost 来启动,所以要给 rabbitmqctl 加上 -n rabbit@localhost。请看输出中的 Memory high watermark 这一行。

声明交换机、队列和 DLX

让 /root/eaimw/mq/topology.py 声明 2 个交换机、2 个队列、2 个绑定,以及 claim.in 的 DLX 参数(运行两次也安全)。

声明是幂等的——用相同的属性再次声明不会有任何事,属性不同,Broker 就会以 406 PRECONDITION_FAILED 关闭通道。DLX 是两个队列参数。

能收到确认的发布者

让 /root/eaimw/mq/publish.py 携带 persistent、message_id、content_type 发布,并通过发布确认和 mandatory 发现无法路由。

先调用 ch.confirm_delivery(),basic_publish 就会等待 Broker 的确认。用 mandatory=True 发送,当没有可接收的队列时,就会抛出 UnroutableError。

处理之后才 ack

让 /root/eaimw/mq/consumer.py 消费 claim.in 并传给理赔系统,只有在收到 201 之后才 ack(其他情况退回)。

用 ch.consume(queue, inactivity_timeout=0.5) 循环,空闲期间会收到 None。超过 --max-seconds 就退出并关闭连接——没有 ack 的,Broker 会退回。不要使用 auto_ack。

业务拒绝走 DLQ

对 422(资料不全)不要退回,而是用 reject(requeue=False) 发往 DLX。

被退回的消息会立刻再次到来。如果把给一百次结果也一样的消息退回,消费者就会一直围着它转。如果 requeue=False,就会去往队列上设置的 DLX。

临时错误再来一次,还不行就走 DLQ

503、连接失败退回一次,如果是 redelivered 却又失败,就发往 DLX。

method 帧的 redelivered 为 True,就说明是已经被退回过一次的消息。如果无限允许退回,当目标长时间挂掉时,整个队列就会只围着这些消息转。

用 prefetch 限制扛下的量

用 --prefetch(默认 5)设置 basic_qos,使 unacked 不超过 5 条。

ch.basic_qos(prefetch_count=N) 以通道为单位,限制未 ack 的投递数。请在开始消费之前调用。

同一个 message_id 只处理一次

把已处理的 message_id 留在 SQLite(EAI_MQ_DB)中,再次到来时不调用,只做 ack。

处理成功后就记录并 ack。即使在记录与 ack 之间死掉,下次到来的同一条消息也会因为有记录而被过滤。记录必须在重新启动消费者之后仍然保留。