Understand the Problem
In this chapter, we explore a popular question in system design interviews: design a distributed message queue. In modern architecture, systems are broken up into small and independent building blocks with well-defined interfaces between them. Message queues provide communication and coordination for those building blocks. What benefits do message queues bring?
-
Decoupling. Message queues eliminate the tight coupling between components so they can be updated independently.
-
Improved scalability. We can scale producers and consumers independently based on traffic load. For example, during peak hours, more consumers can be added to handle the increased traffic.
-
Increased availability. If one part of the system goes offline, the other components can continue to interact with the queue.
-
Better performance. Message queues make asynchronous communication easy. Producers can add messages to a queue without waiting for the response and consumers consume messages whenever they are available. They don’t need to wait for each other.
Figure 1 shows some of the most popular distributed message queues on the market.
Message queues vs event streaming platforms
Strictly speaking, Apache Kafka and Pulsar are not message queues as they are event streaming platforms. However, there is a convergence of features that starts to blur the distinction between message queues (RocketMQ, ActiveMQ, RabbitMQ, ZeroMQ, etc.) and event streaming platforms (Kafka, Pulsar). For example, RabbitMQ, which is a typical message queue, added an optional streaming feature to allow repeated message consumption and long message retention, and its implementation uses an append-only log, much like an event streaming platform would. Apache Pulsar is primarily a Kafka competitor, but it is also flexible and performant enough to be used as a typical distributed message queue.
In this chapter, we will design a distributed message queue with additional features, such as long data retention, repeated consumption of messages, etc., that are typically only available on event streaming platforms. These additional features make the design more complicated. Throughout the chapter, we will highlight places where the design could be simplified if the focus of your interview centers around the more traditional distributed message queues.
Step 1 – Understand the problem and establish design scope
In a nutshell, the basic functionality of a message queue is straightforward: producers send messages to a queue, and consumers consume messages from it. Beyond this basic functionality, there are other considerations including performance, message delivery semantics, data retention, etc. The following set of questions will help clarify requirements and narrow down the scope.
What’s the format and average size of messages? Is it text only? Is multimedia allowed?
Text messages only. Messages are generally measured in the range of kilobytes (KBs).
Can messages be repeatedly consumed?
Yes, messages can be repeatedly consumed by different consumers. Note that this is an added feature. A traditional distributed message queue does not retain a message once it has been successfully delivered to a consumer. Therefore, a message cannot be repeatedly consumed in a traditional message queue.
Are messages consumed in the same order they were produced?
Yes, messages should be consumed in the same order they were produced. Note that this is an added feature. A traditional distributed message queue does not usually guarantee delivery orders.
Does data need to be persisted and what is the data retention?
Yes, let’s assume data retention is two weeks. This is an added feature. A traditional distributed message queue does not retain messages.
How many producers and consumers are we going to support?
The more the better.
What’s the data delivery semantic we need to support? For example, at-most-once, at-least-once, and exactly once.
We definitely want to support at-least-once. Ideally, we should support all of them and make them configurable.
What’s the target throughput and end-to-end latency?
It should support high throughput for use cases like log aggregation. It should also support low latency delivery for more traditional message queue use cases.
With the above conversation, let’s assume we have the following functional requirements:
-
Producers send messages to a message queue.
-
Consumers consume messages from a message queue.
-
Messages can be consumed repeatedly or only once.
-
Historical data can be truncated.
-
Message size is in the kilobyte range.
-
Ability to deliver messages to consumers in the order they were added to the queue.
-
Data delivery semantics (at-least-once, at-most-once, or exactly-once) can be configured by users.
Non-functional requirements
-
High throughput or low latency, configurable based on use cases.
-
Scalable. The system should be distributed in nature. It should be able to support a sudden surge in message volume.
-
Persistent and durable. Data should be persisted on disk and replicated across multiple nodes.
Adjustments for traditional message queues
Traditional message queues like RabbitMQ do not have as strong a retention requirement as event streaming platforms. Traditional queues retain messages in memory just long enough for them to be consumed. They provide on-disk overflow capacity 1 which is several orders of magnitude smaller than the capacity required for event streaming platforms. Traditional message queues do not typically maintain message ordering. The messages can be consumed in a different order than they were produced. These differences greatly simplify the design which we will discuss where appropriate.
Finished reading?
Mark it complete to track your progress.