Compare commits
2 Commits
main
...
2919282a5c
| Author | SHA1 | Date | |
|---|---|---|---|
| 2919282a5c | |||
| 8b5587cf88 |
16
README.md
16
README.md
@@ -10,14 +10,12 @@ module's `scripts/run-all.sh`, never typed by hand.
|
|||||||
| [`kafka-basics/`](kafka-basics/README.md) | [Spring Boot 4.1 and Apache Kafka: Producer, Consumer and Serialisation from Scratch](https://ankurm.com/spring-boot-4-1-kafka-producer-consumer-serialisation/) | The on-ramp: what the starter gives you, the two Jackson serializer families, where a key lands and why, and which defaults are Kafka's rather than Spring's |
|
| [`kafka-basics/`](kafka-basics/README.md) | [Spring Boot 4.1 and Apache Kafka: Producer, Consumer and Serialisation from Scratch](https://ankurm.com/spring-boot-4-1-kafka-producer-consumer-serialisation/) | The on-ramp: what the starter gives you, the two Jackson serializer families, where a key lands and why, and which defaults are Kafka's rather than Spring's |
|
||||||
| [`kafka-error-handling/`](kafka-error-handling/README.md) | [Kafka Error Handling with Spring Kafka 4.1: DLT, Retry Topics and Poison Pills](https://ankurm.com/spring-kafka-4-1-error-handling-dlt-retry-topics/) | What the default error handler really does, why a poison pill stops a partition, and what blocking retries cost that retry topics do not |
|
| [`kafka-error-handling/`](kafka-error-handling/README.md) | [Kafka Error Handling with Spring Kafka 4.1: DLT, Retry Topics and Poison Pills](https://ankurm.com/spring-kafka-4-1-error-handling-dlt-retry-topics/) | What the default error handler really does, why a poison pill stops a partition, and what blocking retries cost that retry topics do not |
|
||||||
| [`rabbitmq/`](rabbitmq/README.md) | [Spring Boot and RabbitMQ: Exchanges, Queues, Bindings and a Working Dead-Letter Queue](https://ankurm.com/spring-boot-rabbitmq-exchanges-dead-letter-queue/) | All four exchange types against a real broker, manual acknowledgement, and a dead-letter path exercised through both rejection and TTL expiry |
|
| [`rabbitmq/`](rabbitmq/README.md) | [Spring Boot and RabbitMQ: Exchanges, Queues, Bindings and a Working Dead-Letter Queue](https://ankurm.com/spring-boot-rabbitmq-exchanges-dead-letter-queue/) | All four exchange types against a real broker, manual acknowledgement, and a dead-letter path exercised through both rejection and TTL expiry |
|
||||||
| [`broker-comparison/`](broker-comparison/README.md) | [Kafka vs RabbitMQ vs Pulsar for Java Teams: A Decision Framework with Benchmarks](https://ankurm.com/kafka-vs-rabbitmq-vs-pulsar-java-decision-framework/) | The three brokers measured side by side on ordering, replay, consumer scaling and operational footprint, ending in a decision table where every row has a transcript behind it |
|
|
||||||
|
|
||||||
The brokers make an instructive set. Kafka's consumer holds an offset and the broker remembers
|
The two brokers make an instructive pair. Kafka's consumer holds an offset and the broker
|
||||||
nothing about individual records; RabbitMQ's broker owns the message until it is acknowledged and
|
remembers nothing about individual records; RabbitMQ's broker owns the message until it is
|
||||||
can route, expire and dead-letter it on its own; Pulsar keeps the log but lets each subscriber
|
acknowledged and can route, expire and dead-letter it on its own. Almost every difference in how
|
||||||
choose how it is read. Almost every difference in how you handle failure follows from those three
|
you handle failure follows from that one sentence — which is why Kafka needs a retry topic
|
||||||
sentences — which is why Kafka needs a retry topic to do what RabbitMQ does with a queue
|
to do what RabbitMQ does with a queue argument.
|
||||||
argument, and why `broker-comparison` is mostly an argument about where a message lives.
|
|
||||||
|
|
||||||
They are meant to be read in order. `kafka-basics` establishes that the default acknowledgement
|
They are meant to be read in order. `kafka-basics` establishes that the default acknowledgement
|
||||||
mode is `BATCH` and therefore that delivery is at-least-once; everything the error-handling
|
mode is `BATCH` and therefore that delivery is at-least-once; everything the error-handling
|
||||||
@@ -31,8 +29,8 @@ Central and from Boot's own `spring-boot-dependencies` POM, rather than from rel
|
|||||||
announcements.
|
announcements.
|
||||||
|
|
||||||
Every module runs its broker without Docker, so the transcripts can be regenerated on any
|
Every module runs its broker without Docker, so the transcripts can be regenerated on any
|
||||||
machine with a JDK: Kafka via the in-process KRaft broker in `spring-kafka-test`, RabbitMQ and
|
machine with a JDK: Kafka via the in-process KRaft broker in `spring-kafka-test`, RabbitMQ via a
|
||||||
Pulsar via real brokers started by the modules' own scripts. Each module also ships a Testcontainers
|
real broker started by the module's own scripts. Each module also ships a Testcontainers
|
||||||
configuration for the cases where the deployed image is what you need to test against.
|
configuration for the cases where the deployed image is what you need to test against.
|
||||||
|
|
||||||
## Licence
|
## Licence
|
||||||
|
|||||||
@@ -1,80 +0,0 @@
|
|||||||
# broker-comparison
|
|
||||||
|
|
||||||
Companion project for **Kafka vs RabbitMQ vs Pulsar for Java Teams: A Decision Framework with
|
|
||||||
Benchmarks** on [ankurm.com](https://ankurm.com).
|
|
||||||
|
|
||||||
Three brokers, three questions, one set of measurements. Ordering, replay and consumer scaling are
|
|
||||||
asked identically of all three by driving the client libraries directly, so that what is compared
|
|
||||||
is the brokers' delivery models rather than three sets of Spring defaults. What each Spring
|
|
||||||
integration adds on top is [chapter 6](docs/06-what-spring-adds.md).
|
|
||||||
|
|
||||||
## Verified stack
|
|
||||||
|
|
||||||
| Component | Version | Source of the number |
|
|
||||||
|---|---|---|
|
|
||||||
| JDK (tests) | 25.0.4.1+1 (Temurin) | `java -version` |
|
|
||||||
| JDK (brokers) | 21.0.12.1+1 (Temurin) | `java -version` |
|
|
||||||
| Spring Boot | 4.1.1 | `maven-metadata.xml` on Maven Central |
|
|
||||||
| kafka-clients | 4.2.1 | `spring-boot-dependencies-4.1.1.pom` |
|
|
||||||
| amqp-client | 5.30.0 | `spring-boot-dependencies-4.1.1.pom` |
|
|
||||||
| pulsar-client | 4.2.4 | `spring-boot-dependencies-4.1.1.pom` |
|
|
||||||
| Kafka broker | 4.2.1 (KRaft) | `kafka_2.13-4.2.1.tgz` |
|
|
||||||
| RabbitMQ broker | 3.10.25 | generic-unix tarball |
|
|
||||||
| Pulsar broker | 4.2.4 standalone | `apache-pulsar-4.2.4-bin.tar.gz` |
|
|
||||||
|
|
||||||
RabbitMQ is pinned at 3.10.25 because the box that produced these transcripts has Erlang 24;
|
|
||||||
3.11 and later need Erlang 25. The AMQP 0-9-1 semantics measured here are unchanged in 4.x.
|
|
||||||
|
|
||||||
## Quickstart
|
|
||||||
|
|
||||||
The Kafka measurements need no external broker — `spring-kafka-test` starts a real KRaft broker
|
|
||||||
in-process. The other two need a broker each, and the scripts start both without Docker and
|
|
||||||
without root.
|
|
||||||
|
|
||||||
```bash
|
|
||||||
mvn -Dgroups=kafka test
|
|
||||||
|
|
||||||
ERL_ROOT=... RABBITMQ_HOME=... scripts/rabbit-broker.sh
|
|
||||||
mvn -Dgroups=rabbit test
|
|
||||||
|
|
||||||
PULSAR_HOME=... scripts/pulsar-broker.sh
|
|
||||||
mvn -Dgroups=pulsar test
|
|
||||||
|
|
||||||
scripts/footprint.sh kafka|rabbit|pulsar # the operational numbers
|
|
||||||
scripts/run-all.sh # regenerates every docs/output/ file
|
|
||||||
```
|
|
||||||
|
|
||||||
## Documentation
|
|
||||||
|
|
||||||
| Chapter | What it settles |
|
|
||||||
|---|---|
|
|
||||||
| [01 Three models](docs/01-three-models.md) | Where a message lives, and why everything else follows |
|
|
||||||
| [02 Ordering](docs/02-ordering.md) | What survives the second consumer, in all three |
|
|
||||||
| [03 Replay](docs/03-replay.md) | Reading it twice, and the one broker that cannot |
|
|
||||||
| [04 Consumer scaling](docs/04-consumer-scaling.md) | The partition ceiling, and the client buffer that defeats fan-out in the other two |
|
|
||||||
| [05 Operational footprint](docs/05-operational-footprint.md) | Startup, memory, ports, processes, configuration surface |
|
|
||||||
| [06 What Spring adds](docs/06-what-spring-adds.md) | The three integrations, and the Boot 4 starter trap |
|
|
||||||
| [07 The decision table](docs/07-the-decision-table.md) | Every row backed by a measurement, and how to choose |
|
|
||||||
|
|
||||||
## Captured output
|
|
||||||
|
|
||||||
| File | What it shows |
|
|
||||||
|---|---|
|
|
||||||
| [`kafka-ordering.txt`](docs/output/kafka-ordering.txt) | Per-key order kept, global order not — and why A, B and C were the wrong keys to demonstrate it |
|
|
||||||
| [`kafka-replay.txt`](docs/output/kafka-replay.txt) | The same 12 records read three times |
|
|
||||||
| [`kafka-consumer-scaling.txt`](docs/output/kafka-consumer-scaling.txt) | 5 consumers, 3 partitions, 2 idle |
|
|
||||||
| [`rabbit-ordering.txt`](docs/output/rabbit-ordering.txt) | Message 1 finishing tenth |
|
|
||||||
| [`rabbit-replay.txt`](docs/output/rabbit-replay.txt) | 12 messages, then 0, then an empty queue |
|
|
||||||
| [`rabbit-consumer-scaling.txt`](docs/output/rabbit-consumer-scaling.txt) | 40/0/0/0/0 without `basicQos`, 8 each with it |
|
|
||||||
| [`pulsar-ordering.txt`](docs/output/pulsar-ordering.txt) | `Shared` spreading a key, `Key_Shared` pinning it |
|
|
||||||
| [`pulsar-replay.txt`](docs/output/pulsar-replay.txt) | `seek(earliest)` and a new subscription |
|
|
||||||
| [`pulsar-consumer-scaling.txt`](docs/output/pulsar-consumer-scaling.txt) | The 1000-message receiver queue defeating fan-out |
|
|
||||||
| [`footprint.txt`](docs/output/footprint.txt) | Startup, memory, ports and config surface for all three |
|
|
||||||
| [`tests.txt`](docs/output/tests.txt) | The nine tests behind all of the above |
|
|
||||||
|
|
||||||
## What is not measured here
|
|
||||||
|
|
||||||
Throughput. A messages-per-second number from one 2-core container says nothing about any of
|
|
||||||
these brokers, and publishing one would be worse than publishing nothing. The measurements above
|
|
||||||
are structural: they hold on any hardware, because they are properties of the delivery models
|
|
||||||
rather than of the machine.
|
|
||||||
@@ -1,49 +0,0 @@
|
|||||||
[README](../README.md) · next: [Ordering](02-ordering.md)
|
|
||||||
|
|
||||||
# 1. Three storage models, and everything that follows from them
|
|
||||||
|
|
||||||
Almost every difference between these brokers follows from one sentence about where a message
|
|
||||||
lives.
|
|
||||||
|
|
||||||
- **Kafka** is an append-only log, partitioned. The broker keeps every record for the retention
|
|
||||||
period and remembers nothing about individual consumers except an offset. Reading does not
|
|
||||||
remove anything.
|
|
||||||
- **RabbitMQ** is a router with queues. The broker owns each message until a consumer
|
|
||||||
acknowledges it, and then deletes it. It can route, expire, and dead-letter on its own.
|
|
||||||
- **Pulsar** is a log too, but the cursor and the storage are separate services: brokers are
|
|
||||||
stateless and BookKeeper holds the data. Acknowledgement moves a cursor, and a message is
|
|
||||||
deleted once every subscription has passed it — unless a retention policy says otherwise.
|
|
||||||
|
|
||||||
Read those three sentences again before reading any comparison table, including the one in
|
|
||||||
[chapter 7](07-the-decision-table.md). They predict most of it.
|
|
||||||
|
|
||||||
| | Kafka | RabbitMQ | Pulsar |
|
|
||||||
|---|---|---|---|
|
|
||||||
| unit of parallelism | partition | queue | subscription |
|
|
||||||
| set when | topic is created | queue is declared | consumer subscribes |
|
|
||||||
| consuming | moves an offset | deletes the message | moves a cursor |
|
|
||||||
| a second reader | new consumer group | another queue, bound in advance | new subscription |
|
|
||||||
|
|
||||||
The last row is the one people underestimate. In Kafka and Pulsar you can add a reader that sees
|
|
||||||
history you have already processed. In RabbitMQ you cannot add one after the fact at all — the
|
|
||||||
messages are gone — so the decision to have a second consumer has to be made *before* the
|
|
||||||
messages arrive.
|
|
||||||
|
|
||||||
## What was measured, and how
|
|
||||||
|
|
||||||
Three questions, asked identically of all three brokers by driving the client libraries directly
|
|
||||||
rather than through three different Spring abstractions:
|
|
||||||
|
|
||||||
| Question | Chapter |
|
|
||||||
|---|---|
|
|
||||||
| What ordering survives when you add a second consumer? | [02](02-ordering.md) |
|
|
||||||
| Can you read the same message twice? | [03](03-replay.md) |
|
|
||||||
| How far do consumers scale, and what stops them? | [04](04-consumer-scaling.md) |
|
|
||||||
| What does each one cost to run? | [05](05-operational-footprint.md) |
|
|
||||||
|
|
||||||
Versions are the ones a Spring Boot 4.1.1 application gets: kafka-clients 4.2.1, amqp-client
|
|
||||||
5.30.0, pulsar-client 4.2.4, all read from `spring-boot-dependencies-4.1.1.pom`. The brokers are
|
|
||||||
Kafka 4.2.1, RabbitMQ 3.10.25 and Pulsar 4.2.4, each started from its own distribution by a script
|
|
||||||
in `scripts/`.
|
|
||||||
|
|
||||||
next: [Ordering](02-ordering.md)
|
|
||||||
@@ -1,66 +0,0 @@
|
|||||||
prev: [Three models](01-three-models.md) · [README](../README.md) · next: [Replay](03-replay.md)
|
|
||||||
|
|
||||||
# 2. Ordering
|
|
||||||
|
|
||||||
Every broker here is FIFO with one producer and one consumer. The question that matters is what
|
|
||||||
survives the second consumer, because that is the first thing you add.
|
|
||||||
|
|
||||||
## Kafka: ordered within a partition, and a key picks the partition
|
|
||||||
|
|
||||||
Twelve records, three keys, three partitions
|
|
||||||
([`docs/output/kafka-ordering.txt`](output/kafka-ordering.txt)):
|
|
||||||
|
|
||||||
```
|
|
||||||
F -> partition [2] values [2, 5, 8, 11]
|
|
||||||
A -> partition [1] values [1, 4, 7, 10]
|
|
||||||
D -> partition [0] values [3, 6, 9, 12]
|
|
||||||
```
|
|
||||||
|
|
||||||
Each key's values come back in the order they were produced. The global sequence does not: the
|
|
||||||
consumer drains one partition's buffer before moving to the next, so the delivered order is
|
|
||||||
2, 5, 8, 11, 1, 4, 7, 10, 3, 6, 9, 12.
|
|
||||||
|
|
||||||
That is the whole Kafka ordering guarantee, and it is usually enough, because "in order" nearly
|
|
||||||
always means "in order per customer / per account / per device" rather than globally.
|
|
||||||
|
|
||||||
**A detail worth stealing.** The keys in that transcript are D, A and F, not A, B and C. With
|
|
||||||
three partitions, murmur2 sends A, B *and* C all to partition 1. Picking three obvious keys to
|
|
||||||
demonstrate partitioning would have produced a transcript in which global order was accidentally
|
|
||||||
preserved. Before you conclude that your keys spread evenly, compute
|
|
||||||
`Utils.toPositive(Utils.murmur2(key)) % partitions` for the ones you actually use.
|
|
||||||
|
|
||||||
## RabbitMQ: FIFO per queue, and no key at all
|
|
||||||
|
|
||||||
One consumer sees the queue in order. Two consumers do not
|
|
||||||
([`docs/output/rabbit-ordering.txt`](output/rabbit-ordering.txt)):
|
|
||||||
|
|
||||||
```
|
|
||||||
one queue, one consumer : [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12]
|
|
||||||
one queue, two consumers : [2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 1, 12]
|
|
||||||
```
|
|
||||||
|
|
||||||
Message 1 went to the slower consumer and finished tenth. There is no key and no partition, so
|
|
||||||
there is nothing to pin related messages to one consumer. Ordering across related messages in
|
|
||||||
RabbitMQ means one queue and one consumer — which means no horizontal scaling of that queue —
|
|
||||||
or the consistent-hash exchange, which is a plugin and effectively reintroduces partitions by
|
|
||||||
hand.
|
|
||||||
|
|
||||||
## Pulsar: the subscription type decides, per subscriber
|
|
||||||
|
|
||||||
Two consumers on a `Shared` subscription and two on `Key_Shared`, same topic, same messages
|
|
||||||
([`docs/output/pulsar-ordering.txt`](output/pulsar-ordering.txt)):
|
|
||||||
|
|
||||||
```
|
|
||||||
Shared A -> [consumer-1, consumer-2] B -> [consumer-1, consumer-2] C -> [consumer-1, consumer-2]
|
|
||||||
Key_Shared A -> [consumer-2] B -> [consumer-1] C -> [consumer-2]
|
|
||||||
```
|
|
||||||
|
|
||||||
Under `Shared`, every key was handled by both consumers, so two messages with the same key can be
|
|
||||||
in flight at once and no ordering claim survives. Under `Key_Shared`, each key went to exactly one
|
|
||||||
consumer — Kafka's guarantee, but chosen by the subscriber rather than fixed by a partition count
|
|
||||||
someone picked when the topic was created, and recomputed as consumers join and leave.
|
|
||||||
|
|
||||||
That last property is the strongest single argument for Pulsar: two teams reading the same topic
|
|
||||||
can make different ordering-versus-throughput trades without negotiating.
|
|
||||||
|
|
||||||
next: [Replay](03-replay.md)
|
|
||||||
@@ -1,59 +0,0 @@
|
|||||||
prev: [Ordering](02-ordering.md) · [README](../README.md) · next: [Consumer scaling](04-consumer-scaling.md)
|
|
||||||
|
|
||||||
# 3. Replay
|
|
||||||
|
|
||||||
"Can I read that again?" is asked after a bad deploy, and by then it is too late to change
|
|
||||||
brokers.
|
|
||||||
|
|
||||||
## Kafka
|
|
||||||
|
|
||||||
```
|
|
||||||
group replay-group-1, first read : 12 records
|
|
||||||
group replay-group-2, brand new group : 12 records
|
|
||||||
group replay-group-1, after seekToBeginning: 12 records
|
|
||||||
```
|
|
||||||
|
|
||||||
([`docs/output/kafka-replay.txt`](output/kafka-replay.txt)) Consuming moves an offset and removes
|
|
||||||
nothing. A new group, a reset offset and a `seekToBeginning` all read the same records, for as
|
|
||||||
long as retention keeps them.
|
|
||||||
|
|
||||||
## RabbitMQ
|
|
||||||
|
|
||||||
```
|
|
||||||
first drain of the queue : 12 messages
|
|
||||||
second drain of the queue : 0 messages
|
|
||||||
queue depth afterwards : 0
|
|
||||||
```
|
|
||||||
|
|
||||||
([`docs/output/rabbit-replay.txt`](output/rabbit-replay.txt)) Acknowledging deletes. There is no
|
|
||||||
offset to rewind and no second reader that can see what the first consumed. Reading a message
|
|
||||||
twice has to be arranged in advance — a second queue bound to the same exchange, or a copy written
|
|
||||||
somewhere else — and it cannot be arranged afterwards.
|
|
||||||
|
|
||||||
This is not a defect. A router that deletes what it has delivered is much cheaper to operate than
|
|
||||||
a log, and most work queues genuinely do not need history. It only becomes a defect at the moment
|
|
||||||
you need history and do not have it.
|
|
||||||
|
|
||||||
RabbitMQ streams (3.9+) are a separate, log-shaped feature that does support replay. They are a
|
|
||||||
different thing living in the same broker, with their own client protocol and their own semantics
|
|
||||||
— worth knowing about, and not what you get from `queueDeclare`.
|
|
||||||
|
|
||||||
## Pulsar
|
|
||||||
|
|
||||||
```
|
|
||||||
subscription replay-sub, first read : 12 messages
|
|
||||||
subscription replay-sub, after seek(earliest): 12 messages
|
|
||||||
subscription replay-sub-2, brand new : 12 messages
|
|
||||||
```
|
|
||||||
|
|
||||||
([`docs/output/pulsar-replay.txt`](output/pulsar-replay.txt)) `seek(MessageId)` and
|
|
||||||
`seek(timestamp)` rewind a live subscription; a new subscription starting from `Earliest` reads
|
|
||||||
everything still stored.
|
|
||||||
|
|
||||||
The difference from Kafka is the **default**, and it is the one that surprises people: Kafka keeps
|
|
||||||
a record for the retention period regardless of who read it, while Pulsar deletes a message once
|
|
||||||
every subscription has acknowledged it, unless a retention policy on the namespace says otherwise.
|
|
||||||
A Pulsar namespace with the default retention and one well-behaved subscription keeps nothing —
|
|
||||||
so the replay you are counting on requires a policy you have to set.
|
|
||||||
|
|
||||||
next: [Consumer scaling](04-consumer-scaling.md)
|
|
||||||
@@ -1,80 +0,0 @@
|
|||||||
prev: [Replay](03-replay.md) · [README](../README.md) · next: [Operational footprint](05-operational-footprint.md)
|
|
||||||
|
|
||||||
# 4. Consumer scaling
|
|
||||||
|
|
||||||
The question is what happens when you add the fifth consumer.
|
|
||||||
|
|
||||||
## Kafka: partitions are a hard ceiling
|
|
||||||
|
|
||||||
Five consumers in one group on a three-partition topic
|
|
||||||
([`docs/output/kafka-consumer-scaling.txt`](output/kafka-consumer-scaling.txt)):
|
|
||||||
|
|
||||||
```
|
|
||||||
consumer-1 -> partitions [0]
|
|
||||||
consumer-2 -> partitions [1]
|
|
||||||
consumer-3 -> partitions [2]
|
|
||||||
consumer-4 -> no partitions (idle)
|
|
||||||
consumer-5 -> no partitions (idle)
|
|
||||||
```
|
|
||||||
|
|
||||||
A partition belongs to at most one consumer in a group, so partition count is the ceiling. Two of
|
|
||||||
the five processes are running, connected, healthy, and doing nothing at all — and they will keep
|
|
||||||
reporting healthy forever.
|
|
||||||
|
|
||||||
Raising the partition count later is possible and is not free: it changes which partition a key
|
|
||||||
maps to, so the per-key ordering guarantee is broken across the change for every key that moves.
|
|
||||||
In practice the partition count is a capacity decision made at design time, on incomplete
|
|
||||||
information, that you then live with.
|
|
||||||
|
|
||||||
## RabbitMQ: no ceiling, but prefetch decides whether it is real
|
|
||||||
|
|
||||||
Five consumers on one queue, forty messages
|
|
||||||
([`docs/output/rabbit-consumer-scaling.txt`](output/rabbit-consumer-scaling.txt)):
|
|
||||||
|
|
||||||
```
|
|
||||||
no basicQos at all (unlimited prefetch, the AMQP default):
|
|
||||||
consumer-1 -> 40 consumer-2..5 -> 0
|
|
||||||
|
|
||||||
basicQos(1):
|
|
||||||
consumer-1..5 -> 8 each
|
|
||||||
```
|
|
||||||
|
|
||||||
There is no structural ceiling — but with the AMQP default the broker pushes as many messages as
|
|
||||||
a consumer will accept, so the first consumer to connect can be handed the entire backlog while
|
|
||||||
four idle processes wait. `basicQos` is not a tuning knob to postpone; it is what makes the
|
|
||||||
fan-out exist.
|
|
||||||
|
|
||||||
Spring AMQP sets it for you: `AbstractMessageListenerContainer.DEFAULT_PREFETCH_COUNT` is 250.
|
|
||||||
Better than unlimited, and still enough to concentrate any backlog smaller than 250 messages on
|
|
||||||
one consumer.
|
|
||||||
|
|
||||||
## Pulsar: no ceiling either, and the same trap under a different name
|
|
||||||
|
|
||||||
Five consumers, one *non-partitioned* topic, forty messages
|
|
||||||
([`docs/output/pulsar-consumer-scaling.txt`](output/pulsar-consumer-scaling.txt)):
|
|
||||||
|
|
||||||
```
|
|
||||||
receiverQueueSize default (1000):
|
|
||||||
one or two consumers take everything; the rest receive nothing
|
|
||||||
|
|
||||||
receiverQueueSize(1):
|
|
||||||
consumer-1..5 -> 8 each
|
|
||||||
```
|
|
||||||
|
|
||||||
The topic has no partitions and five consumers still share the work, which in Kafka would require
|
|
||||||
at least five partitions decided in advance. But the default receiver queue is a thousand
|
|
||||||
messages, so the first consumers to connect pull the whole backlog into their own buffers first.
|
|
||||||
Same failure as RabbitMQ's unlimited prefetch, four times the default.
|
|
||||||
|
|
||||||
## The shape of the answer
|
|
||||||
|
|
||||||
- Kafka gives you ordering per key and pays for it with a fixed parallelism ceiling.
|
|
||||||
- RabbitMQ gives you unlimited competing consumers and pays for it with no ordering across them.
|
|
||||||
- Pulsar lets each subscription choose, which is genuinely the best of both — at the operational
|
|
||||||
cost measured in [the next chapter](05-operational-footprint.md).
|
|
||||||
|
|
||||||
And all three brokers have a client-side buffer that will quietly defeat your fan-out if you leave
|
|
||||||
it at the default and your backlog is small. That is the one finding here that applies to
|
|
||||||
whichever you pick.
|
|
||||||
|
|
||||||
next: [Operational footprint](05-operational-footprint.md)
|
|
||||||
@@ -1,55 +0,0 @@
|
|||||||
prev: [Consumer scaling](04-consumer-scaling.md) · [README](../README.md) · next: [What Spring adds](06-what-spring-adds.md)
|
|
||||||
|
|
||||||
# 5. Operational footprint
|
|
||||||
|
|
||||||
"Ops burden" cannot be benchmarked. The number of moving parts, the memory floor and the size of
|
|
||||||
the configuration surface can be, and they are what people are actually asking about.
|
|
||||||
|
|
||||||
Each broker started from its shipped distribution with default configuration and a 700 MB heap
|
|
||||||
cap, on one 2-core / 3.8 GB Linux box, Temurin JDK 21
|
|
||||||
([`docs/output/footprint.txt`](output/footprint.txt)):
|
|
||||||
|
|
||||||
| | Kafka 4.2.1 | RabbitMQ 3.10.25 | Pulsar 4.2.4 standalone |
|
|
||||||
|---|---|---|---|
|
|
||||||
| launch to first accepted connection | 10.0 s | 15.1 s | 20.8 s |
|
|
||||||
| resident memory at idle | 333 MB | 115 MB | 610 MB |
|
|
||||||
| server processes | 1 | 1 (+ `epmd`) | 1 |
|
|
||||||
| listening ports | 9092, 9093 | 5672, 4369, 25672 | 6650, 8080, 2181 |
|
|
||||||
| unpacked distribution | 135 MB | 26 MB | 344 MB |
|
|
||||||
| settings in the shipped default config | 24 | **0** | 357 |
|
|
||||||
|
|
||||||
Read those numbers as orders of magnitude, not as a benchmark: one sample, one small machine.
|
|
||||||
|
|
||||||
Three things they say clearly.
|
|
||||||
|
|
||||||
**RabbitMQ ships no configuration file at all.** The boot log records `Config file(s): (none)`;
|
|
||||||
`rabbitmq.conf.example` is entirely commented out. Everything works out of the box, and the
|
|
||||||
smallest resident footprint here is the one that is not a JVM. Against that, `epmd` on 4369 and
|
|
||||||
the inter-node port on 25672 are the Erlang distribution, which is also how clustering works and
|
|
||||||
how clustering goes wrong.
|
|
||||||
|
|
||||||
**Pulsar standalone is three systems in one process.** Port 2181 is ZooKeeper and the on-disk
|
|
||||||
data is BookKeeper's; the standalone distribution hides that behind one command. A real
|
|
||||||
deployment does not: you operate ZooKeeper (or, from Pulsar 3.x, an alternative metadata store),
|
|
||||||
BookKeeper bookies and Pulsar brokers as three tiers with three scaling stories. The 357-setting
|
|
||||||
`standalone.conf` is the honest signal here — that is the configuration surface, and running it
|
|
||||||
seriously means learning most of it.
|
|
||||||
|
|
||||||
**Kafka since KRaft is genuinely simpler than it was.** One process, two ports, 24 settings, no
|
|
||||||
ZooKeeper. The old "Kafka means also running ZooKeeper" objection is a version behind; anyone
|
|
||||||
comparing on that basis is comparing to 2022.
|
|
||||||
|
|
||||||
## What the numbers do not include
|
|
||||||
|
|
||||||
The footprint above is one node at idle. The thing that actually determines ops burden is what
|
|
||||||
happens at three in the morning:
|
|
||||||
|
|
||||||
- **Kafka**: rebalances, consumer lag as the primary signal, partition-count decisions you cannot
|
|
||||||
cleanly reverse, and a broad ecosystem of tools that assume Kafka.
|
|
||||||
- **RabbitMQ**: the management plugin is genuinely good, queue depth is a direct and obvious
|
|
||||||
signal, and the hard problem is network partitions in a cluster — which the Erlang distribution
|
|
||||||
makes fast to detect and awkward to resolve.
|
|
||||||
- **Pulsar**: the fewest people on your team will have run it. That is not a technical property
|
|
||||||
and it is usually the deciding one.
|
|
||||||
|
|
||||||
next: [What Spring adds](06-what-spring-adds.md)
|
|
||||||
@@ -1,52 +0,0 @@
|
|||||||
prev: [Operational footprint](05-operational-footprint.md) · [README](../README.md) · next: [The decision table](07-the-decision-table.md)
|
|
||||||
|
|
||||||
# 6. What Spring adds, and what Boot 4 no longer gives you for free
|
|
||||||
|
|
||||||
The measurements in this module drive the client libraries directly, so that what is being
|
|
||||||
compared is the brokers rather than three sets of Spring defaults. In an application you would use
|
|
||||||
the Spring integrations, and they are not equivalent to one another.
|
|
||||||
|
|
||||||
| | Kafka | RabbitMQ | Pulsar |
|
|
||||||
|---|---|---|---|
|
|
||||||
| project | Spring for Apache Kafka 4.1.1 | Spring AMQP 4.1.1 | Spring for Apache Pulsar 2.0.7 |
|
|
||||||
| starter | `spring-boot-starter-kafka` | `spring-boot-starter-amqp` | `spring-boot-starter-pulsar` |
|
|
||||||
| listener | `@KafkaListener` | `@RabbitListener` | `@PulsarListener` |
|
|
||||||
| template | `KafkaTemplate` | `RabbitTemplate` | `PulsarTemplate` |
|
|
||||||
| retry / DLQ | `DefaultErrorHandler`, `@RetryableTopic` | dead-letter exchange, container error handler | `DeadLetterPolicy` on the listener |
|
|
||||||
|
|
||||||
All three are Boot-managed at 4.1.1, so you do not pin their versions.
|
|
||||||
|
|
||||||
## The Boot 4 trap that applies to all three
|
|
||||||
|
|
||||||
**In Spring Boot 4, depending on a messaging library directly rather than through its Boot starter
|
|
||||||
means you have no auto-configuration.** The auto-configuration classes moved out of
|
|
||||||
`spring-boot-autoconfigure` into per-technology modules that only the starters bring:
|
|
||||||
|
|
||||||
| depending on | what you lose | how it presents |
|
|
||||||
|---|---|---|
|
|
||||||
| `org.springframework.kafka:spring-kafka` | `spring-boot-kafka` | `No qualifying bean of type KafkaTemplate<...>` — the context starts fine |
|
|
||||||
| `org.springframework.amqp:spring-rabbit` | `spring-boot-amqp` | no `RabbitTemplate`, no `RabbitAdmin` |
|
|
||||||
| `org.springframework.pulsar:spring-pulsar` | `spring-boot-pulsar` | no `PulsarTemplate` |
|
|
||||||
|
|
||||||
Every Boot 3 tutorial gets this wrong now, and the symptom is a missing bean rather than anything
|
|
||||||
that names the cause. Use the starters.
|
|
||||||
|
|
||||||
## The Jackson fork, in both Kafka and RabbitMQ
|
|
||||||
|
|
||||||
Boot 4 moved to Jackson 3 (`tools.jackson`), and both messaging projects ship converters for
|
|
||||||
both generations. **A `2` in the class name means the previous Jackson** — which is the opposite
|
|
||||||
of the convention you would guess:
|
|
||||||
|
|
||||||
- Spring Kafka: `JsonSerializer`/`JsonDeserializer` are Jackson 2; `JacksonJsonSerializer`,
|
|
||||||
`JacksonJsonDeserializer` and `JacksonJsonSerde` are Jackson 3.
|
|
||||||
- Spring AMQP: `Jackson2JsonMessageConverter` is Jackson 2; `JacksonJsonMessageConverter` is
|
|
||||||
Jackson 3.
|
|
||||||
|
|
||||||
Picking the wrong one gives you `SerializationException: Can't serialize data`, whose cause is
|
|
||||||
`Java 8 date/time type java.time.Instant not supported by default` — a message about date types
|
|
||||||
for what is really a wrong-library problem.
|
|
||||||
|
|
||||||
The [kafka-basics](../kafka-basics/README.md) and [rabbitmq](../rabbitmq/README.md) modules in
|
|
||||||
this repository cover both stacks in detail.
|
|
||||||
|
|
||||||
next: [The decision table](07-the-decision-table.md)
|
|
||||||
@@ -1,59 +0,0 @@
|
|||||||
prev: [What Spring adds](06-what-spring-adds.md) · [README](../README.md)
|
|
||||||
|
|
||||||
# 7. The decision table
|
|
||||||
|
|
||||||
Every row below is either a measurement in `docs/output/` or a statement about operating the
|
|
||||||
thing. Nothing here is a vendor claim.
|
|
||||||
|
|
||||||
| | Kafka 4.2.1 | RabbitMQ 3.10.25 | Pulsar 4.2.4 |
|
|
||||||
|---|---|---|---|
|
|
||||||
| ordering with one consumer | FIFO per partition | FIFO per queue | FIFO |
|
|
||||||
| ordering with many consumers | per key, always | **none** | per key with `Key_Shared`, none with `Shared` |
|
|
||||||
| who chooses that | topic design | plugin, or one consumer | each subscription, independently |
|
|
||||||
| replay after the fact | yes, within retention | **no** | yes, if retention is configured |
|
|
||||||
| second independent reader | new consumer group, any time | must be arranged in advance | new subscription, any time |
|
|
||||||
| consumer parallelism ceiling | partition count | none | none |
|
|
||||||
| changing that ceiling | repartition; breaks key→partition | nothing to change | nothing to change |
|
|
||||||
| the client-side trap | none by default | unlimited prefetch (Spring: 250) | `receiverQueueSize` 1000 |
|
|
||||||
| memory at idle | 333 MB | **115 MB** | 610 MB |
|
|
||||||
| processes to operate | 1 (KRaft) | 1 + `epmd` | broker + BookKeeper + metadata store |
|
|
||||||
| shipped default settings | 24 | **0** | 357 |
|
|
||||||
| Spring project maturity | very high | very high | good, much smaller community |
|
|
||||||
|
|
||||||
## Choosing
|
|
||||||
|
|
||||||
**Choose RabbitMQ** when the work is a queue of tasks: each message is an instruction, order
|
|
||||||
between instructions does not matter, and once it is done it is done. It is the smallest thing
|
|
||||||
that works, it has the best out-of-the-box operability of the three, and queue depth is a metric
|
|
||||||
anyone can interpret. Choose it in the knowledge that you are giving up replay permanently —
|
|
||||||
adding it later is a broker migration, not a configuration change.
|
|
||||||
|
|
||||||
**Choose Kafka** when the messages are events other people will want to read: when more than one
|
|
||||||
consumer will exist, when someone will need to reprocess history after a bug, or when ordering per
|
|
||||||
entity is part of the contract. Pay for it with a partition count decided too early, consumers
|
|
||||||
that scale only as far as that number, and an operational model where lag is the thing you watch.
|
|
||||||
KRaft has removed the old ZooKeeper objection.
|
|
||||||
|
|
||||||
**Choose Pulsar** when you genuinely need what neither of the others gives you: per-subscription
|
|
||||||
choice of ordering versus fan-out, consumer counts not bounded by a number chosen at topic
|
|
||||||
creation, or multi-tenancy with real isolation. It is the most capable design here. It is also
|
|
||||||
three systems, 357 settings, five times RabbitMQ's memory floor, and the one your team has least
|
|
||||||
experience with — and that last point decides more incidents than the first three prevent.
|
|
||||||
|
|
||||||
**If the honest answer is "we do not know yet"**, that argues for Kafka, on grounds that have
|
|
||||||
nothing to do with the technology: you can hire for it, your monitoring vendor supports it, and
|
|
||||||
the failure modes are documented by thousands of people who hit them first. That is a real
|
|
||||||
engineering argument and it is fine to make it out loud.
|
|
||||||
|
|
||||||
## The question that dissolves the choice
|
|
||||||
|
|
||||||
A surprising number of these decisions are made for a system that has one producer, one consumer
|
|
||||||
and fewer than a hundred messages a second. At that volume all three brokers work, none of the
|
|
||||||
measurements above will ever be reached, and the decision is entirely about what your team can
|
|
||||||
operate at three in the morning.
|
|
||||||
|
|
||||||
The measurements matter when you can name which row you are relying on. If you cannot, you are
|
|
||||||
picking an operational burden, not a broker — and the cheapest one to operate is the one in the
|
|
||||||
`115 MB` cell.
|
|
||||||
|
|
||||||
[README](../README.md)
|
|
||||||
@@ -1,40 +0,0 @@
|
|||||||
== Operational footprint, measured on one 2-core / 3.8 GB Linux box ==
|
|
||||||
|
|
||||||
Each broker started from its shipped distribution with default configuration and a
|
|
||||||
700 MB heap cap, on Temurin JDK 21. Times are one sample on a small box: treat them as
|
|
||||||
orders of magnitude, not as a benchmark.
|
|
||||||
|
|
||||||
broker : kafka
|
|
||||||
time from launch to first
|
|
||||||
accepted connection : 10035 ms
|
|
||||||
resident memory at idle : 333 MB
|
|
||||||
server processes : 1
|
|
||||||
listening ports : 9092 (broker), 9093 (controller)
|
|
||||||
unpacked distribution size : 135 MB
|
|
||||||
settings in shipped default
|
|
||||||
configuration file : 24 (server.properties)
|
|
||||||
|
|
||||||
broker : rabbit
|
|
||||||
time from launch to first
|
|
||||||
accepted connection : 15053 ms
|
|
||||||
resident memory at idle : 115 MB
|
|
||||||
server processes : 1
|
|
||||||
listening ports : 5672 (AMQP), 4369 (epmd), 25672 (inter-node)
|
|
||||||
unpacked distribution size : 26 MB
|
|
||||||
settings in shipped default
|
|
||||||
configuration file : 0 (rabbitmq.conf.example)
|
|
||||||
epmd processes : 1
|
|
||||||
NOTE: RabbitMQ ships no configuration file at all. rabbitmq.conf.example is
|
|
||||||
entirely commented out and the boot log records "Config file(s): (none)".
|
|
||||||
Everything in that file is documentation, not a default.
|
|
||||||
|
|
||||||
broker : pulsar
|
|
||||||
time from launch to first
|
|
||||||
accepted connection : 20801 ms
|
|
||||||
resident memory at idle : 610 MB
|
|
||||||
server processes : 1
|
|
||||||
listening ports : 6650 (binary), 8080 (admin/REST), 2181 (ZooKeeper)
|
|
||||||
unpacked distribution size : 344 MB
|
|
||||||
settings in shipped default
|
|
||||||
configuration file : 357 (standalone.conf)
|
|
||||||
|
|
||||||
@@ -1,16 +0,0 @@
|
|||||||
== Kafka: partitions are the ceiling on consumer parallelism ==
|
|
||||||
|
|
||||||
topic 'orders-scaling', 3 partitions, 5 consumers in one group
|
|
||||||
|
|
||||||
consumer-1 -> partitions [0]
|
|
||||||
consumer-2 -> partitions [1]
|
|
||||||
consumer-3 -> partitions [2]
|
|
||||||
consumer-4 -> no partitions (idle)
|
|
||||||
consumer-5 -> no partitions (idle)
|
|
||||||
|
|
||||||
consumers with no partitions: 2
|
|
||||||
|
|
||||||
A partition is assigned to at most one consumer in a group, so the number
|
|
||||||
of partitions is a hard ceiling on consumer parallelism. Adding consumers
|
|
||||||
beyond it adds idle processes, not throughput.
|
|
||||||
|
|
||||||
@@ -1,30 +0,0 @@
|
|||||||
== Kafka: ordering is per partition ==
|
|
||||||
|
|
||||||
produced 12 records, keys D/A/F round-robin, values 1..12 in order
|
|
||||||
(keys chosen so that murmur2 spreads them: D->0, A->1, F->2. A, B and C
|
|
||||||
all hash to partition 1 with three partitions, which is worth knowing
|
|
||||||
before you decide your keys are well distributed.)
|
|
||||||
|
|
||||||
consumed in this order:
|
|
||||||
F=2@p2
|
|
||||||
F=5@p2
|
|
||||||
F=8@p2
|
|
||||||
F=11@p2
|
|
||||||
A=1@p1
|
|
||||||
A=4@p1
|
|
||||||
A=7@p1
|
|
||||||
A=10@p1
|
|
||||||
D=3@p0
|
|
||||||
D=6@p0
|
|
||||||
D=9@p0
|
|
||||||
D=12@p0
|
|
||||||
|
|
||||||
per key:
|
|
||||||
F -> partition [2] values [2, 5, 8, 11]
|
|
||||||
A -> partition [1] values [1, 4, 7, 10]
|
|
||||||
D -> partition [0] values [3, 6, 9, 12]
|
|
||||||
|
|
||||||
Order is preserved within each key because a key hashes to one partition.
|
|
||||||
Across keys it is not: the values above are not 1..12 in order, because a
|
|
||||||
consumer drains one partition's buffer before the next.
|
|
||||||
|
|
||||||
@@ -1,10 +0,0 @@
|
|||||||
== Kafka: the log is the storage ==
|
|
||||||
|
|
||||||
group replay-group-1, first read : 12 records
|
|
||||||
group replay-group-2, brand new group : 12 records
|
|
||||||
group replay-group-1, after seekToBeginning: 12 records
|
|
||||||
|
|
||||||
Consuming does not remove anything. A consumer group is a cursor over a log
|
|
||||||
that the broker keeps until retention expires, so a new group, a reset
|
|
||||||
offset or a seek all read the same records again.
|
|
||||||
|
|
||||||
@@ -1,30 +0,0 @@
|
|||||||
== Pulsar: no partition ceiling, but the receiver queue decides who gets the work ==
|
|
||||||
|
|
||||||
one non-partitioned topic, 40 messages, 5 consumers, Shared subscription
|
|
||||||
|
|
||||||
receiverQueueSize left at the default (1000):
|
|
||||||
consumer-1 -> 40 messages
|
|
||||||
consumer-2 -> 0 messages
|
|
||||||
consumer-3 -> 0 messages
|
|
||||||
consumer-4 -> 0 messages
|
|
||||||
consumer-5 -> 0 messages
|
|
||||||
consumers that received nothing : 4
|
|
||||||
|
|
||||||
receiverQueueSize(1):
|
|
||||||
consumer-1 -> 8 messages
|
|
||||||
consumer-2 -> 8 messages
|
|
||||||
consumer-3 -> 8 messages
|
|
||||||
consumer-4 -> 8 messages
|
|
||||||
consumer-5 -> 8 messages
|
|
||||||
consumers that received nothing : 0
|
|
||||||
|
|
||||||
The topic has no partitions and five consumers can still share the work --
|
|
||||||
in Kafka the same shape needs at least five partitions, chosen when the
|
|
||||||
topic was created. But the default receiver queue is 1000 messages, so the
|
|
||||||
first consumers to connect pull the whole 40-message backlog into their own
|
|
||||||
buffers before the rest ask for anything, and the subscription looks
|
|
||||||
broken. Which consumers win is a race and varies between runs -- one
|
|
||||||
consumer taking all forty, or two taking twenty each -- but the consumers
|
|
||||||
that lose it see nothing at all. This is the same trap as RabbitMQ's
|
|
||||||
unbounded prefetch, with a different name and a much larger default.
|
|
||||||
|
|
||||||
@@ -1,24 +0,0 @@
|
|||||||
== Pulsar: the subscription type decides ==
|
|
||||||
|
|
||||||
12 messages, keys A/B/C, values 1..12 in order, two consumers per
|
|
||||||
subscription, receiverQueueSize 1 so the first consumer cannot take the
|
|
||||||
whole backlog.
|
|
||||||
|
|
||||||
Shared subscription, which consumers saw each key:
|
|
||||||
A -> [consumer-1, consumer-2]
|
|
||||||
B -> [consumer-1, consumer-2]
|
|
||||||
C -> [consumer-1, consumer-2]
|
|
||||||
|
|
||||||
Key_Shared subscription, which consumers saw each key:
|
|
||||||
A -> [consumer-1]
|
|
||||||
B -> [consumer-1]
|
|
||||||
C -> [consumer-1]
|
|
||||||
|
|
||||||
A Shared subscription round-robins individual messages, so messages with
|
|
||||||
the same key end up on different consumers and can be processed at the
|
|
||||||
same time: there is no per-key order left to speak of. Key_Shared hashes
|
|
||||||
the key to one consumer, which is Kafka's guarantee -- except that the
|
|
||||||
assignment belongs to the subscription and is recomputed as consumers come
|
|
||||||
and go, rather than being fixed by a partition count chosen when the topic
|
|
||||||
was created.
|
|
||||||
|
|
||||||
@@ -1,14 +0,0 @@
|
|||||||
== Pulsar: acknowledged, but still there ==
|
|
||||||
|
|
||||||
subscription replay-sub, first read : 12 messages
|
|
||||||
subscription replay-sub, after seek(earliest): 12 messages
|
|
||||||
subscription replay-sub-2, brand new : 12 messages
|
|
||||||
|
|
||||||
Acknowledgement moves a cursor; the message itself lives in the managed
|
|
||||||
ledger. seek(MessageId) and seek(timestamp) rewind a live subscription,
|
|
||||||
which Kafka can also do by resetting offsets. What differs is the default:
|
|
||||||
Pulsar deletes a message once every subscription has acknowledged it,
|
|
||||||
unless a retention policy on the namespace says otherwise, whereas Kafka
|
|
||||||
keeps it for the retention period regardless of who read it. A Pulsar
|
|
||||||
topic with no retention policy and no subscriptions keeps nothing.
|
|
||||||
|
|
||||||
@@ -1,34 +0,0 @@
|
|||||||
== RabbitMQ: consumers scale, and prefetch decides whether they actually do ==
|
|
||||||
|
|
||||||
one queue, 40 messages, 5 consumers, each taking 10 ms per message
|
|
||||||
|
|
||||||
no basicQos at all (unlimited prefetch, the AMQP default):
|
|
||||||
consumer-1 -> 40 messages
|
|
||||||
consumer-2 -> 0 messages
|
|
||||||
consumer-3 -> 0 messages
|
|
||||||
consumer-4 -> 0 messages
|
|
||||||
consumer-5 -> 0 messages
|
|
||||||
consumers that received nothing : 4
|
|
||||||
|
|
||||||
basicQos(1):
|
|
||||||
consumer-1 -> 8 messages
|
|
||||||
consumer-2 -> 9 messages
|
|
||||||
consumer-3 -> 8 messages
|
|
||||||
consumer-4 -> 8 messages
|
|
||||||
consumer-5 -> 7 messages
|
|
||||||
consumers that received nothing : 0
|
|
||||||
|
|
||||||
Every consumer on a queue competes for the same messages, so adding
|
|
||||||
consumers adds throughput and there is no structural ceiling of the kind
|
|
||||||
Kafka's partition count imposes. But with the AMQP default the broker
|
|
||||||
pushes as many messages as a consumer will take, so whichever consumer
|
|
||||||
connects first can be handed the entire backlog while the others sit idle.
|
|
||||||
basicQos is not a tuning knob you get to postpone; it is what makes the
|
|
||||||
fan-out real. Spring AMQP sets it for you --
|
|
||||||
AbstractMessageListenerContainer.DEFAULT_PREFETCH_COUNT is 250 -- which is
|
|
||||||
better than unlimited and still large enough to concentrate a small
|
|
||||||
backlog on one consumer.
|
|
||||||
|
|
||||||
The price of all this is the ordering measurement: there is no key, so
|
|
||||||
nothing constrains related messages to one consumer.
|
|
||||||
|
|
||||||
@@ -1,18 +0,0 @@
|
|||||||
== RabbitMQ: FIFO per queue, per consumer ==
|
|
||||||
|
|
||||||
one queue, one consumer, 12 messages
|
|
||||||
completion order : [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12]
|
|
||||||
|
|
||||||
one queue, two consumers, prefetch 1, consumer-1 slower than consumer-2
|
|
||||||
completion order : [2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 1, 12]
|
|
||||||
out-of-order steps: 1
|
|
||||||
|
|
||||||
A queue is FIFO and a single consumer sees it that way. The moment a second
|
|
||||||
consumer is added the broker hands the next message to whichever consumer
|
|
||||||
is free, so the order in which work finishes is no longer the order in
|
|
||||||
which it was published. There is no key: RabbitMQ has no notion of a
|
|
||||||
partition to which related messages could be pinned. Ordering across
|
|
||||||
related messages means one queue and one consumer, and therefore no
|
|
||||||
horizontal scaling for that queue -- or a consistent-hash exchange, which
|
|
||||||
is a plugin.
|
|
||||||
|
|
||||||
@@ -1,12 +0,0 @@
|
|||||||
== RabbitMQ: there is nothing to replay ==
|
|
||||||
|
|
||||||
first drain of the queue : 12 messages
|
|
||||||
second drain of the queue : 0 messages
|
|
||||||
queue depth afterwards : 0
|
|
||||||
|
|
||||||
Acknowledging a message deletes it. The broker is a router with buffers,
|
|
||||||
not a log: there is no offset to rewind and no second reader that can see
|
|
||||||
what the first one consumed. Reading the same message twice means
|
|
||||||
arranging it in advance -- a second queue bound to the same exchange, or a
|
|
||||||
copy written somewhere else -- and it cannot be arranged after the fact.
|
|
||||||
|
|
||||||
@@ -1,12 +0,0 @@
|
|||||||
[INFO] Running com.ankurm.brokers.KafkaComparisonTest
|
|
||||||
[INFO] Tests run: 3, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 26.88 s -- in com.ankurm.brokers.KafkaComparisonTest
|
|
||||||
[INFO] Tests run: 3, Failures: 0, Errors: 0, Skipped: 0
|
|
||||||
[INFO] BUILD SUCCESS
|
|
||||||
[INFO] Running com.ankurm.brokers.RabbitComparisonTest
|
|
||||||
[INFO] Tests run: 3, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 4.732 s -- in com.ankurm.brokers.RabbitComparisonTest
|
|
||||||
[INFO] Tests run: 3, Failures: 0, Errors: 0, Skipped: 0
|
|
||||||
[INFO] BUILD SUCCESS
|
|
||||||
[INFO] Running com.ankurm.brokers.PulsarComparisonTest
|
|
||||||
[INFO] Tests run: 3, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 55.00 s -- in com.ankurm.brokers.PulsarComparisonTest
|
|
||||||
[INFO] Tests run: 3, Failures: 0, Errors: 0, Skipped: 0
|
|
||||||
[INFO] BUILD SUCCESS
|
|
||||||
@@ -1,66 +0,0 @@
|
|||||||
<project xmlns="http://maven.apache.org/POM/4.0.0">
|
|
||||||
<modelVersion>4.0.0</modelVersion>
|
|
||||||
|
|
||||||
<!-- Every client version below is managed by spring-boot-dependencies 4.1.1, so the three
|
|
||||||
brokers are compared at the versions a Boot 4.1 application would actually get:
|
|
||||||
kafka-clients 4.2.1, amqp-client 5.30.0, pulsar-client 4.2.4. Read from
|
|
||||||
spring-boot-dependencies-4.1.1.pom, not from release announcements. -->
|
|
||||||
<parent>
|
|
||||||
<groupId>org.springframework.boot</groupId>
|
|
||||||
<artifactId>spring-boot-starter-parent</artifactId>
|
|
||||||
<version>4.1.1</version>
|
|
||||||
<relativePath/>
|
|
||||||
</parent>
|
|
||||||
|
|
||||||
<groupId>com.ankurm</groupId>
|
|
||||||
<artifactId>broker-comparison</artifactId>
|
|
||||||
<version>1.0</version>
|
|
||||||
<packaging>jar</packaging>
|
|
||||||
|
|
||||||
<properties>
|
|
||||||
<java.version>25</java.version>
|
|
||||||
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
|
|
||||||
</properties>
|
|
||||||
|
|
||||||
<dependencies>
|
|
||||||
<dependency>
|
|
||||||
<groupId>org.springframework.boot</groupId>
|
|
||||||
<artifactId>spring-boot-starter</artifactId>
|
|
||||||
</dependency>
|
|
||||||
|
|
||||||
<!-- The three client libraries, driven directly. The comparison is between the brokers'
|
|
||||||
delivery models, and going through three different Spring abstractions would put three
|
|
||||||
different sets of defaults between the measurement and the thing being measured.
|
|
||||||
What each Boot starter adds on top is covered in docs/06-what-spring-adds.md. -->
|
|
||||||
<dependency>
|
|
||||||
<groupId>org.apache.kafka</groupId>
|
|
||||||
<artifactId>kafka-clients</artifactId>
|
|
||||||
</dependency>
|
|
||||||
<dependency>
|
|
||||||
<groupId>com.rabbitmq</groupId>
|
|
||||||
<artifactId>amqp-client</artifactId>
|
|
||||||
</dependency>
|
|
||||||
<dependency>
|
|
||||||
<groupId>org.apache.pulsar</groupId>
|
|
||||||
<artifactId>pulsar-client</artifactId>
|
|
||||||
</dependency>
|
|
||||||
|
|
||||||
<dependency>
|
|
||||||
<groupId>org.springframework.boot</groupId>
|
|
||||||
<artifactId>spring-boot-starter-test</artifactId>
|
|
||||||
<scope>test</scope>
|
|
||||||
</dependency>
|
|
||||||
<!-- Provides the in-process KRaft broker used by the Kafka measurements: a real Kafka
|
|
||||||
broker, no Docker. -->
|
|
||||||
<dependency>
|
|
||||||
<groupId>org.springframework.kafka</groupId>
|
|
||||||
<artifactId>spring-kafka-test</artifactId>
|
|
||||||
<scope>test</scope>
|
|
||||||
</dependency>
|
|
||||||
</dependencies>
|
|
||||||
|
|
||||||
<!-- Each broker is started by its own script and the matching tagged tests are run on
|
|
||||||
their own; see scripts/run-all.sh, which passes -Dgroups=kafka|rabbit|pulsar. Running
|
|
||||||
"mvn test" with no group runs all three and needs all three brokers up. -->
|
|
||||||
|
|
||||||
</project>
|
|
||||||
@@ -1,73 +0,0 @@
|
|||||||
#!/usr/bin/env bash
|
|
||||||
# Measures the operational footprint of one broker: how long it takes to accept a connection,
|
|
||||||
# how much memory it holds at idle, how many processes and listening ports it opens, and how
|
|
||||||
# large its shipped default configuration is.
|
|
||||||
#
|
|
||||||
# These are proxies, not a benchmark. "Ops burden" is not directly measurable, but the number of
|
|
||||||
# moving parts, the memory floor and the size of the configuration surface are, and they are what
|
|
||||||
# people are actually asking about when they ask which broker is heavier to run.
|
|
||||||
#
|
|
||||||
# scripts/footprint.sh kafka|rabbit|pulsar
|
|
||||||
set -euo pipefail
|
|
||||||
BROKER="${1:?kafka|rabbit|pulsar}"
|
|
||||||
|
|
||||||
extra=""
|
|
||||||
port_open() { (echo > "/dev/tcp/127.0.0.1/$1") 2>/dev/null; }
|
|
||||||
rss_mb() { ps -eo rss,cmd | grep -F "$1" | grep -v grep | awk '{s+=$1} END {printf "%.0f", s/1024}'; }
|
|
||||||
procs() { ps -eo cmd | grep -F "$1" | grep -cv grep; }
|
|
||||||
settings() { grep -cE '^[a-zA-Z]' "$1" 2>/dev/null || echo 0; }
|
|
||||||
|
|
||||||
start=$(date +%s%3N)
|
|
||||||
case "$BROKER" in
|
|
||||||
kafka)
|
|
||||||
: "${KAFKA_HOME:?set KAFKA_HOME}"
|
|
||||||
# Wipe whatever log.dirs points at, so the format below is not refused with
|
|
||||||
# "Invalid cluster.id ... Expected X, but read Y" from a previous run.
|
|
||||||
KAFKA_LOG_DIRS=$(grep -E '^log.dirs=' "$KAFKA_HOME/config/server.properties" | cut -d= -f2)
|
|
||||||
rm -rf ${KAFKA_LOG_DIRS//,/ }
|
|
||||||
"$KAFKA_HOME/bin/kafka-storage.sh" format --standalone -t "$("$KAFKA_HOME/bin/kafka-storage.sh" random-uuid)" \
|
|
||||||
-c "$KAFKA_HOME/config/server.properties" --ignore-formatted >/dev/null
|
|
||||||
start=$(date +%s%3N) # after formatting: kafka-storage.sh is two more JVM starts and is a
|
|
||||||
# one-off, not part of what a restart costs
|
|
||||||
setsid nohup "$KAFKA_HOME/bin/kafka-server-start.sh" "$KAFKA_HOME/config/server.properties" \
|
|
||||||
> /tmp/kafka-start.log 2>&1 < /dev/null &
|
|
||||||
for _ in $(seq 1 90); do port_open 9092 && break; sleep 1; done
|
|
||||||
ready=$(( $(date +%s%3N) - start ))
|
|
||||||
marker="kafka.Kafka"; ports="9092 (broker), 9093 (controller)"; conf="$KAFKA_HOME/config/server.properties"
|
|
||||||
dist="$KAFKA_HOME"
|
|
||||||
;;
|
|
||||||
rabbit)
|
|
||||||
: "${RABBITMQ_HOME:?set RABBITMQ_HOME}"; : "${ERL_ROOT:?set ERL_ROOT}"
|
|
||||||
scripts/rabbit-broker.sh >/dev/null
|
|
||||||
ready=$(( $(date +%s%3N) - start ))
|
|
||||||
# The broker is an Erlang VM, so the process to measure is beam.smp, not anything named
|
|
||||||
# "rabbitmq". epmd is a second, tiny process and is counted separately below.
|
|
||||||
marker="beam.smp"; ports="5672 (AMQP), 4369 (epmd), 25672 (inter-node)"
|
|
||||||
conf="$RABBITMQ_HOME/etc/rabbitmq/rabbitmq.conf.example"; dist="$RABBITMQ_HOME"
|
|
||||||
extra="epmd processes : $(procs epmd)
|
|
||||||
NOTE: RabbitMQ ships no configuration file at all. rabbitmq.conf.example is
|
|
||||||
entirely commented out and the boot log records \"Config file(s): (none)\".
|
|
||||||
Everything in that file is documentation, not a default."
|
|
||||||
;;
|
|
||||||
pulsar)
|
|
||||||
: "${PULSAR_HOME:?set PULSAR_HOME}"
|
|
||||||
scripts/pulsar-broker.sh >/dev/null
|
|
||||||
ready=$(( $(date +%s%3N) - start ))
|
|
||||||
marker="PulsarStandaloneStarter"; ports="6650 (binary), 8080 (admin/REST), 2181 (ZooKeeper)"
|
|
||||||
conf="$PULSAR_HOME/conf/standalone.conf"; dist="$PULSAR_HOME"
|
|
||||||
;;
|
|
||||||
*) echo "unknown broker $BROKER" >&2; exit 2 ;;
|
|
||||||
esac
|
|
||||||
|
|
||||||
cat <<EOF
|
|
||||||
broker : $BROKER
|
|
||||||
time from launch to first
|
|
||||||
accepted connection : ${ready} ms
|
|
||||||
resident memory at idle : $(rss_mb "$marker") MB
|
|
||||||
server processes : $(procs "$marker")
|
|
||||||
listening ports : $ports
|
|
||||||
unpacked distribution size : $(du -sm "$dist" | cut -f1) MB
|
|
||||||
settings in shipped default
|
|
||||||
configuration file : $(settings "$conf") ($(basename "$conf"))
|
|
||||||
${extra:-}
|
|
||||||
EOF
|
|
||||||
@@ -1,25 +0,0 @@
|
|||||||
#!/usr/bin/env bash
|
|
||||||
# Start a real Apache Pulsar broker with no Docker and no root.
|
|
||||||
#
|
|
||||||
# The standalone distribution bundles ZooKeeper, BookKeeper and the broker in one process, which
|
|
||||||
# is why a single tarball is enough. -nss skips the stream storage (BookKeeper's table service,
|
|
||||||
# needed only by stateful Pulsar Functions) and -nfw skips the functions worker; both save
|
|
||||||
# several hundred megabytes of heap and a few seconds of startup.
|
|
||||||
#
|
|
||||||
# PULSAR_HOME must point at an unpacked apache-pulsar-<version>-bin directory. Pulsar 4.2
|
|
||||||
# supports Java 17 and 21; the transcripts here were produced on Temurin 21.
|
|
||||||
set -eu
|
|
||||||
: "${PULSAR_HOME:?set PULSAR_HOME to an unpacked apache-pulsar-*-bin directory}"
|
|
||||||
export PULSAR_MEM="${PULSAR_MEM:--Xms384m -Xmx700m -XX:MaxDirectMemorySize=384m}"
|
|
||||||
|
|
||||||
cd "$PULSAR_HOME"
|
|
||||||
setsid nohup bin/pulsar standalone -nss -nfw > /tmp/pulsar-start.log 2>&1 < /dev/null &
|
|
||||||
for _ in $(seq 1 90); do
|
|
||||||
if [ "$(curl -s -o /dev/null -w '%{http_code}' http://127.0.0.1:8080/admin/v2/brokers/version)" = "200" ]; then
|
|
||||||
echo "broker ready on 6650 (admin 8080), version $(curl -s http://127.0.0.1:8080/admin/v2/brokers/version)"
|
|
||||||
exit 0
|
|
||||||
fi
|
|
||||||
sleep 2
|
|
||||||
done
|
|
||||||
echo "broker did not start; see /tmp/pulsar-start.log" >&2
|
|
||||||
exit 1
|
|
||||||
@@ -1,31 +0,0 @@
|
|||||||
#!/usr/bin/env bash
|
|
||||||
# Start a real RabbitMQ broker with no Docker and no root.
|
|
||||||
#
|
|
||||||
# RabbitMQ is an Erlang application, so an Erlang runtime and epmd on the path are the whole
|
|
||||||
# dependency. Point ERL_ROOT at an Erlang installation and RABBITMQ_HOME at an unpacked
|
|
||||||
# rabbitmq-server-generic-unix tarball. Ubuntu 22.04 ships Erlang 24, whose newest compatible
|
|
||||||
# broker is 3.10.25; 3.11 and later need Erlang 25.
|
|
||||||
#
|
|
||||||
# The one non-obvious step is starting epmd yourself: rabbitmq-server's own attempt to start it
|
|
||||||
# fails in a container without a resolvable hostname, and the failure surfaces as a forty-line
|
|
||||||
# Erlang crash dump ending in {'EXIT',nodistribution}.
|
|
||||||
set -eu
|
|
||||||
: "${ERL_ROOT:?set ERL_ROOT to an Erlang installation directory}"
|
|
||||||
: "${RABBITMQ_HOME:?set RABBITMQ_HOME to an unpacked rabbitmq-server-generic-unix directory}"
|
|
||||||
|
|
||||||
export PATH="$ERL_ROOT/bin:$ERL_ROOT/erts-"*/bin":$PATH"
|
|
||||||
export RABBITMQ_MNESIA_BASE="${RABBITMQ_MNESIA_BASE:-/tmp/rmq/data}"
|
|
||||||
export RABBITMQ_LOG_BASE="${RABBITMQ_LOG_BASE:-/tmp/rmq/log}"
|
|
||||||
export RABBITMQ_NODENAME="${RABBITMQ_NODENAME:-rabbit@localhost}"
|
|
||||||
export HOME="${HOME:-/tmp/rmq}"
|
|
||||||
mkdir -p "$RABBITMQ_MNESIA_BASE" "$RABBITMQ_LOG_BASE"
|
|
||||||
|
|
||||||
epmd -daemon 2>/dev/null || true
|
|
||||||
sleep 1
|
|
||||||
setsid nohup "$RABBITMQ_HOME/sbin/rabbitmq-server" > /tmp/rmq/boot.log 2>&1 < /dev/null &
|
|
||||||
for _ in $(seq 1 60); do
|
|
||||||
if (echo > /dev/tcp/127.0.0.1/5672) 2>/dev/null; then echo "broker ready on 5672"; exit 0; fi
|
|
||||||
sleep 1
|
|
||||||
done
|
|
||||||
echo "broker did not start; see /tmp/rmq/boot.log" >&2
|
|
||||||
exit 1
|
|
||||||
@@ -1,37 +0,0 @@
|
|||||||
#!/usr/bin/env bash
|
|
||||||
# Regenerates every file under docs/output/. Each broker is started, measured and left running
|
|
||||||
# only for its own group of tests, because three brokers do not fit comfortably in memory on a
|
|
||||||
# small machine at the same time.
|
|
||||||
#
|
|
||||||
# Required environment:
|
|
||||||
# KAFKA_HOME unpacked kafka_2.13-4.2.1
|
|
||||||
# ERL_ROOT an Erlang 24 installation (RabbitMQ 3.10.x)
|
|
||||||
# RABBITMQ_HOME unpacked rabbitmq-server-generic-unix-3.10.25
|
|
||||||
# PULSAR_HOME unpacked apache-pulsar-4.2.4-bin
|
|
||||||
set -euo pipefail
|
|
||||||
cd "$(dirname "$0")/.."
|
|
||||||
mkdir -p docs/output
|
|
||||||
|
|
||||||
# Kafka needs no external broker for the measurements: spring-kafka-test starts a real KRaft
|
|
||||||
# broker in-process. KAFKA_HOME is only used by the footprint measurement.
|
|
||||||
mvn -B -Dgroups=kafka test 2>&1 | grep -E 'Running |Tests run:|BUILD ' > docs/output/tests.txt
|
|
||||||
|
|
||||||
scripts/rabbit-broker.sh
|
|
||||||
mvn -B -Dgroups=rabbit test 2>&1 | grep -E 'Running |Tests run:|BUILD ' >> docs/output/tests.txt
|
|
||||||
|
|
||||||
scripts/pulsar-broker.sh
|
|
||||||
mvn -B -Dgroups=pulsar test 2>&1 | grep -E 'Running |Tests run:|BUILD ' >> docs/output/tests.txt
|
|
||||||
|
|
||||||
{
|
|
||||||
echo "== Operational footprint, measured on one 2-core / 3.8 GB Linux box =="
|
|
||||||
echo
|
|
||||||
echo "Each broker started from its shipped distribution with default configuration and a"
|
|
||||||
echo "700 MB heap cap, on Temurin JDK 21. Times are one sample on a small box: treat them as"
|
|
||||||
echo "orders of magnitude, not as a benchmark."
|
|
||||||
echo
|
|
||||||
scripts/footprint.sh kafka; echo
|
|
||||||
scripts/footprint.sh rabbit; echo
|
|
||||||
scripts/footprint.sh pulsar
|
|
||||||
} > docs/output/footprint.txt
|
|
||||||
|
|
||||||
echo "regenerated:"; ls -1 docs/output
|
|
||||||
@@ -1,15 +0,0 @@
|
|||||||
package com.ankurm.brokers;
|
|
||||||
|
|
||||||
/**
|
|
||||||
* There is deliberately no Spring application here. The measurements in src/test drive the three
|
|
||||||
* client libraries directly so that the thing being compared is the broker's delivery model and
|
|
||||||
* not three different sets of Spring defaults.
|
|
||||||
*
|
|
||||||
* <p>What each Spring Boot starter adds on top — and what it silently does not add, which
|
|
||||||
* in Boot 4 is more than people expect — is in docs/06-what-spring-adds.md.
|
|
||||||
*/
|
|
||||||
public final class BrokerComparisonApplication {
|
|
||||||
|
|
||||||
private BrokerComparisonApplication() {
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,23 +0,0 @@
|
|||||||
package com.ankurm.brokers;
|
|
||||||
|
|
||||||
import java.io.IOException;
|
|
||||||
import java.nio.file.Files;
|
|
||||||
import java.nio.file.Path;
|
|
||||||
|
|
||||||
/** Writes a transcript under docs/output/. Nothing in the article is typed by hand. */
|
|
||||||
public final class Capture {
|
|
||||||
|
|
||||||
private Capture() {
|
|
||||||
}
|
|
||||||
|
|
||||||
public static void write(String fileName, String heading, String body) {
|
|
||||||
Path dir = Path.of(System.getProperty("user.dir"), "docs", "output");
|
|
||||||
try {
|
|
||||||
Files.createDirectories(dir);
|
|
||||||
Files.writeString(dir.resolve(fileName), "== " + heading + " ==\n\n" + body + "\n");
|
|
||||||
}
|
|
||||||
catch (IOException ex) {
|
|
||||||
throw new IllegalStateException("could not write " + fileName, ex);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,232 +0,0 @@
|
|||||||
package com.ankurm.brokers;
|
|
||||||
|
|
||||||
import java.time.Duration;
|
|
||||||
import java.util.ArrayList;
|
|
||||||
import java.util.LinkedHashMap;
|
|
||||||
import java.util.List;
|
|
||||||
import java.util.Map;
|
|
||||||
import java.util.Properties;
|
|
||||||
import java.util.Set;
|
|
||||||
|
|
||||||
import org.apache.kafka.clients.consumer.ConsumerConfig;
|
|
||||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
|
||||||
import org.apache.kafka.clients.consumer.ConsumerRecords;
|
|
||||||
import org.apache.kafka.clients.consumer.KafkaConsumer;
|
|
||||||
import org.apache.kafka.clients.producer.KafkaProducer;
|
|
||||||
import org.apache.kafka.clients.producer.ProducerConfig;
|
|
||||||
import org.apache.kafka.clients.producer.ProducerRecord;
|
|
||||||
import org.apache.kafka.common.TopicPartition;
|
|
||||||
import org.apache.kafka.common.serialization.StringDeserializer;
|
|
||||||
import org.apache.kafka.common.serialization.StringSerializer;
|
|
||||||
import org.junit.jupiter.api.AfterAll;
|
|
||||||
import org.junit.jupiter.api.BeforeAll;
|
|
||||||
import org.junit.jupiter.api.Tag;
|
|
||||||
import org.junit.jupiter.api.Test;
|
|
||||||
|
|
||||||
import org.springframework.kafka.test.EmbeddedKafkaKraftBroker;
|
|
||||||
|
|
||||||
import static org.assertj.core.api.Assertions.assertThat;
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Kafka measured against a real broker: an in-process KRaft broker from spring-kafka-test, one
|
|
||||||
* node, one topic, three partitions.
|
|
||||||
*
|
|
||||||
* <p>The three questions are the same for all three brokers: what ordering survives, what can be
|
|
||||||
* replayed, and how far consumers scale.
|
|
||||||
*/
|
|
||||||
@Tag("kafka")
|
|
||||||
class KafkaComparisonTest {
|
|
||||||
|
|
||||||
private static final String ORDERING = "orders-ordering";
|
|
||||||
|
|
||||||
private static final String REPLAY = "orders-replay";
|
|
||||||
|
|
||||||
private static final String SCALING = "orders-scaling";
|
|
||||||
|
|
||||||
private static EmbeddedKafkaKraftBroker broker;
|
|
||||||
|
|
||||||
@BeforeAll
|
|
||||||
static void startBroker() {
|
|
||||||
broker = new EmbeddedKafkaKraftBroker(1, 3, ORDERING, REPLAY, SCALING);
|
|
||||||
broker.afterPropertiesSet();
|
|
||||||
}
|
|
||||||
|
|
||||||
@AfterAll
|
|
||||||
static void stopBroker() {
|
|
||||||
broker.destroy();
|
|
||||||
}
|
|
||||||
|
|
||||||
// ---------------------------------------------------------------- ordering
|
|
||||||
|
|
||||||
@Test
|
|
||||||
void orderingIsPerPartitionOnly() {
|
|
||||||
produce(ORDERING, 12);
|
|
||||||
|
|
||||||
List<ConsumerRecord<String, String>> received = consumeAll(ORDERING, "ordering-group", 12);
|
|
||||||
|
|
||||||
Map<String, List<String>> perKey = new LinkedHashMap<>();
|
|
||||||
Map<String, Set<Integer>> partitionsPerKey = new LinkedHashMap<>();
|
|
||||||
List<String> globalOrder = new ArrayList<>();
|
|
||||||
for (ConsumerRecord<String, String> record : received) {
|
|
||||||
perKey.computeIfAbsent(record.key(), (k) -> new ArrayList<>()).add(record.value());
|
|
||||||
partitionsPerKey.computeIfAbsent(record.key(), (k) -> new java.util.TreeSet<>())
|
|
||||||
.add(record.partition());
|
|
||||||
globalOrder.add(record.key() + "=" + record.value() + "@p" + record.partition());
|
|
||||||
}
|
|
||||||
|
|
||||||
// Per key, the sequence numbers must come back in the order they were produced. Compared
|
|
||||||
// as numbers: as strings, "10" sorts before "2" and the assertion is meaningless.
|
|
||||||
perKey.forEach((key, values) -> assertThat(values.stream().map(Integer::parseInt).toList())
|
|
||||||
.isSorted());
|
|
||||||
// Each key went to exactly one partition, which is why that holds...
|
|
||||||
partitionsPerKey.forEach((key, partitions) -> assertThat(partitions).hasSize(1));
|
|
||||||
// ...and the three keys are spread over all three partitions, so the transcript below is
|
|
||||||
// actually testing something.
|
|
||||||
assertThat(partitionsPerKey.values().stream().flatMap(Set::stream).distinct().count())
|
|
||||||
.isEqualTo(3);
|
|
||||||
List<Integer> asDelivered = globalOrder.stream()
|
|
||||||
.map((entry) -> Integer.parseInt(entry.split("=")[1].split("@")[0])).toList();
|
|
||||||
assertThat(asDelivered).isNotEqualTo(
|
|
||||||
java.util.stream.IntStream.rangeClosed(1, 12).boxed().toList());
|
|
||||||
|
|
||||||
StringBuilder body = new StringBuilder();
|
|
||||||
body.append("produced 12 records, keys D/A/F round-robin, values 1..12 in order\n")
|
|
||||||
.append("(keys chosen so that murmur2 spreads them: D->0, A->1, F->2. A, B and C\n")
|
|
||||||
.append(" all hash to partition 1 with three partitions, which is worth knowing\n")
|
|
||||||
.append(" before you decide your keys are well distributed.)\n\n");
|
|
||||||
body.append("consumed in this order:\n ").append(String.join("\n ", globalOrder))
|
|
||||||
.append("\n\nper key:\n");
|
|
||||||
perKey.forEach((key, values) -> body.append(" ").append(key).append(" -> partition ")
|
|
||||||
.append(partitionsPerKey.get(key)).append(" values ").append(values).append('\n'));
|
|
||||||
body.append("\nOrder is preserved within each key because a key hashes to one partition.\n")
|
|
||||||
.append("Across keys it is not: the values above are not 1..12 in order, because a\n")
|
|
||||||
.append("consumer drains one partition's buffer before the next.\n");
|
|
||||||
Capture.write("kafka-ordering.txt", "Kafka: ordering is per partition", body.toString());
|
|
||||||
}
|
|
||||||
|
|
||||||
// ---------------------------------------------------------------- replay
|
|
||||||
|
|
||||||
@Test
|
|
||||||
void anythingStillWithinRetentionCanBeReadAgain() {
|
|
||||||
produce(REPLAY, 12);
|
|
||||||
|
|
||||||
int first = consumeAll(REPLAY, "replay-group-1", 12).size();
|
|
||||||
int again = consumeAll(REPLAY, "replay-group-2", 12).size();
|
|
||||||
|
|
||||||
List<ConsumerRecord<String, String>> third;
|
|
||||||
try (KafkaConsumer<String, String> consumer = consumer("replay-group-1")) {
|
|
||||||
consumer.subscribe(List.of(REPLAY));
|
|
||||||
consumer.poll(Duration.ofSeconds(2));
|
|
||||||
consumer.seekToBeginning(consumer.assignment());
|
|
||||||
third = drain(consumer, 12);
|
|
||||||
}
|
|
||||||
|
|
||||||
assertThat(first).isEqualTo(12);
|
|
||||||
assertThat(again).isEqualTo(12);
|
|
||||||
assertThat(third).hasSize(12);
|
|
||||||
|
|
||||||
Capture.write("kafka-replay.txt", "Kafka: the log is the storage",
|
|
||||||
"""
|
|
||||||
group replay-group-1, first read : %d records
|
|
||||||
group replay-group-2, brand new group : %d records
|
|
||||||
group replay-group-1, after seekToBeginning: %d records
|
|
||||||
|
|
||||||
Consuming does not remove anything. A consumer group is a cursor over a log
|
|
||||||
that the broker keeps until retention expires, so a new group, a reset
|
|
||||||
offset or a seek all read the same records again.
|
|
||||||
""".formatted(first, again, third.size()));
|
|
||||||
}
|
|
||||||
|
|
||||||
// ---------------------------------------------------------------- consumer scaling
|
|
||||||
|
|
||||||
@Test
|
|
||||||
void consumerParallelismIsCappedByPartitionCount() throws Exception {
|
|
||||||
produce(SCALING, 12);
|
|
||||||
|
|
||||||
List<KafkaConsumer<String, String>> consumers = new ArrayList<>();
|
|
||||||
Map<String, Set<Integer>> assignment = new LinkedHashMap<>();
|
|
||||||
try {
|
|
||||||
for (int i = 1; i <= 5; i++) {
|
|
||||||
KafkaConsumer<String, String> consumer = consumer("scaling-group");
|
|
||||||
consumer.subscribe(List.of(SCALING));
|
|
||||||
consumers.add(consumer);
|
|
||||||
}
|
|
||||||
// Poll each consumer until the group has settled and every member knows its share.
|
|
||||||
for (int round = 0; round < 8; round++) {
|
|
||||||
for (KafkaConsumer<String, String> consumer : consumers) {
|
|
||||||
consumer.poll(Duration.ofMillis(500));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
for (int i = 0; i < consumers.size(); i++) {
|
|
||||||
Set<Integer> partitions = new java.util.TreeSet<>();
|
|
||||||
for (TopicPartition tp : consumers.get(i).assignment()) {
|
|
||||||
partitions.add(tp.partition());
|
|
||||||
}
|
|
||||||
assignment.put("consumer-" + (i + 1), partitions);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
finally {
|
|
||||||
consumers.forEach(KafkaConsumer::close);
|
|
||||||
}
|
|
||||||
|
|
||||||
long idle = assignment.values().stream().filter(Set::isEmpty).count();
|
|
||||||
assertThat(assignment).hasSize(5);
|
|
||||||
assertThat(idle).isEqualTo(2);
|
|
||||||
|
|
||||||
StringBuilder body = new StringBuilder("topic 'orders-scaling', 3 partitions, 5 consumers in one group\n\n");
|
|
||||||
assignment.forEach((name, partitions) -> body.append(" ").append(name).append(" -> ")
|
|
||||||
.append(partitions.isEmpty() ? "no partitions (idle)" : "partitions " + partitions)
|
|
||||||
.append('\n'));
|
|
||||||
body.append("\nconsumers with no partitions: ").append(idle).append('\n')
|
|
||||||
.append("\nA partition is assigned to at most one consumer in a group, so the number\n")
|
|
||||||
.append("of partitions is a hard ceiling on consumer parallelism. Adding consumers\n")
|
|
||||||
.append("beyond it adds idle processes, not throughput.\n");
|
|
||||||
Capture.write("kafka-consumer-scaling.txt",
|
|
||||||
"Kafka: partitions are the ceiling on consumer parallelism", body.toString());
|
|
||||||
}
|
|
||||||
|
|
||||||
// ---------------------------------------------------------------- helpers
|
|
||||||
|
|
||||||
private void produce(String topic, int count) {
|
|
||||||
Properties props = new Properties();
|
|
||||||
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, broker.getBrokersAsString());
|
|
||||||
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
|
|
||||||
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
|
|
||||||
// D, A and F, not A, B and C: with three partitions murmur2 maps A, B and C all to
|
|
||||||
// partition 1, so the obvious choice of keys would have produced a transcript in which
|
|
||||||
// global order happened to be preserved and proved nothing. D->0, A->1, F->2.
|
|
||||||
String[] keys = { "D", "A", "F" };
|
|
||||||
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
|
|
||||||
for (int i = 1; i <= count; i++) {
|
|
||||||
producer.send(new ProducerRecord<>(topic, keys[i % keys.length], String.valueOf(i)));
|
|
||||||
}
|
|
||||||
producer.flush();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private KafkaConsumer<String, String> consumer(String group) {
|
|
||||||
Properties props = new Properties();
|
|
||||||
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, broker.getBrokersAsString());
|
|
||||||
props.put(ConsumerConfig.GROUP_ID_CONFIG, group);
|
|
||||||
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
|
|
||||||
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
|
|
||||||
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
|
|
||||||
return new KafkaConsumer<>(props);
|
|
||||||
}
|
|
||||||
|
|
||||||
private List<ConsumerRecord<String, String>> consumeAll(String topic, String group, int expected) {
|
|
||||||
try (KafkaConsumer<String, String> consumer = consumer(group)) {
|
|
||||||
consumer.subscribe(List.of(topic));
|
|
||||||
return drain(consumer, expected);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private List<ConsumerRecord<String, String>> drain(KafkaConsumer<String, String> consumer, int expected) {
|
|
||||||
List<ConsumerRecord<String, String>> received = new ArrayList<>();
|
|
||||||
for (int i = 0; i < 20 && received.size() < expected; i++) {
|
|
||||||
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
|
|
||||||
records.forEach(received::add);
|
|
||||||
}
|
|
||||||
return received;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,298 +0,0 @@
|
|||||||
package com.ankurm.brokers;
|
|
||||||
|
|
||||||
import java.util.ArrayList;
|
|
||||||
import java.util.LinkedHashMap;
|
|
||||||
import java.util.List;
|
|
||||||
import java.util.Map;
|
|
||||||
import java.util.Set;
|
|
||||||
import java.util.concurrent.TimeUnit;
|
|
||||||
|
|
||||||
import org.apache.pulsar.client.api.Consumer;
|
|
||||||
import org.apache.pulsar.client.api.Message;
|
|
||||||
import org.apache.pulsar.client.api.MessageId;
|
|
||||||
import org.apache.pulsar.client.api.Producer;
|
|
||||||
import org.apache.pulsar.client.api.PulsarClient;
|
|
||||||
import org.apache.pulsar.client.api.Schema;
|
|
||||||
import org.apache.pulsar.client.api.SubscriptionInitialPosition;
|
|
||||||
import org.apache.pulsar.client.api.SubscriptionType;
|
|
||||||
import org.junit.jupiter.api.AfterAll;
|
|
||||||
import org.junit.jupiter.api.BeforeAll;
|
|
||||||
import org.junit.jupiter.api.Tag;
|
|
||||||
import org.junit.jupiter.api.Test;
|
|
||||||
|
|
||||||
import static org.assertj.core.api.Assertions.assertThat;
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Pulsar measured against a real standalone broker on localhost:6650, started by
|
|
||||||
* scripts/pulsar-broker.sh.
|
|
||||||
*
|
|
||||||
* <p>Pulsar is the interesting third option because it does not force the trade the other two
|
|
||||||
* force. The subscription type is chosen per subscriber rather than being a property of the
|
|
||||||
* topic, so ordering and consumer fan-out are separate decisions.
|
|
||||||
*/
|
|
||||||
@Tag("pulsar")
|
|
||||||
class PulsarComparisonTest {
|
|
||||||
|
|
||||||
private static PulsarClient client;
|
|
||||||
|
|
||||||
@BeforeAll
|
|
||||||
static void connect() throws Exception {
|
|
||||||
client = PulsarClient.builder().serviceUrl("pulsar://127.0.0.1:6650").build();
|
|
||||||
}
|
|
||||||
|
|
||||||
@AfterAll
|
|
||||||
static void close() throws Exception {
|
|
||||||
client.close();
|
|
||||||
}
|
|
||||||
|
|
||||||
// ---------------------------------------------------------------- ordering
|
|
||||||
|
|
||||||
@Test
|
|
||||||
void sharedSpreadsAKeyAndKeySharedPinsIt() throws Exception {
|
|
||||||
String topic = "persistent://public/default/ordering-" + System.nanoTime();
|
|
||||||
produce(topic, 12);
|
|
||||||
|
|
||||||
Map<String, Set<String>> shared = consumersPerKey(topic, "sub-shared", SubscriptionType.Shared, 12);
|
|
||||||
Map<String, Set<String>> keyShared = consumersPerKey(topic, "sub-keyshared",
|
|
||||||
SubscriptionType.Key_Shared, 12);
|
|
||||||
|
|
||||||
// Key_Shared: every key was handled by exactly one consumer, so per-key order is safe.
|
|
||||||
keyShared.forEach((key, handlers) -> assertThat(handlers).hasSize(1));
|
|
||||||
// Shared: at least one key was handled by both, so two messages with that key were in
|
|
||||||
// flight on two consumers and nothing orders them.
|
|
||||||
assertThat(shared.values().stream().anyMatch((handlers) -> handlers.size() > 1)).isTrue();
|
|
||||||
|
|
||||||
Capture.write("pulsar-ordering.txt", "Pulsar: the subscription type decides",
|
|
||||||
"""
|
|
||||||
12 messages, keys A/B/C, values 1..12 in order, two consumers per
|
|
||||||
subscription, receiverQueueSize 1 so the first consumer cannot take the
|
|
||||||
whole backlog.
|
|
||||||
|
|
||||||
Shared subscription, which consumers saw each key:
|
|
||||||
%s
|
|
||||||
Key_Shared subscription, which consumers saw each key:
|
|
||||||
%s
|
|
||||||
A Shared subscription round-robins individual messages, so messages with
|
|
||||||
the same key end up on different consumers and can be processed at the
|
|
||||||
same time: there is no per-key order left to speak of. Key_Shared hashes
|
|
||||||
the key to one consumer, which is Kafka's guarantee -- except that the
|
|
||||||
assignment belongs to the subscription and is recomputed as consumers come
|
|
||||||
and go, rather than being fixed by a partition count chosen when the topic
|
|
||||||
was created.
|
|
||||||
""".formatted(render3(shared), render3(keyShared)));
|
|
||||||
}
|
|
||||||
|
|
||||||
// ---------------------------------------------------------------- replay
|
|
||||||
|
|
||||||
@Test
|
|
||||||
void aSubscriptionCanBeRewound() throws Exception {
|
|
||||||
String topic = "persistent://public/default/replay-" + System.nanoTime();
|
|
||||||
produce(topic, 12);
|
|
||||||
|
|
||||||
int first;
|
|
||||||
int afterSeek;
|
|
||||||
try (Consumer<String> consumer = consumer(topic, "replay-sub", SubscriptionType.Exclusive)) {
|
|
||||||
first = drain(consumer, 12).size();
|
|
||||||
consumer.seek(MessageId.earliest);
|
|
||||||
afterSeek = drain(consumer, 12).size();
|
|
||||||
}
|
|
||||||
|
|
||||||
int newSubscription;
|
|
||||||
try (Consumer<String> consumer = consumer(topic, "replay-sub-2", SubscriptionType.Exclusive)) {
|
|
||||||
newSubscription = drain(consumer, 12).size();
|
|
||||||
}
|
|
||||||
|
|
||||||
assertThat(first).isEqualTo(12);
|
|
||||||
assertThat(afterSeek).isEqualTo(12);
|
|
||||||
assertThat(newSubscription).isEqualTo(12);
|
|
||||||
|
|
||||||
Capture.write("pulsar-replay.txt", "Pulsar: acknowledged, but still there",
|
|
||||||
"""
|
|
||||||
subscription replay-sub, first read : %d messages
|
|
||||||
subscription replay-sub, after seek(earliest): %d messages
|
|
||||||
subscription replay-sub-2, brand new : %d messages
|
|
||||||
|
|
||||||
Acknowledgement moves a cursor; the message itself lives in the managed
|
|
||||||
ledger. seek(MessageId) and seek(timestamp) rewind a live subscription,
|
|
||||||
which Kafka can also do by resetting offsets. What differs is the default:
|
|
||||||
Pulsar deletes a message once every subscription has acknowledged it,
|
|
||||||
unless a retention policy on the namespace says otherwise, whereas Kafka
|
|
||||||
keeps it for the retention period regardless of who read it. A Pulsar
|
|
||||||
topic with no retention policy and no subscriptions keeps nothing.
|
|
||||||
""".formatted(first, afterSeek, newSubscription));
|
|
||||||
}
|
|
||||||
|
|
||||||
// ---------------------------------------------------------------- consumer scaling
|
|
||||||
|
|
||||||
@Test
|
|
||||||
void sharedSubscriptionsHaveNoPartitionCeilingButHaveAReceiverQueue() throws Exception {
|
|
||||||
String topicA = "persistent://public/default/scaling-default-" + System.nanoTime();
|
|
||||||
produce(topicA, 40);
|
|
||||||
Map<String, Integer> withDefaultQueue = drainAcross(topicA, "scaling-sub", 5, 40, 0);
|
|
||||||
|
|
||||||
String topicB = "persistent://public/default/scaling-q1-" + System.nanoTime();
|
|
||||||
produce(topicB, 40);
|
|
||||||
Map<String, Integer> withQueueOfOne = drainAcross(topicB, "scaling-sub", 5, 40, 1);
|
|
||||||
|
|
||||||
long idleDefault = withDefaultQueue.values().stream().filter((n) -> n == 0).count();
|
|
||||||
long idleTuned = withQueueOfOne.values().stream().filter((n) -> n == 0).count();
|
|
||||||
|
|
||||||
assertThat(idleTuned).isZero();
|
|
||||||
assertThat(idleDefault).isGreaterThan(0);
|
|
||||||
|
|
||||||
Capture.write("pulsar-consumer-scaling.txt",
|
|
||||||
"Pulsar: no partition ceiling, but the receiver queue decides who gets the work",
|
|
||||||
"""
|
|
||||||
one non-partitioned topic, 40 messages, 5 consumers, Shared subscription
|
|
||||||
|
|
||||||
receiverQueueSize left at the default (1000):
|
|
||||||
%s consumers that received nothing : %d
|
|
||||||
|
|
||||||
receiverQueueSize(1):
|
|
||||||
%s consumers that received nothing : %d
|
|
||||||
|
|
||||||
The topic has no partitions and five consumers can still share the work --
|
|
||||||
in Kafka the same shape needs at least five partitions, chosen when the
|
|
||||||
topic was created. But the default receiver queue is 1000 messages, so the
|
|
||||||
first consumers to connect pull the whole 40-message backlog into their own
|
|
||||||
buffers before the rest ask for anything, and the subscription looks
|
|
||||||
broken. Which consumers win is a race and varies between runs -- one
|
|
||||||
consumer taking all forty, or two taking twenty each -- but the consumers
|
|
||||||
that lose it see nothing at all. This is the same trap as RabbitMQ's
|
|
||||||
unbounded prefetch, with a different name and a much larger default.
|
|
||||||
""".formatted(render2(withDefaultQueue), idleDefault,
|
|
||||||
render2(withQueueOfOne), idleTuned));
|
|
||||||
}
|
|
||||||
|
|
||||||
// ---------------------------------------------------------------- helpers
|
|
||||||
|
|
||||||
private void produce(String topic, int count) throws Exception {
|
|
||||||
String[] keys = { "A", "B", "C" };
|
|
||||||
try (Producer<String> producer = client.newProducer(Schema.STRING).topic(topic).create()) {
|
|
||||||
for (int i = 1; i <= count; i++) {
|
|
||||||
producer.newMessage().key(keys[i % keys.length]).value(String.valueOf(i)).send();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private Consumer<String> consumer(String topic, String subscription, SubscriptionType type)
|
|
||||||
throws Exception {
|
|
||||||
return consumer(topic, subscription, type, 0);
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @param receiverQueueSize 0 leaves Pulsar's default of 1000 in place; anything else sets it.
|
|
||||||
* The default matters more than it looks: it is how many messages a single consumer will pull
|
|
||||||
* into its own buffer before another consumer on the same subscription gets a chance.
|
|
||||||
*/
|
|
||||||
private Consumer<String> consumer(String topic, String subscription, SubscriptionType type,
|
|
||||||
int receiverQueueSize) throws Exception {
|
|
||||||
var builder = client.newConsumer(Schema.STRING)
|
|
||||||
.topic(topic)
|
|
||||||
.subscriptionName(subscription)
|
|
||||||
.subscriptionType(type)
|
|
||||||
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest);
|
|
||||||
if (receiverQueueSize > 0) {
|
|
||||||
builder = builder.receiverQueueSize(receiverQueueSize);
|
|
||||||
}
|
|
||||||
return builder.subscribe();
|
|
||||||
}
|
|
||||||
|
|
||||||
private List<String> drain(Consumer<String> consumer, int expected) throws Exception {
|
|
||||||
List<String> received = new ArrayList<>();
|
|
||||||
while (received.size() < expected) {
|
|
||||||
Message<String> message = consumer.receive(2, TimeUnit.SECONDS);
|
|
||||||
if (message == null) {
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
received.add(message.getValue());
|
|
||||||
consumer.acknowledge(message);
|
|
||||||
}
|
|
||||||
return received;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Which consumers saw each key, with a receiver queue of one so nothing hoards. */
|
|
||||||
private Map<String, Set<String>> consumersPerKey(String topic, String subscription,
|
|
||||||
SubscriptionType type, int expected) throws Exception {
|
|
||||||
Map<String, Set<String>> perKey = new LinkedHashMap<>();
|
|
||||||
List<Consumer<String>> consumers = new ArrayList<>();
|
|
||||||
try {
|
|
||||||
consumers.add(consumer(topic, subscription, type, 1));
|
|
||||||
consumers.add(consumer(topic, subscription, type, 1));
|
|
||||||
int received = 0;
|
|
||||||
while (received < expected) {
|
|
||||||
boolean any = false;
|
|
||||||
for (int i = 0; i < consumers.size(); i++) {
|
|
||||||
Message<String> message = consumers.get(i).receive(500, TimeUnit.MILLISECONDS);
|
|
||||||
if (message != null) {
|
|
||||||
perKey.computeIfAbsent(message.getKey(), (k) -> new java.util.TreeSet<>())
|
|
||||||
.add("consumer-" + (i + 1));
|
|
||||||
consumers.get(i).acknowledge(message);
|
|
||||||
received++;
|
|
||||||
any = true;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if (!any) {
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
finally {
|
|
||||||
for (Consumer<String> consumer : consumers) {
|
|
||||||
consumer.close();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return perKey;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Drains a Shared subscription across {@code n} consumers and reports the distribution. */
|
|
||||||
private Map<String, Integer> drainAcross(String topic, String subscription, int n, int expected,
|
|
||||||
int receiverQueueSize) throws Exception {
|
|
||||||
Map<String, Integer> counts = new LinkedHashMap<>();
|
|
||||||
List<Consumer<String>> consumers = new ArrayList<>();
|
|
||||||
try {
|
|
||||||
for (int i = 1; i <= n; i++) {
|
|
||||||
consumers.add(consumer(topic, subscription, SubscriptionType.Shared, receiverQueueSize));
|
|
||||||
counts.put("consumer-" + i, 0);
|
|
||||||
}
|
|
||||||
int received = 0;
|
|
||||||
while (received < expected) {
|
|
||||||
boolean any = false;
|
|
||||||
for (int i = 0; i < consumers.size(); i++) {
|
|
||||||
Message<String> message = consumers.get(i).receive(300, TimeUnit.MILLISECONDS);
|
|
||||||
if (message != null) {
|
|
||||||
consumers.get(i).acknowledge(message);
|
|
||||||
counts.merge("consumer-" + (i + 1), 1, Integer::sum);
|
|
||||||
received++;
|
|
||||||
any = true;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if (!any) {
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
finally {
|
|
||||||
for (Consumer<String> consumer : consumers) {
|
|
||||||
consumer.close();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return counts;
|
|
||||||
}
|
|
||||||
|
|
||||||
private static String render3(Map<String, Set<String>> perKey) {
|
|
||||||
StringBuilder out = new StringBuilder();
|
|
||||||
perKey.entrySet().stream().sorted(Map.Entry.comparingByKey()).forEach((entry) ->
|
|
||||||
out.append(" ").append(entry.getKey()).append(" -> ").append(entry.getValue())
|
|
||||||
.append('\n'));
|
|
||||||
return out.toString();
|
|
||||||
}
|
|
||||||
|
|
||||||
private static String render2(Map<String, Integer> counts) {
|
|
||||||
StringBuilder out = new StringBuilder();
|
|
||||||
counts.entrySet().stream().sorted(Map.Entry.comparingByKey()).forEach((entry) ->
|
|
||||||
out.append(" ").append(entry.getKey()).append(" -> ").append(entry.getValue())
|
|
||||||
.append(" messages\n"));
|
|
||||||
return out.toString();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,286 +0,0 @@
|
|||||||
package com.ankurm.brokers;
|
|
||||||
|
|
||||||
import java.io.IOException;
|
|
||||||
import java.nio.charset.StandardCharsets;
|
|
||||||
import java.util.ArrayList;
|
|
||||||
import java.util.LinkedHashMap;
|
|
||||||
import java.util.List;
|
|
||||||
import java.util.Map;
|
|
||||||
import java.util.concurrent.ConcurrentHashMap;
|
|
||||||
import java.util.concurrent.CopyOnWriteArrayList;
|
|
||||||
import java.util.concurrent.CountDownLatch;
|
|
||||||
import java.util.concurrent.TimeUnit;
|
|
||||||
import java.util.concurrent.atomic.AtomicInteger;
|
|
||||||
|
|
||||||
import com.rabbitmq.client.AMQP;
|
|
||||||
import com.rabbitmq.client.Channel;
|
|
||||||
import com.rabbitmq.client.Connection;
|
|
||||||
import com.rabbitmq.client.ConnectionFactory;
|
|
||||||
import com.rabbitmq.client.DefaultConsumer;
|
|
||||||
import com.rabbitmq.client.Envelope;
|
|
||||||
import com.rabbitmq.client.GetResponse;
|
|
||||||
import org.junit.jupiter.api.AfterAll;
|
|
||||||
import org.junit.jupiter.api.BeforeAll;
|
|
||||||
import org.junit.jupiter.api.Tag;
|
|
||||||
import org.junit.jupiter.api.Test;
|
|
||||||
|
|
||||||
import static org.assertj.core.api.Assertions.assertThat;
|
|
||||||
|
|
||||||
/**
|
|
||||||
* RabbitMQ measured against a real broker on localhost:5672, started by
|
|
||||||
* scripts/rabbit-broker.sh.
|
|
||||||
*
|
|
||||||
* <p>The same three questions as the Kafka measurements, against a broker that owns the message
|
|
||||||
* until it is acknowledged rather than a log that keeps it.
|
|
||||||
*/
|
|
||||||
@Tag("rabbit")
|
|
||||||
class RabbitComparisonTest {
|
|
||||||
|
|
||||||
private static Connection connection;
|
|
||||||
|
|
||||||
@BeforeAll
|
|
||||||
static void connect() throws Exception {
|
|
||||||
ConnectionFactory factory = new ConnectionFactory();
|
|
||||||
factory.setHost("127.0.0.1");
|
|
||||||
factory.setPort(5672);
|
|
||||||
connection = factory.newConnection();
|
|
||||||
}
|
|
||||||
|
|
||||||
@AfterAll
|
|
||||||
static void close() throws Exception {
|
|
||||||
connection.close();
|
|
||||||
}
|
|
||||||
|
|
||||||
// ---------------------------------------------------------------- ordering
|
|
||||||
|
|
||||||
@Test
|
|
||||||
void oneConsumerIsFifoAndTwoConsumersAreNot() throws Exception {
|
|
||||||
String single = declare("order-single");
|
|
||||||
publish(single, 12);
|
|
||||||
List<String> fifo = getAll(single, 12);
|
|
||||||
|
|
||||||
String shared = declare("order-shared");
|
|
||||||
publish(shared, 12);
|
|
||||||
List<String> completionOrder = consumeWithTwoConsumers(shared, 12);
|
|
||||||
|
|
||||||
assertThat(fifo).containsExactly("1", "2", "3", "4", "5", "6", "7", "8", "9", "10", "11", "12");
|
|
||||||
|
|
||||||
int inversions = 0;
|
|
||||||
for (int i = 1; i < completionOrder.size(); i++) {
|
|
||||||
if (Integer.parseInt(completionOrder.get(i)) < Integer.parseInt(completionOrder.get(i - 1))) {
|
|
||||||
inversions++;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
Capture.write("rabbit-ordering.txt", "RabbitMQ: FIFO per queue, per consumer",
|
|
||||||
"""
|
|
||||||
one queue, one consumer, 12 messages
|
|
||||||
completion order : %s
|
|
||||||
|
|
||||||
one queue, two consumers, prefetch 1, consumer-1 slower than consumer-2
|
|
||||||
completion order : %s
|
|
||||||
out-of-order steps: %d
|
|
||||||
|
|
||||||
A queue is FIFO and a single consumer sees it that way. The moment a second
|
|
||||||
consumer is added the broker hands the next message to whichever consumer
|
|
||||||
is free, so the order in which work finishes is no longer the order in
|
|
||||||
which it was published. There is no key: RabbitMQ has no notion of a
|
|
||||||
partition to which related messages could be pinned. Ordering across
|
|
||||||
related messages means one queue and one consumer, and therefore no
|
|
||||||
horizontal scaling for that queue -- or a consistent-hash exchange, which
|
|
||||||
is a plugin.
|
|
||||||
""".formatted(fifo, completionOrder, inversions));
|
|
||||||
}
|
|
||||||
|
|
||||||
// ---------------------------------------------------------------- replay
|
|
||||||
|
|
||||||
@Test
|
|
||||||
void anAcknowledgedMessageIsGone() throws Exception {
|
|
||||||
String queue = declare("replay-probe");
|
|
||||||
publish(queue, 12);
|
|
||||||
|
|
||||||
List<String> first = getAll(queue, 12);
|
|
||||||
List<String> second = getAll(queue, 12);
|
|
||||||
|
|
||||||
int depth = depth(queue);
|
|
||||||
|
|
||||||
assertThat(first).hasSize(12);
|
|
||||||
assertThat(second).isEmpty();
|
|
||||||
assertThat(depth).isZero();
|
|
||||||
|
|
||||||
Capture.write("rabbit-replay.txt", "RabbitMQ: there is nothing to replay",
|
|
||||||
"""
|
|
||||||
first drain of the queue : %d messages
|
|
||||||
second drain of the queue : %d messages
|
|
||||||
queue depth afterwards : %d
|
|
||||||
|
|
||||||
Acknowledging a message deletes it. The broker is a router with buffers,
|
|
||||||
not a log: there is no offset to rewind and no second reader that can see
|
|
||||||
what the first one consumed. Reading the same message twice means
|
|
||||||
arranging it in advance -- a second queue bound to the same exchange, or a
|
|
||||||
copy written somewhere else -- and it cannot be arranged after the fact.
|
|
||||||
""".formatted(first.size(), second.size(), depth));
|
|
||||||
}
|
|
||||||
|
|
||||||
// ---------------------------------------------------------------- consumer scaling
|
|
||||||
|
|
||||||
@Test
|
|
||||||
void consumersScaleWithoutAPartitionCeilingButPrefetchDecidesWho() throws Exception {
|
|
||||||
String unlimited = declare("scaling-unlimited");
|
|
||||||
publish(unlimited, 40);
|
|
||||||
Map<String, Integer> withDefaultPrefetch = drainAcross(unlimited, 5, 40, 0);
|
|
||||||
|
|
||||||
String throttled = declare("scaling-prefetch-1");
|
|
||||||
publish(throttled, 40);
|
|
||||||
Map<String, Integer> withPrefetchOne = drainAcross(throttled, 5, 40, 1);
|
|
||||||
|
|
||||||
long idleDefault = withDefaultPrefetch.values().stream().filter((n) -> n == 0).count();
|
|
||||||
long idleTuned = withPrefetchOne.values().stream().filter((n) -> n == 0).count();
|
|
||||||
|
|
||||||
assertThat(idleTuned).isZero();
|
|
||||||
|
|
||||||
Capture.write("rabbit-consumer-scaling.txt",
|
|
||||||
"RabbitMQ: consumers scale, and prefetch decides whether they actually do",
|
|
||||||
"""
|
|
||||||
one queue, 40 messages, 5 consumers, each taking 10 ms per message
|
|
||||||
|
|
||||||
no basicQos at all (unlimited prefetch, the AMQP default):
|
|
||||||
%s consumers that received nothing : %d
|
|
||||||
|
|
||||||
basicQos(1):
|
|
||||||
%s consumers that received nothing : %d
|
|
||||||
|
|
||||||
Every consumer on a queue competes for the same messages, so adding
|
|
||||||
consumers adds throughput and there is no structural ceiling of the kind
|
|
||||||
Kafka's partition count imposes. But with the AMQP default the broker
|
|
||||||
pushes as many messages as a consumer will take, so whichever consumer
|
|
||||||
connects first can be handed the entire backlog while the others sit idle.
|
|
||||||
basicQos is not a tuning knob you get to postpone; it is what makes the
|
|
||||||
fan-out real. Spring AMQP sets it for you --
|
|
||||||
AbstractMessageListenerContainer.DEFAULT_PREFETCH_COUNT is 250 -- which is
|
|
||||||
better than unlimited and still large enough to concentrate a small
|
|
||||||
backlog on one consumer.
|
|
||||||
|
|
||||||
The price of all this is the ordering measurement: there is no key, so
|
|
||||||
nothing constrains related messages to one consumer.
|
|
||||||
""".formatted(render(withDefaultPrefetch), idleDefault,
|
|
||||||
render(withPrefetchOne), idleTuned));
|
|
||||||
}
|
|
||||||
|
|
||||||
// ---------------------------------------------------------------- helpers
|
|
||||||
|
|
||||||
private String declare(String name) throws Exception {
|
|
||||||
try (Channel channel = connection.createChannel()) {
|
|
||||||
channel.queueDelete(name);
|
|
||||||
channel.queueDeclare(name, false, false, true, null);
|
|
||||||
}
|
|
||||||
return name;
|
|
||||||
}
|
|
||||||
|
|
||||||
private void publish(String queue, int count) throws Exception {
|
|
||||||
try (Channel channel = connection.createChannel()) {
|
|
||||||
for (int i = 1; i <= count; i++) {
|
|
||||||
channel.basicPublish("", queue, null, String.valueOf(i).getBytes(StandardCharsets.UTF_8));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private List<String> getAll(String queue, int max) throws Exception {
|
|
||||||
List<String> received = new ArrayList<>();
|
|
||||||
try (Channel channel = connection.createChannel()) {
|
|
||||||
for (int i = 0; i < max; i++) {
|
|
||||||
GetResponse response = channel.basicGet(queue, true);
|
|
||||||
if (response == null) {
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
received.add(new String(response.getBody(), StandardCharsets.UTF_8));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return received;
|
|
||||||
}
|
|
||||||
|
|
||||||
private int depth(String queue) throws Exception {
|
|
||||||
try (Channel channel = connection.createChannel()) {
|
|
||||||
return channel.queueDeclarePassive(queue).getMessageCount();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private List<String> consumeWithTwoConsumers(String queue, int count) throws Exception {
|
|
||||||
List<String> completion = new CopyOnWriteArrayList<>();
|
|
||||||
CountDownLatch done = new CountDownLatch(count);
|
|
||||||
List<Channel> channels = new ArrayList<>();
|
|
||||||
for (int i = 1; i <= 2; i++) {
|
|
||||||
long delay = (i == 1) ? 60 : 5;
|
|
||||||
Channel channel = connection.createChannel();
|
|
||||||
channel.basicQos(1);
|
|
||||||
channels.add(channel);
|
|
||||||
channel.basicConsume(queue, false, new DefaultConsumer(channel) {
|
|
||||||
@Override
|
|
||||||
public void handleDelivery(String tag, Envelope envelope, AMQP.BasicProperties props,
|
|
||||||
byte[] body) throws IOException {
|
|
||||||
try {
|
|
||||||
Thread.sleep(delay);
|
|
||||||
}
|
|
||||||
catch (InterruptedException ex) {
|
|
||||||
Thread.currentThread().interrupt();
|
|
||||||
}
|
|
||||||
completion.add(new String(body, StandardCharsets.UTF_8));
|
|
||||||
getChannel().basicAck(envelope.getDeliveryTag(), false);
|
|
||||||
done.countDown();
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
done.await(30, TimeUnit.SECONDS);
|
|
||||||
for (Channel channel : channels) {
|
|
||||||
channel.close();
|
|
||||||
}
|
|
||||||
return List.copyOf(completion);
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Drains a queue across {@code n} competing consumers. prefetch 0 means "do not set it". */
|
|
||||||
private Map<String, Integer> drainAcross(String queue, int n, int expected, int prefetch)
|
|
||||||
throws Exception {
|
|
||||||
Map<String, AtomicInteger> perConsumer = new ConcurrentHashMap<>();
|
|
||||||
CountDownLatch done = new CountDownLatch(expected);
|
|
||||||
List<Channel> channels = new ArrayList<>();
|
|
||||||
for (int i = 1; i <= n; i++) {
|
|
||||||
String name = "consumer-" + i;
|
|
||||||
Channel channel = connection.createChannel();
|
|
||||||
if (prefetch > 0) {
|
|
||||||
channel.basicQos(prefetch);
|
|
||||||
}
|
|
||||||
channels.add(channel);
|
|
||||||
perConsumer.put(name, new AtomicInteger());
|
|
||||||
channel.basicConsume(queue, false, new DefaultConsumer(channel) {
|
|
||||||
@Override
|
|
||||||
public void handleDelivery(String tag, Envelope envelope, AMQP.BasicProperties props,
|
|
||||||
byte[] body) throws IOException {
|
|
||||||
perConsumer.get(name).incrementAndGet();
|
|
||||||
try {
|
|
||||||
Thread.sleep(10);
|
|
||||||
}
|
|
||||||
catch (InterruptedException ex) {
|
|
||||||
Thread.currentThread().interrupt();
|
|
||||||
}
|
|
||||||
getChannel().basicAck(envelope.getDeliveryTag(), false);
|
|
||||||
done.countDown();
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
done.await(30, TimeUnit.SECONDS);
|
|
||||||
for (Channel channel : channels) {
|
|
||||||
channel.close();
|
|
||||||
}
|
|
||||||
Map<String, Integer> counts = new LinkedHashMap<>();
|
|
||||||
perConsumer.forEach((name, count) -> counts.put(name, count.get()));
|
|
||||||
return counts;
|
|
||||||
}
|
|
||||||
|
|
||||||
private static String render(Map<String, Integer> counts) {
|
|
||||||
StringBuilder out = new StringBuilder();
|
|
||||||
counts.entrySet().stream().sorted(Map.Entry.comparingByKey()).forEach((entry) ->
|
|
||||||
out.append(" ").append(entry.getKey()).append(" -> ").append(entry.getValue())
|
|
||||||
.append(" messages\n"));
|
|
||||||
return out.toString();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -4,7 +4,7 @@ Companion project for
|
|||||||
[**Spring Boot and RabbitMQ: Exchanges, Queues, Bindings and a Working Dead-Letter Queue**](https://ankurm.com/spring-boot-rabbitmq-exchanges-dead-letter-queue/)
|
[**Spring Boot and RabbitMQ: Exchanges, Queues, Bindings and a Working Dead-Letter Queue**](https://ankurm.com/spring-boot-rabbitmq-exchanges-dead-letter-queue/)
|
||||||
on ankurm.com.
|
on ankurm.com.
|
||||||
|
|
||||||
Ten tests against a **real RabbitMQ broker**. All four exchange types, manual acknowledgement,
|
Nine tests against a **real RabbitMQ broker**. All four exchange types, manual acknowledgement,
|
||||||
a dead-letter path exercised through `basicNack` and through TTL expiry, and the three ways a
|
a dead-letter path exercised through `basicNack` and through TTL expiry, and the three ways a
|
||||||
topology loses work quietly.
|
topology loses work quietly.
|
||||||
|
|
||||||
@@ -63,7 +63,7 @@ export ERL_ROOT=/path/to/erlang RABBITMQ_HOME=/path/to/rabbitmq_server-3.10.25
|
|||||||
| [`requeue-loop.txt`](docs/output/requeue-loop.txt) | 199 redeliveries, 0 dead-lettered |
|
| [`requeue-loop.txt`](docs/output/requeue-loop.txt) | 199 redeliveries, 0 dead-lettered |
|
||||||
| [`unroutable.txt`](docs/output/unroutable.txt) | `312 NO_ROUTE` |
|
| [`unroutable.txt`](docs/output/unroutable.txt) | `312 NO_ROUTE` |
|
||||||
| [`precondition-failed.txt`](docs/output/precondition-failed.txt) | `406` on an inequivalent argument |
|
| [`precondition-failed.txt`](docs/output/precondition-failed.txt) | `406` on an inequivalent argument |
|
||||||
| [`tests.txt`](docs/output/tests.txt) | 10 tests |
|
| [`tests.txt`](docs/output/tests.txt) | 9 tests |
|
||||||
|
|
||||||
## Five things this module exists to prove
|
## Five things this module exists to prove
|
||||||
|
|
||||||
|
|||||||
@@ -51,21 +51,6 @@ a consumer looked at the work and refused it; `expired` means nobody got to it i
|
|||||||
different incidents with different fixes, and they are indistinguishable without reading the
|
different incidents with different fixes, and they are indistinguishable without reading the
|
||||||
header.
|
header.
|
||||||
|
|
||||||
The third trigger has a detail worth knowing: when `x-max-length` is exceeded RabbitMQ drops from
|
|
||||||
the **head**, so the message that gets dead-lettered is the **oldest** one already queued, not the
|
|
||||||
one that just arrived:
|
|
||||||
|
|
||||||
```
|
|
||||||
=== x-death after x-max-length overflow ===
|
|
||||||
reason maxlen
|
|
||||||
queue orders.bounded
|
|
||||||
body {"orderId":"m-1","detail":"detail for m-1"}
|
|
||||||
```
|
|
||||||
|
|
||||||
Three messages into a queue that holds two, and `m-1` is the one on the DLQ. A bounded queue
|
|
||||||
under sustained overload therefore dead-letters your *backlog* while continuing to accept new
|
|
||||||
work — which is usually what you want for telemetry and exactly wrong for orders.
|
|
||||||
|
|
||||||
`x-death` is a **list**, not a map — one entry per queue the message has been dead-lettered from,
|
`x-death` is a **list**, not a map — one entry per queue the message has been dead-lettered from,
|
||||||
and `count` accumulates. That is how you build a retry limit: read
|
and `count` accumulates. That is how you build a retry limit: read
|
||||||
`x-death[0].count`, and stop republishing past a threshold.
|
`x-death[0].count`, and stop republishing past a threshold.
|
||||||
|
|||||||
@@ -2,21 +2,13 @@
|
|||||||
reason rejected
|
reason rejected
|
||||||
count 1
|
count 1
|
||||||
exchange
|
exchange
|
||||||
time Sat Aug 29 10:30:04 IST 2026
|
time Sat Aug 29 09:55:45 IST 2026
|
||||||
routing-keys [orders.work]
|
routing-keys [orders.work]
|
||||||
queue orders.work
|
queue orders.work
|
||||||
=== x-death after x-message-ttl expiry ===
|
=== x-death after x-message-ttl expiry ===
|
||||||
reason expired
|
reason expired
|
||||||
count 1
|
count 1
|
||||||
exchange
|
exchange
|
||||||
time Sat Aug 29 10:30:06 IST 2026
|
time Sat Aug 29 09:55:47 IST 2026
|
||||||
routing-keys [orders.ttl]
|
routing-keys [orders.ttl]
|
||||||
queue orders.ttl
|
queue orders.ttl
|
||||||
=== x-death after x-max-length overflow ===
|
|
||||||
reason maxlen
|
|
||||||
count 1
|
|
||||||
exchange orders.direct
|
|
||||||
time Sat Aug 29 10:30:06 IST 2026
|
|
||||||
routing-keys [bounded]
|
|
||||||
queue orders.bounded
|
|
||||||
body {"orderId":"m-1","detail":"detail for m-1"}
|
|
||||||
|
|||||||
@@ -1,4 +1,3 @@
|
|||||||
Tests run: 2, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 6.132 s -- in com.ankurm.rabbit.TopologyTrapsTest
|
Tests run: 2, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 5.178 s -- in com.ankurm.rabbit.TopologyTrapsTest
|
||||||
Tests run: 3, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 1.714 s -- in com.ankurm.rabbit.DeadLetterTest
|
Tests run: 3, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 1.797 s -- in com.ankurm.rabbit.DeadLetterTest
|
||||||
Tests run: 1, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 0.061 s -- in com.ankurm.rabbit.MaxLengthTest
|
Tests run: 4, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 4.511 s -- in com.ankurm.rabbit.RoutingTest
|
||||||
Tests run: 4, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 4.476 s -- in com.ankurm.rabbit.RoutingTest
|
|
||||||
|
|||||||
@@ -11,7 +11,6 @@ mvn -B test 2>&1 | tee /tmp/rabbit-test.log > /dev/null
|
|||||||
sed -n '/=== topic exchange ===/,/^order\.high/p' /tmp/rabbit-test.log > docs/output/topic-wildcards.txt
|
sed -n '/=== topic exchange ===/,/^order\.high/p' /tmp/rabbit-test.log > docs/output/topic-wildcards.txt
|
||||||
sed -n '/=== x-death after basicNack/,/^ queue/p' /tmp/rabbit-test.log > docs/output/dead-letter.txt
|
sed -n '/=== x-death after basicNack/,/^ queue/p' /tmp/rabbit-test.log > docs/output/dead-letter.txt
|
||||||
sed -n '/=== x-death after x-message-ttl/,/^ queue/p' /tmp/rabbit-test.log >> docs/output/dead-letter.txt
|
sed -n '/=== x-death after x-message-ttl/,/^ queue/p' /tmp/rabbit-test.log >> docs/output/dead-letter.txt
|
||||||
sed -n '/=== x-death after x-max-length/,/^ body/p' /tmp/rabbit-test.log >> docs/output/dead-letter.txt
|
|
||||||
sed -n '/=== basicNack(requeue=true)/,/still on queue/p' /tmp/rabbit-test.log > docs/output/requeue-loop.txt
|
sed -n '/=== basicNack(requeue=true)/,/still on queue/p' /tmp/rabbit-test.log > docs/output/requeue-loop.txt
|
||||||
sed -n '/=== returned message ===/,/routingKey/p' /tmp/rabbit-test.log > docs/output/unroutable.txt
|
sed -n '/=== returned message ===/,/routingKey/p' /tmp/rabbit-test.log > docs/output/unroutable.txt
|
||||||
sed -n '/=== redeclaring orders.ttl/,+1p' /tmp/rabbit-test.log > docs/output/precondition-failed.txt
|
sed -n '/=== redeclaring orders.ttl/,+1p' /tmp/rabbit-test.log > docs/output/precondition-failed.txt
|
||||||
|
|||||||
@@ -1,73 +0,0 @@
|
|||||||
package com.ankurm.rabbit;
|
|
||||||
|
|
||||||
import com.rabbitmq.client.GetResponse;
|
|
||||||
import org.junit.jupiter.api.Test;
|
|
||||||
import org.springframework.amqp.core.Binding;
|
|
||||||
import org.springframework.amqp.core.BindingBuilder;
|
|
||||||
import org.springframework.amqp.core.DirectExchange;
|
|
||||||
import org.springframework.amqp.core.Queue;
|
|
||||||
import org.springframework.amqp.core.QueueBuilder;
|
|
||||||
import org.springframework.amqp.rabbit.core.RabbitAdmin;
|
|
||||||
import org.springframework.amqp.rabbit.core.RabbitTemplate;
|
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
|
||||||
import org.springframework.boot.test.context.SpringBootTest;
|
|
||||||
|
|
||||||
import java.util.List;
|
|
||||||
import java.util.Map;
|
|
||||||
|
|
||||||
import static org.assertj.core.api.Assertions.assertThat;
|
|
||||||
|
|
||||||
/**
|
|
||||||
* The third dead-letter trigger: a queue that is full.
|
|
||||||
*
|
|
||||||
* <p>Included because the article states all three triggers in one table, and this was the only
|
|
||||||
* row not produced by a run. It is now.
|
|
||||||
*
|
|
||||||
* @see <a href="../../../../../docs/04-dead-lettering.md">docs/04-dead-lettering.md</a>
|
|
||||||
*/
|
|
||||||
@SpringBootTest
|
|
||||||
class MaxLengthTest {
|
|
||||||
|
|
||||||
static final String Q_BOUNDED = "orders.bounded";
|
|
||||||
|
|
||||||
@Autowired
|
|
||||||
RabbitTemplate template;
|
|
||||||
|
|
||||||
@Autowired
|
|
||||||
RabbitAdmin admin;
|
|
||||||
|
|
||||||
@Test
|
|
||||||
@SuppressWarnings("unchecked")
|
|
||||||
void exceedingMaxLengthDeadLettersTheOldestWithReasonMaxlen() throws Exception {
|
|
||||||
this.admin.deleteQueue(Q_BOUNDED);
|
|
||||||
Queue bounded = QueueBuilder.durable(Q_BOUNDED)
|
|
||||||
.maxLength(2)
|
|
||||||
.deadLetterExchange(Topology.DLX)
|
|
||||||
.deadLetterRoutingKey("failed")
|
|
||||||
.build();
|
|
||||||
this.admin.declareQueue(bounded);
|
|
||||||
Binding binding = BindingBuilder.bind(bounded)
|
|
||||||
.to(new DirectExchange(Topology.DIRECT)).with("bounded");
|
|
||||||
this.admin.declareBinding(binding);
|
|
||||||
this.admin.purgeQueue(Topology.Q_DLQ, false);
|
|
||||||
|
|
||||||
// Three messages into a queue that holds two. RabbitMQ drops from the HEAD, so the
|
|
||||||
// FIRST message is the one dead-lettered - the oldest, not the newest.
|
|
||||||
for (String id : List.of("m-1", "m-2", "m-3")) {
|
|
||||||
this.template.convertAndSend(Topology.DIRECT, "bounded", OrderMessage.of(id));
|
|
||||||
}
|
|
||||||
|
|
||||||
var message = this.template.receive(Topology.Q_DLQ, 10_000);
|
|
||||||
assertThat(message).isNotNull();
|
|
||||||
List<Map<String, Object>> deaths =
|
|
||||||
(List<Map<String, Object>>) message.getMessageProperties().getHeader("x-death");
|
|
||||||
System.out.println("=== x-death after x-max-length overflow ===");
|
|
||||||
deaths.get(0).forEach((k, v) -> System.out.printf(" %-16s %s%n", k, v));
|
|
||||||
System.out.println(" body " + new String(message.getBody()));
|
|
||||||
|
|
||||||
assertThat(deaths.get(0)).containsEntry("reason", "maxlen");
|
|
||||||
assertThat(deaths.get(0)).containsEntry("queue", Q_BOUNDED);
|
|
||||||
assertThat(new String(message.getBody())).contains("m-1");
|
|
||||||
}
|
|
||||||
|
|
||||||
}
|
|
||||||
Reference in New Issue
Block a user