High-Level Design
First, let’s discuss the basic functionalities of a message queue.
Figure 2 shows the key components of a message queue and the simplified interactions between these components.
-
Producer sends messages to a message queue.
-
Consumer subscribes to a queue and consumes the subscribed messages.
-
Message queue is a service in the middle that decouples the producers from the consumers, allowing each of them to operate and scale independently.
-
Both producer and consumer are clients in the client/server model, while the message queue is the server. The clients and servers communicate over the network.
Messaging models
The most popular messaging models are point-to-point and publish-subscribe.
Point-to-point
This model is commonly found in traditional message queues. In a point-to-point model, a message is sent to a queue and consumed by one and only one consumer. There can be multiple consumers waiting to consume messages in the queue, but each message can only be consumed by a single consumer. In Figure 3, message A is only consumed by consumer 1.
Once the consumer acknowledges that a message is consumed, it is removed from the queue. There is no data retention in the point-to-point model. In contrast, our design includes a persistence layer that keeps the messages for two weeks, which allows messages to be repeatedly consumed.
While our design could simulate a point-to-point model, its capabilities map more naturally to the publish-subscribe model.
Publish-subscribe
First, let’s introduce a new concept, the topic. Topics are the categories used to organize messages. Each topic has a name that is unique across the entire message queue service. Messages are sent to and read from a specific topic.
In the publish-subscribe model, a message is sent to a topic and received by the consumers subscribing to this topic. As shown in Figure 4, message A is consumed by both consumer 1 and consumer 2.
Our distributed message queue supports both models. The publish-subscribe model is implemented by topics, and the point-to-point model can be simulated by the concept of the consumer group, which will be introduced in the consumer group section.
Topics, partitions, and brokers
As mentioned earlier, messages are persisted by topics. What if the data volume in a topic is too large for a single server to handle?
One approach to solve this problem is called partition (sharding). As Figure 5 shows, we divide a topic into partitions and deliver messages evenly across partitions. Think of a partition as a small subset of the messages for a topic. Partitions are evenly distributed across the servers in the message queue cluster. These servers that hold partitions are called brokers. The distribution of partitions among brokers is the key element to support high scalability. We can scale the topic capacity by expanding the number of partitions.
Each topic partition operates in the form of a queue with the FIFO (first in, first out) mechanism. This means we can keep the order of messages inside a partition. The position of a message in the partition is called an offset.
When a message is sent by a producer, it is actually sent to one of the partitions for the topic. Each message has an optional message key (for example, a user’s ID), and all messages for the same message key are sent to the same partition. If the message key is absent, the message is randomly sent to one of the partitions.
When a consumer subscribes to a topic, it pulls data from one or more of these partitions. When there are multiple consumers subscribing to a topic, each consumer is responsible for a subset of the partitions for the topic. The consumers form a consumer group for a topic.
The message queue cluster with brokers and partitions is represented in Figure 6.
Consumer group
As mentioned earlier, we need to support both point-to-point and publish-subscribe models. A consumer group is a set of consumers, working together to consume messages from topics.
Consumers can be organized into groups. Each consumer group can subscribe to multiple topics and maintain its own consuming offsets. For example, we can group consumers by use cases, one group for billing and the other for accounting.
The instances in the same group can consume traffic in parallel, as in Figure 7.
-
Consumer group 1 subscribes to topic A.
-
Consumer group 2 subscribes to both topics A and B.
-
Topic A is subscribed by both consumer groups 1 and 2, which means the same message is consumed by multiple consumers. This pattern supports the subscribe/publish model.
However, there is one problem. Reading data in parallel improves the throughput, but the consumption order of messages in the same partition cannot be guaranteed. For example, if consumer 1 and consumer 2 both read from partition 1, we will not be able to guarantee the message consumption order in partition 1.
The good news is we can fix this by adding a constraint, that a single partition can only be consumed by one consumer in the same group. If the number of consumers of a group is larger than the number of partitions of a topic, some consumers will not get data from this topic. For example, in Figure 7, consumer 3 in group 2 cannot consume messages from topic B because it is consumed by Consumer-4 in the same consumer group, already.
With this constraint, if we put all consumers in the same consumer group, then messages in the same partition are consumed by only one consumer, which is equivalent to the point-to-point model. Since a partition is the smallest storage unit, we can allocate enough partitions in advance to avoid the need to dynamically increase the number of partitions. To handle high scale, we just need to add consumers.
High-Level architecture
Figure 8 shows the updated high-level design.
Clients
-
Producer: pushes messages to specific topics.
-
Consumer group: subscribes to topics and consumes messages.
Core service and storage
-
Broker: holds multiple partitions. A partition holds a subset of messages for a topic.
-
Storage:
-
Data storage: messages are persisted in data storage in partitions.
-
State storage: consumer states are managed by state storage.
-
Metadata storage: configuration and properties of topics are persisted in metadata storage.
-
-
Coordination service:
-
Service discovery: which brokers are alive.
-
Leader election: one of the brokers is selected as the active controller. There is only one active controller in the cluster. The active controller is responsible for assigning partitions.
-
Apache Zookeeper 2 or etcd 3 are commonly used to elect a controller.
-
Finished reading?
Mark it complete to track your progress.