理赔走队列 — 不丢失,不重复处理
目标
用 RabbitMQ 异步传递保险理赔。声明拓扑(交换机、队列、DLX),并制作能收到发布确认的发布者,以及具备手动 ack、DLQ、重新投递、prefetch 和去重的消费者。
为什么重要
队列承诺“不丢失”,但这个承诺只有在做了发布确认和手动 ack 时才成立,而作为代价,同一条消息会到达两次。如果把该拒绝的消息退回,就会变成毒消息,没有 prefetch,消费者就会在缓慢的目标前全部扛下。异步对接的事故,大部分都出在这五点上。
步骤
- 用
bash /opt/lab/fixtures/eaimw/mq/mq-up.sh启动 Broker,并用rabbitmqctl -n rabbit@localhost status > /root/eaimw/mq/status.txt留下状态。确认内存告警阈值是绝对值 512MiB(0.5369 gb)。 /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)。运行两次也不能出错。/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 的退出码结束。/root/eaimw/mq/consumer.py --claim <청구시스템URL> --max-seconds <초>(占位符依次为理赔系统 URL、秒数):消费claim.in,把正文通过POST /v1/claims传出,如果是 201,这时才 ack(手动 ack)。其他结果先退回(nack,requeue)。超过--max-seconds就结束。- 业务拒绝(422,资料不全)不要退回,而是用
reject(requeue=False)发往 DLX。 - 临时错误(503、连接失败、超时)退回一次,如果被重新投递的(
redelivered)消息又失败,就发往 DLX。 - 用
--prefetch(默认 5)设置basic_qos,使即使在缓慢的理赔系统前,未 ack 的消息也不超过 5 条。 - 把处理记录留在 SQLite(
EAI_MQ_DB,默认/root/eaimw/mq/processed.db)中,对于已经处理过的message_id,不调用理赔系统,只做 ack(即使重新启动消费者也能记住)。
参考
- 所有脚本都读取
EAI_AMQP_URL(默认amqp://guest:guest@127.0.0.1:5672/%2F):pika.BlockingConnection(pika.URLParameters(URL))。评分器会传入临时 vhost 的地址。guest 账户默认只能从回环地址连接。 - 理赔系统夹具:
nohup python3 /opt/lab/fixtures/eaimw/partner.py claim > /root/eaimw/mq/claim.out 2>&1 &(9202)。docs为空的理赔返回 422,claimId以BUSY开头则返回 503,统计看/_stats的by_guid(按 claimId 的调用数)。 - 查看队列:
rabbitmqctl -n rabbit@localhost list_queues name messages messages_unacknowledged。 - pika:
ch.confirm_delivery()之后,basic_publish(..., mandatory=True)遇到无法路由就会抛出pika.exceptions.UnroutableError。消费用for m, props, body in ch.consume("claim.in", inactivity_timeout=0.5):(空闲期间m为 None)。 - 常见错误:
auto_ack=True;对业务拒绝做 requeue;只创建 durable 队列而不用 persistent。
在 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 之间死掉,下次到来的同一条消息也会因为有记录而被过滤。记录必须在重新启动消费者之后仍然保留。