When not to use FIFO messages queues
Queues
If you are relatively new to system design, software architecture or whatever you want to name the discipline of deciding which components you system be composed of, definitely you have seen that messages queues are a key component of a diverse of systems. From simple single producer single consumer architectures where a process put messages in a queue and (maybe?) another process read from it to architectures where messages are generated from different origins, some of them even generated by the system itself, and being consumed by multiple consumers maybe because you want to speed up the processing time. And in any of aforementioned scenarios, you must decide if that message queue will be an standard queue or a FIFO like queue.
Standard queue
Standard queues are the most widely used queue system because its write high throughput and basically no complex setup to work in production in no time. These type of queues allows to write a lot of messages because it doesn’t give the user guarantees about the messages it hold. Why? given its distributed nature, it skips a few mechanisms to ensure a message is delivered only once at a consumer but jey, at least it tries.
To be fair, the amount of companies likely to face drawbacks of this type of messages queue are not that many so it is worth of using it without giving a thought for more than a few hours. or minutes… But it is always nice to remember that you might face some distributed problem issues. The most important one: process a message twice or more. So design your system to be idempotent.
FIFO queue
FIFO (First-in First-out) message queues are the queues teached in computer science foundations. A first item arrives at the queue and no matter how many items more arrive, the first item to be retrieve will be always the one that arrives first. So far so good. But what happens when providers implement FIFO queues in distributed systems is that some guarantees about FIFO can’t be achieve while not sacrificing a bit of write and read throughput.
FIFO messages queues will ensure that the message will be delivered exactly once at a consumer. In fact, the message will be delivered once and to only one consumer, so naturally this takes a little more time per message than standard queues. If we are using message queues, latency are not likely the first problem we are facing because we have a background process consuming messages and taking action based in the content of the message.
Usual FIFO queues would allow one consumer in order to guarantee only once delivery but this seems kind of inappropiate in a distributed system. Why would we like a single consumer in a high performance environment? Odd. A very good proposal is separate, logically, messages based in an identifier that allows a different consumer to read messages from the same queue. This identifier is known as a message group ID and let’s group messages to be processed by the same consumer and, therefore, any other messages be processed by another consumer.
Why would we want to have FIFO queues if with standard queues + dedup / idempotency mechanism we would achieve the same goal then?
It is not up to me to tell you why you should use one over another but the typical use case for FIFO queues is critical transactions that cannot, under any circumstance, be processed twice. Even idempotency mechanism can fail we race conditions enter the room. Better be safe and slow than attend after-work meeting explaining why we are charging our customers twice randomly with no clear explanation.
So, if we want to implement a payment processor system, each event / message generated by the same customer can be assigned the same message group and be processed by the same consumer, one by one, ensuring that they are consumed in order.
Humbling experience
In my company, we designed a system receives recording phone calls and emit events to a FIFO queue with the path of such recording phone call and start an analytics pipeline for that audio. We did some tests on a development environment sending about 40k in recordings and observing the amount of time that took to process such amount. Seems OK there.
Now, time to deploy into production. Boom. First hour in production, the amount of recording calls where higher than expected but our autoscaling policy on the queue kicks in and start scaling out more instances to process the messages on the queue but then disaster occured. Our system was unable to decrease the amount of messages. In fact, the queue did nothing more than increasing the amount of messages with sign of decreasing. What is happening? we asked.
All of the messages from the queue were arriving at the same message group, so only one consumer were able to read them because, when a consumer reads / polls for messages, it locks effectively the group, making impossible the other consumers to read from the queue. Clearly it did not matter how many VM were deployed and ready to consume the queue, none of them were able to do a thing.
What could be the solution? Sending messages to random message groups was a naive idea but might work. No, changing the type of queue from FIFO to Standard was a better, and the chosed, idea.
Since I started working there, the usage of FIFO queues were the to-go option because it allows to avoid idempotency, the messages were processed only once. What a paradise, huh? But standard queues were the clear best fit for our use case.
The moral of the story
Conclusion
Study the technology you use but most importantly, study your business case and verify that you have checked if each component of your proposed system supports the goal.