DEV Community

Dilip V P
Dilip V P

Posted on

Queues Deliver at Least Once, So Every Consumer Must Be Idempotent

A food delivery app's SMS service has been down for three hours. You place an order anyway, and it goes through. In fact, not one order has failed in those three hours.

The messages for the SMS service are waiting in a queue. That delay is exactly what the system was designed to do.

This article builds the same app twice, first the simple way and then with a message queue. Then it works through the two problems that queues bring with them. The same message can arrive twice, and some messages fail forever.

The simple version: six steps in a row

You tap Place Order, and the order service does six steps, one after another:

  1. check the payment
  2. save the order
  3. tell the restaurant
  4. find a rider
  5. send an SMS
  6. add reward points

Each step is a synchronous call. The order service sends a request and waits for the answer before it moves on to the next step.

Now the SMS service goes down. The SMS call fails, so the whole order fails, even though the payment worked and a rider was already assigned. This is called tight coupling. One service cannot finish unless another service is working at the same moment.

Here is how often that happens. Say each of the six services is up 99.9 percent of the time. The order works only when all six work at once. In a simple model where they fail independently, the chain is up 0.999 to the sixth power, which is about 99.4 percent. Over a 30 day month, that is about 259 minutes of failed orders, more than four hours. One service on its own at 99.9 percent is down for about 43 minutes a month.

The critical path rule

The fix starts with one rule.

A step belongs on the waiting path only if the customer needs its result before the screen says Order placed.

That waiting path is called the critical path. Only two of the six steps pass the test. The payment must succeed, and the order must be saved. The customer does not need the SMS, the reward points, the restaurant or even the rider before the screen confirms the order.

Save now, send later

So the order service does only those two steps. Then it writes one short message that says order 4472 was placed, and the phone shows Order placed right away.

That message cannot wait in the order service's memory, because a restart would wipe it. It goes to a separate system that stores messages safely until they are read. That system is called a message broker. RabbitMQ, Amazon SQS and Kafka are common examples.

The broker puts a copy of the message into four queues, one for each later step. The order service is the producer. The services that read the queues are the consumers. A consumer is not the customer. It is a program running on a server, like the rider service.

Now replay the outage. The SMS service goes down, but the order service never notices, because it no longer calls the SMS service. Orders keep going through, and the messages pile up in the SMS queue. When the SMS service comes back, it works through the backlog. Three hours of downtime cost zero failed orders.

A queue also absorbs bursts. At dinner time, orders can arrive faster than the SMS service can send them. The queue grows, then drains when traffic falls. But if messages arrive faster than they are processed all day, the queue only keeps growing. A queue can hold a temporary excess. It cannot fix a consumer that is too slow all the time.

The same message, twice

Follow one message into the rider service. The consumer reads the message and assigns a rider named Alex. Then it tells the broker it is done. That signal is called an acknowledgement, or ack, and after the ack the broker deletes the message.

If no ack arrives, the broker assumes the consumer failed and delivers the message again. In Amazon SQS, the broker waits for a time limit called the visibility timeout, and the default is 30 seconds.

Now watch what one crash does:

1. rider service A reads "order 4472 was placed"
2. rider service A saves Alex as the rider of order 4472
3. rider service A crashes before it sends the ack
4. 30 seconds pass with no ack, so the broker delivers the message again
5. rider service B reads the same message and assigns Sam
Enter fullscreen mode Exit fullscreen mode

Two riders now drive to the same restaurant, even though each step did what its code says.

This is called at-least-once delivery. A crash does not lose the message, but the message can arrive more than once. The SQS documentation says it directly: for standard queues, "more than one copy of a message might be delivered".

The other choice is at-most-once delivery. The consumer acks before it does the work. A crash then cannot cause a duplicate, but it can leave the order with no rider at all.

Exactly-once delivery sounds better, but a consumer cannot promise it. Saving the rider and sending the ack happen in two different systems, the database and the broker, and a crash can always land between them. Kafka offers exactly-once guarantees, but its transactions cover only data that stays inside Kafka. The moment the consumer writes to its own database, there are two systems again.

