System Design Interview

Understand the Problem

Scope 4 min readLesson 1 of 4

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.

Figure 1 Popular distributed message queues
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.

Clarifying the scope
Candidate

What’s the format and average size of messages? Is it text only? Is multimedia allowed?

Interviewer

Text messages only. Messages are generally measured in the range of kilobytes (KBs).

Candidate

Can messages be repeatedly consumed?

Interviewer

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.

Candidate

Are messages consumed in the same order they were produced?

Interviewer

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.

Candidate

Does data need to be persisted and what is the data retention?

Interviewer

Yes, let’s assume data retention is two weeks. This is an added feature. A traditional distributed message queue does not retain messages.

Candidate

How many producers and consumers are we going to support?

Interviewer

The more the better.

Candidate

What’s the data delivery semantic we need to support? For example, at-most-once, at-least-once, and exactly once.

Interviewer

We definitely want to support at-least-once. Ideally, we should support all of them and make them configurable.

Candidate

What’s the target throughput and end-to-end latency?

Interviewer

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.