Building an EAI Middleware Layer
Claims Through a Queue — Nothing Lost, Nothing Processed Twice
Goal
Hand insurance claims over asynchronously with RabbitMQ. Declare the topology (exchange, queue, DLX), and build a publisher that receives confirms and a consumer with manual ack, a DLQ, redelivery, prefetch and duplicate filtering.
Why it matters
A queue promises "we won't lose it," but that promise holds only when you do publisher confirms and manual ack, and in return the same message comes twice. If you return a message that should be rejected, it becomes a poison message, and without prefetch the consumer takes everything on in front of a slow target. Most asynchronous integration incidents come from these five things.
Steps
- Start the broker with
bash /opt/lab/fixtures/eaimw/mq/mq-up.shand save the status withrabbitmqctl -n rabbit@localhost status > /root/eaimw/mq/status.txt. Check that the memory alarm threshold is the absolute value 512MiB (0.5369 gb). /root/eaimw/mq/topology.py: exchangeseai.claim(direct, durable) andeai.dlx(direct, durable); queuesclaim.in(durable, with argumentsx-dead-letter-exchange=eai.dlxandx-dead-letter-routing-key=claim.dead) andclaim.dead(durable); bindingseai.claim→claim.in(keyclaim) andeai.dlx→claim.dead(keyclaim.dead). It must not error even when run twice./root/eaimw/mq/publish.py <청구JSON파일>(claim JSON file): publish the JSON as the body as it is, toeai.claimwith the routing keyclaim. Persistent (delivery_mode 2),message_idis the JSON'sguid, andcontent_typeisapplication/json. Turn on publisher confirms and send withmandatory, so that it ends with a non-zero code on unroutable or rejected./root/eaimw/mq/consumer.py --claim <청구시스템URL> --max-seconds <초>(claims system URL, seconds): consumeclaim.in, pass the body on withPOST /v1/claims, and ack only then if it is 201 (manual ack). For any other result, return it for now (nack, requeue). It ends when--max-secondshas passed.- Do not return business rejections (422, missing documents); send them to the DLX with
reject(requeue=False). - Return a transient error (503, connection failure, timeout) once, and if a message that was redelivered (
redelivered) fails again, send it to the DLX. - With
--prefetch(default 5), applybasic_qos, so that even in front of a slow claims system, unacknowledged messages do not exceed 5. - Leave a processing record in SQLite (
EAI_MQ_DB, default/root/eaimw/mq/processed.db), so that for amessage_idalready processed, you only ack without calling the claims system (it remembers even if you restart the consumer).
Notes
- Every script reads
EAI_AMQP_URL(defaultamqp://guest:guest@127.0.0.1:5672/%2F):pika.BlockingConnection(pika.URLParameters(URL)). The grader passes the address of a temporary vhost. The guest account can by default connect only from loopback. - Claims system fixture:
nohup python3 /opt/lab/fixtures/eaimw/partner.py claim > /root/eaimw/mq/claim.out 2>&1 &(9202). A claim with emptydocsgets 422, aclaimIdstarting withBUSYgets 503, and statistics are in/_stats, underby_guid(number of calls per claimId). - Looking into queues:
rabbitmqctl -n rabbit@localhost list_queues name messages messages_unacknowledged. - pika: after
ch.confirm_delivery(),basic_publish(..., mandatory=True)raisespika.exceptions.UnroutableErrorwhen unroutable. To consume,for m, props, body in ch.consume("claim.in", inactivity_timeout=0.5):(during idle timemis None). - Common mistakes:
auto_ack=True, requeuing a business rejection, creating only a durable queue without persistent.
Start the broker inside the Pod
Start RabbitMQ with mq-up.sh and save the output of rabbitmqctl status to /root/eaimw/mq/status.txt.
The helper starts the node under the name rabbit@localhost, so give rabbitmqctl -n rabbit@localhost. Look at the Memory high watermark line in the output.
Declare the exchanges, queues and DLX
/root/eaimw/mq/topology.py declares 2 exchanges, 2 queues, 2 bindings and claim.in's DLX arguments (safe to run twice).
Declarations are idempotent — declaring again with the same properties does nothing, and if the properties differ, the broker closes the channel with 406 PRECONDITION_FAILED. The DLX is two queue arguments.
A publisher that receives confirms
/root/eaimw/mq/publish.py publishes with persistent, message_id and content_type attached, and finds out about unroutable messages with publisher confirms and mandatory.
If you call ch.confirm_delivery() first, basic_publish waits for the broker's confirmation. If you send with mandatory=True, an UnroutableError occurs when there is no queue to receive it.
Ack only after processing
/root/eaimw/mq/consumer.py consumes claim.in, passes it to the claims system, and acks only after receiving 201 (anything else is returned).
If you loop with ch.consume(queue, inactivity_timeout=0.5), None comes during idle time. When --max-seconds has passed, exit and close the connection — the broker returns what was not acked. Do not use auto_ack.
Business rejections go to the DLQ
Do not return a 422 (missing documents); send it to the DLX with reject(requeue=False).
A returned message comes back right away. If you return a message that gives the same result a hundred times over, the consumer loops holding on to only that. With requeue=False, it goes to the DLX attached to the queue.
Transient errors once more, then the DLQ if it still fails
Return a 503 or connection failure once, and if it is redelivered and fails again, send it to the DLX.
If the method frame's redelivered is True, the message has already been returned once. If you allow returning without limit, when the target is dead for a long time the whole queue circles only those messages.
Limit how much is taken on with prefetch
Apply basic_qos with --prefetch (default 5) so that unacked does not exceed 5.
ch.basic_qos(prefetch_count=N) limits the number of unacknowledged deliveries per channel. Call it before starting to consume.
The same message_id only once
Leave the processed message_id in SQLite (EAI_MQ_DB), and when it comes again, do not call; only ack.
On success, record and ack. Even if it dies between the record and the ack, the same message that comes next is filtered out thanks to the record. The record must remain even if you restart the consumer.