The fix lives in the consumer

In the database, the consumer runs one conditional update. It sets the rider only if the order has no rider yet:

UPDATE orders
SET rider = 'Alex'
WHERE id = 4472 AND rider IS NULL;
Enter fullscreen mode Exit fullscreen mode

The first delivery updates one row. A second delivery, from any instance, updates zero rows, so the consumer skips the work and sends the ack.

Doing the work twice now has the same effect as doing it once. That property is called idempotency, and it gives the rule for every queue you will use:

Queues deliver at least once, so every consumer must be idempotent.

Deleted or kept

Brokers come in two main designs. After a consumer reads a message, the message is either deleted or kept.

In a queue, like SQS or a standard RabbitMQ queue, the message is deleted after the ack. If you run three instances of the rider service, they share one queue, and each message normally goes to just one of them. That pattern is called competing consumers.

In a log, like Kafka, each message is added to the end and kept after it is read, for a retention period. Kafka's default is seven days (168 hours). Each reading service remembers its own position in the log, and that position is called an offset.

The difference shows up when something new arrives. Next month, you add a fraud check. On a log, the fraud check can start from the oldest stored message and read the last seven days of orders, because nothing was deleted. A new queue would only see orders from now on. Reading old messages again is called replay.

So the choice comes down to two sentences:

  • Use a queue when each message is a task for one worker.
  • Use a log when many services need the same events, or when you may need to replay them.

The message that fails forever

Now look at the last problem. An order arrives with a phone number that has only five digits. The SMS provider rejects it, so the SMS service sends no ack, and the message comes back again and again.

A message that can never succeed is called a poison message. Every retry wastes work. In a queue that keeps strict order, it also blocks every message behind it.

So split failures into two kinds:

  • A temporary failure, like a timeout from the provider, may work on the next try.
  • A permanent failure, like a bad phone number, fails every time.

Temporary failures get retries with a longer wait each time, plus a small random delay:

attempt 1 fails, then wait about 1 second
attempt 2 fails, then wait about 2 seconds
attempt 3 fails, then wait about 4 seconds
Enter fullscreen mode Exit fullscreen mode

That is exponential backoff, and the random part is called jitter. The jitter matters. Without it, every client that failed at the same moment also retries at the same moment, and the retries themselves overload the service again. That pattern is called a retry storm.

Permanent failures need an attempt limit. After a few failed deliveries, the broker moves the message to a separate queue called a dead letter queue.

In SQS, you set the limit with maxReceiveCount in the queue's redrive policy, which is the setting that connects a queue to its dead letter queue. In RabbitMQ, quorum queues (its replicated queue type) have a delivery limit. Since RabbitMQ 4.0, the default is 20. Past the limit, RabbitMQ drops the message, unless you configure a dead letter exchange. A dead letter exchange is where RabbitMQ sends failed messages, and it routes them into a dead letter queue. So on RabbitMQ, configure one, or the bad message is lost.

With a dead letter queue in place, the main queue keeps moving, and a person can inspect the bad message, fix the data, and send the message back.

There is one more warning. A dead letter queue that nobody watches is just a slower way to lose messages, so put an alert on it.

When a direct call is still right

None of this means that every call should go through a queue. A direct call is still right when the caller cannot continue without the reply.

The payment check stays synchronous, because the customer is waiting to know whether they paid. Even that call sends an idempotency key, which is a unique ID for this one payment attempt. If the call is retried, the payment service sees the same key and does not charge the customer twice.

Three lines worth keeping

  1. Put a step on the critical path only if the customer needs its result right now.
  2. Queues deliver at least once, so make every consumer idempotent.
  3. Retry temporary failures with backoff, and send permanent failures to a dead letter queue.

Every consumer in this article was a program we wrote. In the next episode, an AI model decides which tools to call. It calls five of them, and only three answer.


The video builds the whole app on screen, and shows the outage, the double delivery and the dead letter queue step by step: https://youtu.be/sLze66ezsHY

Top comments (0)