TT Lab
Get started
Learn Learning paths Courses

Building an EAI Middleware Layer

Queues Don't Lose Messages, but They Deliver Twice

Continue in TT Lab

In one line

Asynchronous integration promises not "process this now" but "process this later, without losing it." A queue is the device that keeps that promise, but for the promise to hold, the publishing side must get confirmation that the broker received the message (publisher confirms), and the consuming side must acknowledge only after processing (manual ack) — and in return, the same message can arrive twice.

Why it was needed

Consider accepting an insurance claim. When a customer uploads a claim in the app, the channel wants to show "Received" immediately. But the claims system takes a few seconds per case for document verification, and at month end it backs up by hours. If you call it synchronously, the channel is bound to the claims system's speed, and if the claims system goes down briefly, even accepting the claim fails. With a queue in between, the channel can answer the moment it puts it on the queue, and the claims system takes items out at its own pace. Even if one side stops, the other keeps working.

In exchange, new questions arise. How do you know it was put in? If the side that took it out dies mid-processing, where does the message go? Does a message that cannot be processed (missing documents) circulate in the queue forever? If the claims system is slow, how many items must a consumer take on? This module answers these questions one by one.

How it works

The three parts of AMQP 0-9-1. A publisher sends not to a queue but to an exchange. The exchange puts the message on queues by looking at the routing key and the bindings. A direct exchange sends to the queues whose binding key exactly equals the routing key (AMQP concepts). Thanks to this one step that keeps the publisher from knowing queue names, even if you later bind one more queue for auditing, the publisher has nothing to fix.

Durable and persistent are different. A durable queue keeps its queue definition even if the broker restarts. For messages to survive, you must publish them as persistent (delivery_mode=2). If you do only one of the two, after a restart the queue exists but is empty.

Publisher confirms. Just because basic_publish returned without error does not mean the broker received it. If you put the channel in confirm mode, the broker returns an ack for each message, and according to the documentation, a persistent message going to a durable queue is confirmed after it is written to disk (Confirms). A message with nowhere to be routed is silently dropped, but if you send it with mandatory, the broker returns it with basic.return before the ack. pika's BlockingChannel reports this as UnroutableError.

Manual ack. In auto-ack mode, the broker considers delivery finished the moment it sends the message, so if the consumer dies midway through processing, that message is lost. In manual ack, it is finished only when the consumer finishes processing and sends an ack. If the channel closes without an ack, the broker automatically returns that message and attaches a redelivered mark when it gives it out again (same document). This is at-least-once delivery, and to put it the other way around, a message that was processed but whose consumer died just before the ack comes once more. The queue does not give idempotency. The consumer has to leave a processing record by message_id and filter duplicates (the same idea as module 8's ledger).

There are two kinds of rejection. A missing-document case (business rejection) gives the same result a hundred times over. If you return it (requeue), the same message circulates in the queue indefinitely — commonly called a poison message. If you call basic.reject (or nack) with requeue=False, the broker drops the message or, if the queue has a dead-letter exchange (DLX) specified, republishes it there. You specify it with the queue arguments x-dead-letter-exchange and x-dead-letter-routing-key, and the republished message carries the reason in the x-death header (rejected, expired, maxlen, delivery_limit) (DLX). A transient error (claims system 503), on the other hand, can simply be retried a little later. In this lab you return it once, and if a redelivered message fails again, send it to the DLQ. (The documentation recommends setting the DLX in production with a policy rather than with arguments — because you can change it without redeploying.)

Prefetch is backpressure. A consumer channel's basic.qos(prefetch_count) is the maximum number of unacknowledged messages it may take on. At 0 (unlimited), the broker pushes all the messages in the queue to the consumer. If the claims system is slow, thousands pile up in the consumer's memory, and even if you start another consumer, it has already taken everything so there is nothing to share. If you set a limit, the rest stay in the queue and new consumers share them.

The broker protects itself too. When memory exceeds the alarm threshold, the broker blocks all publishing connections, and releases them when consumption proceeds and memory falls (memory alarms). The default threshold is a ratio of the detected memory (0.6 in the current documentation), but the documentation warns that in containers the broker cannot always find out the cgroup limit and recommends an absolute value. In fact, measured in the course of this work, inside a 2Gi container on an 8GB machine, the Ubuntu package's 3.12 broker set the threshold at 3.3GB — meaning the container dies of OOM before the alarm sounds. That is why this lab's helper starts it with vm_memory_high_watermark.absolute = 512MiB.

What it looks like in the field

The most common incident is a consumer written with auto ack that loses the message it was processing when it restarts during a deployment. There is nothing in the log. The second is the poison message — one message with a wrong format is redelivered endlessly, eating consumer CPU, and the normal messages behind it are delayed by hours. The third is not doing publisher confirms because of the belief that "we put it in MQ, so it's safe." The broker is blocking publishing because of a memory alarm, the publisher retries looking only at the timeout, and nobody knows what went in in the meantime. The fourth is creating a DLQ and nobody looking at it — a DLQ must come paired with monitoring and a reprocessing procedure.

What we do in the next lab

You start RabbitMQ inside the Pod (with a helper), declare the exchange, queue and DLX topology, and build a publisher that receives confirms. Then you grow the consumer in stages — manual ack, the DLQ for business rejections, a single redelivery for transient errors, prefetch, and filtering duplicates by message_id. The grader creates a temporary vhost so as not to touch the student's queues and runs your scripts there (that is why every script reads EAI_AMQP_URL).