Kafka vs RabbitMQ vs NATS vs SQS: Choosing the Right Message Broker
Lesson 4 of 5
Kafka for ordered event streams with replay; RabbitMQ for flexible task routing with competing consumers; NATS for lightweight service communication; SQS for managed AWS simplicity.
The Quick Start: Decision Table
| Workload | Best Fit | Why |
|---|---|---|
| Ordered event stream, replay required | Kafka | Persistent log, consumer groups, offset seek for replay |
| Event streaming, simpler ops | NATS JetStream | Log semantics, lower infrastructure overhead |
| Task queue with competing consumers | RabbitMQ | Flexible routing (exchanges, bindings), built for job dispatch |
| Task queue in AWS, zero ops | SQS | Fully managed, scales automatically, FIFO for ordering |
| Service-to-service RPC | NATS | Built-in request-reply, sub-millisecond latency |
| Fan-out to multiple subscribers | RabbitMQ (fanout) or Core NATS | Fanout exchange (RabbitMQ) or ephemeral pub/sub (NATS) |
Get the model wrong and the cost shows up in production either direction. Picture a six-broker Kafka cluster carrying password resets at 2 messages/second: six brokers of hardware and an on-call rotation for a workload a single queue could carry. Or the reverse — an order-event stream on SQS, where a message is gone once the consumer deletes it, so a bad deployment that processed events wrongly leaves nothing to replay; inventory has to be reconciled from database snapshots.
Match the broker to the messaging model your workload requires. For anything the table doesn't settle, the Decision Framework section classifies your workload first, checks constraints, then picks the broker. Most production systems use more than one broker for different workloads.
- Event streaming (ordered, replayable): Kafka or NATS JetStream
- Task queues (fire-and-forget, competing consumers): RabbitMQ or SQS
- Service communication (request-reply, low-latency): Core NATS
Apache Kafka: The Distributed Event Log
Kafka[1] is purpose-built for high-throughput, ordered, replayable event streaming. It's overkill for task queues but unmatched for audit logs, CDC pipelines, and event sourcing.
Architecture: Clusters of brokers, data organized into topics (partitioned), each partition is an ordered append-only log. Producers write to partitions (by key for ordering or round-robin). Consumers track offsets per consumer group — enabling multiple groups to read the same data independently.
Why ordering matters: All events for the same order go to the same partition (partition key = order ID). A single consumer processes the partition sequentially, guaranteeing per-entity ordering. Different partitions process in parallel across multiple consumer instances. This is critical for workflows where the sequence matters — payment, then shipment, then delivery.
Why replay works: Consumer offsets are just integers — positions in each partition. Reset the offset to 0 and reprocess all events since the beginning. Reset to yesterday's offset and replay the last 24 hours. No other broker provides this flexibility.
Producer: Idempotent Writes
package main
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/segmentio/kafka-go"
)
type OrderEvent struct {
OrderID string `json:"order_id"`
Amount int64 `json:"amount_cents"`
Timestamp time.Time `json:"timestamp"`
}
func publishOrderEvent(ctx context.Context, w *kafka.Writer, event OrderEvent) error {
payload, err := json.Marshal(event)
if err != nil {
return fmt.Errorf("marshal event: %w", err)
}
msg := kafka.Message{
Key: []byte(event.OrderID), // same order → same partition → ordered
Value: payload,
}
return w.WriteMessages(ctx, msg)
}
func newKafkaWriter(brokers []string, topic string) *kafka.Writer {
return &kafka.Writer{
Addr: kafka.TCP(brokers...),
Topic: topic,
Balancer: &kafka.Hash{}, // partition by key
RequiredAcks: kafka.RequireAll, // wait for all ISR replicas
MaxAttempts: 5,
Async: false, // synchronous: fail immediately on error
}
}RequiredAcks = All: Leader waits for all in-sync replicas to acknowledge before returning. Only way to guarantee no loss if leader dies immediately after.
Consumer: Manual Offset Commit
func consumeOrders(ctx context.Context, brokers []string, topic, groupID string) error {
r := kafka.NewReader(kafka.ReaderConfig{
Brokers: brokers,
Topic: topic,
GroupID: groupID,
CommitInterval: 0, // default: CommitMessages commits synchronously
StartOffset: kafka.LastOffset,
})
defer r.Close()
for {
msg, err := r.FetchMessage(ctx) // does NOT auto-commit (ReadMessage does)
if err != nil {
return fmt.Errorf("fetch message: %w", err) // ctx cancelled or reader closed
}
if err := processOrder(ctx, msg); err != nil {
// Don't skip ahead: committing any later offset on this partition
// commits this one too. Stop; a restart resumes from the last commit.
return fmt.Errorf("process offset %d: %w", msg.Offset, err)
}
// Commit only after successful processing
if err := r.CommitMessages(ctx, msg); err != nil {
return fmt.Errorf("commit offset: %w", err)
}
}
}Auto-commit (default in many clients) commits offsets on a timer. If your consumer crashes between auto-commit and finishing processing, that message is lost. Disable auto-commit for any workload where losing data is unacceptable. Skipping a failed message and committing the next one loses it the same way: a committed offset covers every earlier offset in the partition.
RabbitMQ: Flexible Message Routing
RabbitMQ[2] is built on AMQP[3]. Its strength is sophisticated routing: producers publish to exchanges, exchanges route to queues via bindings. This decouples producers from consumers.
Architecture: Exchanges (routing rules), bindings (how messages get routed), queues (final destination). Four exchange types: direct (exact key match), topic (wildcard patterns), fanout (broadcast to all), headers (match message attributes).
Why it excels at task queues: Competing consumers. Multiple workers consume from the same queue. Each message goes to one worker. No consumer groups, no rebalancing slowness.
Publisher with Confirms
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"time"
amqp "github.com/rabbitmq/amqp091-go"
)
type TaskPayload struct {
TaskID string `json:"task_id"`
Type string `json:"type"`
Payload []byte `json:"payload"`
}
func publishTask(ch *amqp.Channel, exchange, routingKey string, task TaskPayload) error {
// Enable publisher confirms
if err := ch.Confirm(false); err != nil {
return fmt.Errorf("enable confirms: %w", err)
}
body, err := json.Marshal(task)
if err != nil {
return fmt.Errorf("marshal task: %w", err)
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
confirmation, err := ch.PublishWithDeferredConfirmWithContext(
ctx, exchange, routingKey, true, false,
amqp.Publishing{
DeliveryMode: amqp.Persistent, // survive broker restart
ContentType: "application/json",
MessageId: task.TaskID,
Body: body,
},
)
if err != nil {
return fmt.Errorf("publish: %w", err) // else confirmation is nil → Wait() panics
}
acked, err := confirmation.WaitContext(ctx) // wait for broker ack
if err != nil {
return fmt.Errorf("wait for confirm: %w", err)
}
if !acked {
return fmt.Errorf("broker nacked task %s", task.TaskID)
}
return nil
}Consumer with Ack
func consumeFromQueue(ch *amqp.Channel, queue string) error {
msgs, err := ch.Consume(
queue, "", false, false, false, false, nil, // autoAck=false: manual ack only
)
if err != nil {
return fmt.Errorf("consume: %w", err)
}
for msg := range msgs {
if err := processTask(msg.Body); err != nil {
log.Printf("process task: %v — requeueing", err)
// Negative ack: message goes back to queue
if nerr := msg.Nack(false, true); nerr != nil {
return fmt.Errorf("nack: %w", nerr)
}
continue
}
// Positive ack: message removed
if err := msg.Ack(false); err != nil {
return fmt.Errorf("ack: %w", err)
}
}
return fmt.Errorf("consume %s: delivery channel closed", queue) // channel or connection lost
}NATS: Lightweight and Fast
NATS[4] is a minimal pub/sub system. Core NATS is ephemeral (no persistence), making it ideal for low-latency service communication. JetStream adds persistence with simpler ops than Kafka.
Core NATS: Sub-millisecond latency, no persistence. Perfect for synchronous RPC and telemetry.
JetStream: Adds persistent streams (log semantics) and durable consumers (like Kafka consumer groups) without the operational overhead.
Core NATS: Request-Reply
package main
import (
"fmt"
"log"
"time"
"github.com/nats-io/nats.go"
)
func rpcExample(nc *nats.Conn) error {
// Service: subscribe and respond
if _, err := nc.Subscribe("api.users.get", func(msg *nats.Msg) {
response := `{"user_id": "42", "name": "Alice"}`
if err := msg.Respond([]byte(response)); err != nil {
log.Printf("respond: %v", err)
}
}); err != nil {
return fmt.Errorf("subscribe: %w", err)
}
// Client: request with timeout
reply, err := nc.Request("api.users.get", []byte("42"), 2*time.Second)
if err != nil {
return err // timeout or no responder
}
log.Printf("response: %s", string(reply.Data))
return nil
}Built-in load balancing: multiple subscribers to the same subject means each request goes to one of N responders.
JetStream: Persistent Streams
func jetStreamPublish(js nats.JetStreamContext, subject string) error {
_, err := js.Publish(subject, []byte(`{"order_id":"123"}`))
return err
}
func jetStreamConsume(js nats.JetStreamContext) error {
// Create stream
if _, err := js.AddStream(&nats.StreamConfig{
Name: "ORDERS",
Subjects: []string{"orders.>"},
Storage: nats.FileStorage,
MaxAge: 7 * 24 * time.Hour,
}); err != nil {
return fmt.Errorf("add stream: %w", err)
}
// Create durable consumer
if _, err := js.AddConsumer("ORDERS", &nats.ConsumerConfig{
Durable: "order-processor",
AckPolicy: nats.AckExplicitPolicy,
AckWait: 30 * time.Second,
FilterSubject: "orders.created",
}); err != nil {
return fmt.Errorf("add consumer: %w", err)
}
// Consume
sub, err := js.PullSubscribe("orders.created", "order-processor")
if err != nil {
return fmt.Errorf("subscribe: %w", err)
}
msgs, err := sub.Fetch(10)
if err != nil {
return fmt.Errorf("fetch: %w", err)
}
for _, msg := range msgs {
if err := msg.Ack(); err != nil {
log.Printf("ack failed for msg: %v", err)
}
}
return nil
}Amazon SQS: Managed Simplicity
SQS[5] is fully managed: no brokers, no disks, no ops. Trade-off: less control, weaker guarantees.
Standard vs FIFO: Standard queues scale unlimited but have best-effort ordering. FIFO queues guarantee order within a message group. The default limit is 300 msg/sec per API action (≈3,000/sec with 10-message batching); enabling high-throughput FIFO spreads message groups across partitions for up to 70,000 msg/sec (700,000 batched) in supported regions.
No replay: Messages are deleted after consumption. Lost order events cannot be recovered.
Send and Consume
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/sqs"
)
type EmailTask struct {
To string `json:"to"`
Subject string `json:"subject"`
Body string `json:"body"`
}
func sendToSQS(ctx context.Context, queueURL string, task EmailTask) error {
cfg, err := config.LoadDefaultConfig(ctx)
if err != nil {
return fmt.Errorf("load AWS config: %w", err)
}
client := sqs.NewFromConfig(cfg)
body, err := json.Marshal(task)
if err != nil {
return fmt.Errorf("marshal task: %w", err)
}
_, err = client.SendMessage(ctx, &sqs.SendMessageInput{
QueueUrl: aws.String(queueURL),
MessageBody: aws.String(string(body)),
})
return err
}
func pollSQS(ctx context.Context, queueURL string) error {
cfg, err := config.LoadDefaultConfig(ctx)
if err != nil {
return fmt.Errorf("load AWS config: %w", err)
}
client := sqs.NewFromConfig(cfg)
result, err := client.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{
QueueUrl: aws.String(queueURL),
MaxNumberOfMessages: 10,
WaitTimeSeconds: 20, // long polling
VisibilityTimeout: 60, // 60s to process
})
if err != nil {
return fmt.Errorf("receive: %w", err) // else result is nil → panic below
}
for _, msg := range result.Messages {
var task EmailTask
if err := json.Unmarshal([]byte(*msg.Body), &task); err != nil {
log.Printf("unmarshal task: %v — leaving for redrive, a malformed body won't fix itself on retry", err)
continue // visibility timeout will requeue; a DLQ redrive policy stops this from looping forever
}
if err := sendEmail(ctx, task); err != nil {
log.Printf("send email (msg %s): %v — leaving for redelivery", aws.ToString(msg.MessageId), err)
continue // visibility timeout will requeue
}
if _, err := client.DeleteMessage(ctx, &sqs.DeleteMessageInput{
QueueUrl: aws.String(queueURL),
ReceiptHandle: msg.ReceiptHandle,
}); err != nil {
log.Printf("delete message: %v — email already sent, expect a duplicate redelivery", err)
}
}
return nil
}Standard queues offer best-effort ordering only. Messages usually arrive in order but are not guaranteed. Use FIFO queues for strict per-group ordering, at the cost of lower throughput.
Head-to-Head Comparison
| Dimension | Kafka | RabbitMQ | NATS JetStream | SQS |
|---|---|---|---|---|
| Model | Log | Broker (AMQP) | Streams | Queue |
| Delivery | At-least-once (exactly-once with idempotent producers & transactions) | At-least-once | At-least-once | At-least-once (Standard); Exactly-once (FIFO) |
| Ordering | Per-partition | Per-queue | Per-stream | Best-effort (Standard); Per-group (FIFO) |
| Throughput | Scales out with partitions and brokers | One CPU core per queue replica[6]; spread across queues | Core NATS: in memory, no persistence; JetStream: bounded by disk + replication | Nearly unlimited (Standard); FIFO 300/s per action, up to 70K high-throughput |
| Replay | Yes | No | Yes | No |
| Retention | Days to months | Until consumed | Configurable | 60s–14 days |
| Operations | Medium (KRaft only, ZooKeeper removed in 4.0) | Medium | Low | None |
Throughput and latency for the self-hosted three depend on hardware, replication, and acknowledgement settings, so benchmark your own message size and durability settings before trusting any published figure.
Decision Framework
Interactive Broker Decision Finder
Step 1 of 2What is your primary architectural workload shape?
graph TD
Start{"Workload shape?"} -->|Event log<br/>need replay| Replay{"Throughput?"}
Start -->|Task queue<br/>fire-and-forget| TQ{"AWS-only?"}
Start -->|Request-reply<br/>low latency| NATS["Core NATS"]
Start -->|Fan-out<br/>notifications| FanOut{"Persistence?"}
Replay -->|>100K msg/s<br/>+ ecosystem| Kafka["Kafka"]
Replay -->|simpler ops<br/>short retention| JetStream["NATS JetStream"]
TQ -->|Yes| SQS["Amazon SQS"]
TQ -->|No / multi-cloud /<br/>complex routing| RMQ["RabbitMQ"]
FanOut -->|Needed| RMQ
FanOut -->|Ephemeral OK| NATS
Step 1: Classify Your Workload
Event Streaming (ordered, replay required): Kafka or NATS JetStream. Choose Kafka if you need ecosystem (Kafka Connect, Schema Registry) or >100K msgs/sec. Choose JetStream for simpler ops and shorter retention windows.
Task Queues (competing consumers, fire-and-forget): RabbitMQ or SQS. Choose RabbitMQ for complex routing or multi-cloud. Choose SQS for AWS-native simplicity.
Service Communication (request-reply, low-latency): Core NATS. Built-in load balancing, sub-millisecond latency. gRPC is better for strongly-typed APIs.
Fan-Out Notifications: RabbitMQ fanout exchange (each subscriber gets its own queue) or Core NATS (ephemeral, millions/sec if loss is acceptable).
Step 2: Check Your Constraints
- Replay required? Eliminate RabbitMQ and SQS.
- Throughput >100K msgs/sec? Kafka or NATS.
- Latency p99
<5ms? RabbitMQ or Core NATS. - Complex routing rules? RabbitMQ only.
- Fully managed, AWS only? SQS.
- Multi-cloud or on-premises? Eliminate SQS.
Step 3: Use Hybrid Architectures
Most production systems use multiple brokers:
- Kafka: Event backbone (all domain events, 7-day retention)
- SQS or RabbitMQ: Task queues (email, PDF generation, webhooks)
- NATS: Internal service communication (request-reply between microservices)
flowchart TD
Events[Domain Events] --> Kafka[Kafka<br/>Event Log]
Kafka --> Analytics[Analytics<br/>Consumer]
Kafka --> OrderConsumer[Order<br/>Consumer]
OrderConsumer --> NATS[NATS<br/>Request-Reply]
Kafka --> Email[Email<br/>Notifier]
Email --> SQS[SQS<br/>Task Queue]
NATS --> Service[Microservice<br/>API]
Production Checklist
- Monitoring: Track consumer lag (Kafka), queue depth (RabbitMQ/SQS), message age. Alert on lag exceeding processing SLO.
- Dead letter queues: Every queue needs a DLQ. Separate consumer logs failure reason and optionally retries.
- Idempotent consumers: At-least-once delivery means duplicates possible. Deduplicate by message ID or make operations idempotent (upsert, not insert).
- Schema evolution: Use schema registry. Message format changes without coordination break consumers.
- Backpressure: Understand your broker's backpressure model (Kafka offset lag, RabbitMQ prefetch, SQS visibility timeout) and configure it.
- Graceful shutdown: Drain in-flight messages before killing consumer. Don't exit with uncommitted offsets or unacknowledged messages.
Local broker harness for integration tests
The fastest feedback loop for queue logic is a docker-compose with all four brokers on one network, so a single test suite exercises consumer-rebalance, redelivery, and DLQ behaviour without per-broker setup drift:
# docker-compose.brokers.yml — paste alongside your test runner
services:
kafka:
image: confluentinc/cp-kafka:7.6.0
ports: ["9092:9092"]
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
CLUSTER_ID: MkU3OEVBNTcwNTJENDM2Qk
# Tight retention for tests so disk doesn't bloat across runs.
KAFKA_LOG_RETENTION_HOURS: 1
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
rabbitmq:
image: rabbitmq:3.13-management
ports: ["5672:5672", "15672:15672"]
environment:
RABBITMQ_DEFAULT_USER: test
RABBITMQ_DEFAULT_PASS: test
nats:
image: nats:2.10
command: ["-js", "-sd", "/data"]
ports: ["4222:4222"]
localstack: # SQS-compatible
image: localstack/localstack:3.4
ports: ["4566:4566"]
environment:
SERVICES: sqs
DEFAULT_REGION: us-east-1The other half is rebalance safety. If a partition is revoked while a polled batch is still processing, the new owner re-reads the uncommitted records, and a commit issued after the revoke lands on a partition this member no longer owns. kafka-go's Reader exposes no revocation hook, so this uses franz-go, whose BlockRebalanceOnPoll holds a rebalance until you call AllowRebalance — the cost is that a slow batch can exceed the rebalance timeout and get the member kicked from the group[7]:
package consumer
import (
"context"
"errors"
"fmt"
"github.com/twmb/franz-go/pkg/kgo"
)
func Consume(ctx context.Context, brokers []string, group, topic string,
process func(context.Context, *kgo.Record) error) error {
cl, err := kgo.NewClient(
kgo.SeedBrokers(brokers...),
kgo.ConsumerGroup(group),
kgo.ConsumeTopics(topic),
kgo.DisableAutoCommit(),
kgo.BlockRebalanceOnPoll(), // no partition moves between poll and AllowRebalance
)
if err != nil {
return fmt.Errorf("new client: %w", err)
}
defer cl.CloseAllowingRebalance()
for {
fetches := cl.PollRecords(ctx, 100) // small batches finish inside the rebalance timeout
if err := ctx.Err(); err != nil {
return err
}
if fetches.IsClientClosed() {
return errors.New("kafka client closed")
}
var fetchErr error
fetches.EachError(func(t string, p int32, err error) {
fetchErr = errors.Join(fetchErr, fmt.Errorf("fetch %s[%d]: %w", t, p, err))
})
if fetchErr != nil {
return fetchErr
}
records := fetches.Records()
for _, rec := range records {
if err := process(ctx, rec); err != nil {
// Commit nothing from this batch; whoever owns the partition next re-reads it.
return fmt.Errorf("process %s[%d]@%d: %w", rec.Topic, rec.Partition, rec.Offset, err)
}
}
if err := cl.CommitRecords(ctx, records...); err != nil {
return fmt.Errorf("commit: %w", err)
}
cl.AllowRebalance() // only now can partitions move to another member
}
}The two compose: with the compose file up, kill one of two consumer instances mid-batch and assert that every record is still processed at least once (duplicates are expected — the consumer must be idempotent, per the checklist above). Run that in CI and rebalance bugs surface in an integration test instead of during a deploy.
When None of the Big Four Fit
The four brokers above cover most workloads, but three constraints push teams toward Apache Pulsar, Redpanda, or NATS JetStream at multi-tenant scale. The pattern: a constraint the big four force you to architect around becomes the dominant cost in the system, and a specialised broker removes it.
Trigger 1: Geo-Replicated Multi-Region with Tenanted Topics
Kafka MirrorMaker 2 replicates asynchronously and, by default, renames topics on the destination to <source>.<topic> (DefaultReplicationPolicy)[8]. Consumers that fail over read a different topic name than producers wrote to, from offsets MM2 translated through its checkpoints, which makes failover testing painful. Pulsar's geo-replication is built in: brokers replicate asynchronously between clusters under the same topic name, enabled per namespace or per topic[9]; synchronous replication is a different topology, a single cluster stretched across regions with BookKeeper region-aware placement[10]. Bookies (the storage layer) and stateless brokers (the serving layer) scale independently[11]. The split-storage architecture is the core advantage: you scale CPU for fan-out without rebalancing data.
# Replicate retail/orders across three regions. Assumes the clusters are registered
# and the "retail" tenant is allowed on all three.
bin/pulsar-admin namespaces create retail/orders
bin/pulsar-admin namespaces set-clusters retail/orders \
--clusters us-east-1,eu-west-1,ap-southeast-1
bin/pulsar-admin namespaces set-retention retail/orders --size 100G --time 7d
# The broker rejects a backlog quota at or above the retention size.
bin/pulsar-admin namespaces set-backlog-quota retail/orders \
--limit 50G --policy producer_request_hold # hold producers, don't drop
bin/pulsar-admin namespaces set-schema-compatibility-strategy retail/orders \
--compatibility BACKWARD_TRANSITIVESet once, these policies apply to every topic in the namespace, including topics created later. MM2's nearest equivalent is a topic regex per replication flow (us-west->us-east.topics = foo.*), with ACLs and quotas managed separately; a topic the regex misses simply stays local.
Trigger 2: Sub-10ms Tail Latency Without Tuning Page Cache
Kafka writes through the Linux page cache and leaves flushing to the kernel[1], so its tail latency inherits kernel writeback behaviour. Redpanda[12] reimplements the Kafka protocol in C++ on the Seastar framework, pins one thread per CPU core, and uses direct memory access (DMA) for disk I/O rather than the page cache[13]. The goal is a flatter tail-latency distribution; whether it beats a tuned Kafka on your workload is a measurement, not an assumption.
The trade-off is operational, not client-side. Kafka clients for protocol 0.11+ work against Redpanda with a short list of documented exceptions, such as one SCRAM mechanism per user and no request_percentage quota[14], and Kafka Streams and Connect run on those clients. What diverges is the broker: configuration, tuning, and upgrades go through Redpanda's tooling, not Apache Kafka's. For trading systems, ad bidding, and real-time fraud detection, a tighter latency tail can be worth that. For a CDC pipeline feeding a data warehouse, it is not.
// Redpanda is Kafka API compatible — the same kafka-go client works,
// you just swap brokers. Compare latency histograms before migrating.
package main
import (
"context"
"time"
"github.com/segmentio/kafka-go"
)
func newRedpandaWriter(brokers []string, topic string) *kafka.Writer {
return &kafka.Writer{
Addr: kafka.TCP(brokers...),
Topic: topic,
Balancer: &kafka.Hash{},
RequiredAcks: kafka.RequireAll,
// Tight client-side timeout: fail fast rather than queue behind a slow ack.
WriteTimeout: 50 * time.Millisecond,
BatchTimeout: 5 * time.Millisecond, // tight batching for low latency
Async: false,
}
}
func publishLowLatency(ctx context.Context, w *kafka.Writer, key, value []byte) error {
deadline, cancel := context.WithTimeout(ctx, 20*time.Millisecond)
defer cancel()
return w.WriteMessages(deadline, kafka.Message{Key: key, Value: value})
}A decent migration smoke test: run the producer above against both clusters with identical traffic, capture the latency histograms, and verify p99.9 stays inside your SLO under a sustained load target such as 80% of provisioned capacity — pick the figure from your own headroom needs, not a universal number. If Kafka clears the bar, do not migrate — the operational divergence costs more than the latency savings buy.
Trigger 3: Hundreds of Tenants with Per-Tenant Streams
NATS JetStream's underrated strength is that a stream is a cheap entity — you can run tens of thousands of streams on a single cluster, each with its own retention, replication factor, and ACL. The same workload on Kafka requires either tens of thousands of topics (which destroys the metadata layer) or topic-per-tenant with namespaced subjects (which makes ACLs untenable). For SaaS backends with a "one stream per customer" model — webhook delivery, audit logs, per-tenant change feeds — JetStream wins on operational simplicity by an order of magnitude.
# Provision a per-tenant stream with bounded retention and ACLs in three commands.
# Run this from your tenant-onboarding job; it is idempotent.
TENANT_ID="$1"
nats stream add "tenant-${TENANT_ID}" \
--subjects="tenants.${TENANT_ID}.events.>" \
--storage=file \
--retention=limits \
--max-age=720h \
--max-bytes=10737418240 \
--replicas=3 \
--discard=old
nats consumer add "tenant-${TENANT_ID}" "tenant-${TENANT_ID}-webhooks" \
--filter="tenants.${TENANT_ID}.events.>" \
--ack=explicit \
--max-deliver=10 \
--wait=30s \
--backoff=linear --backoff-min=1s --backoff-max=5m
nats auth user add "tenant-${TENANT_ID}-publisher" \
--pub="tenants.${TENANT_ID}.events.>" \
--sub="_INBOX.>" \
--account=tenantsThe same provisioning on Kafka — topic creation, ACL update, consumer group quota — requires multiple admin API calls, an ACL service round-trip, and a metadata controller update that rate-limits new topics. At 10k tenants, the difference between a sub-second NATS provision and a multi-second Kafka provision shows up as onboarding latency for paying customers.
Decision Rule
The honest version: pick a niche broker only when the constraint that triggers it is in your top-three production risks. Pulsar for compliant multi-region with tight RPOs. Redpanda for a p99.9 latency target that Kafka misses in your own measurements. JetStream for tenant-per-stream SaaS topologies. Outside those scenarios, the smaller community and thinner ecosystem will cost more in incident response time than the technical advantage saves in steady state.
Frequently Asked Questions
When should I use Kafka vs RabbitMQ?
Use Kafka for ordered event streaming with replay and consumer groups. Use RabbitMQ for task queues with flexible routing, competing consumers, and no replay requirement.
What is the difference between queue semantics and log semantics?
Queues: one consumer per message, deleted after ack, no replay. Logs: persistent, offset-based, multiple consumer groups, replay by offset reset.
Is Kafka overkill for simple task queues?
Yes. Kafka's operational complexity (KRaft, partitions, rebalancing) justifies only high-throughput workloads or replay requirements. Use RabbitMQ or SQS for simple job dispatch. Note: Kafka 4.0 (March 2025) removed ZooKeeper entirely — KRaft is now the only option for cluster coordination.
Can NATS replace Kafka for event streaming?
NATS JetStream provides log semantics with lower ops overhead than Kafka. Good for moderate throughput; less ecosystem tooling than Kafka. Choose based on retention needs and team ops capacity.
Keep Reading
- Event-Driven Microservices in Go: Kafka, Sagas, and the Outbox Pattern — Kafka implementation with saga orchestration, idempotent consumers, and the outbox pattern
- Microservices Architecture: From Monolith to Production-Ready Services — Where message brokers fit in microservices: service boundaries, communication patterns, and the event backbone
- Scaling Redis for High-Throughput Systems — Redis Streams and Pub/Sub for lightweight patterns when a full broker is overkill
- Idempotency Patterns in Distributed Systems — Every at-least-once broker (Kafka, RabbitMQ, NATS, SQS) demands idempotent consumers; this is the consumer-side rulebook
- Kafka Producer Tuning Cheat Sheet — Once you have picked Kafka, this cheat sheet is the next 30 minutes of work
Coming Next
Choosing the broker is only half the battle. Implementing it in code without losing events or running into distributed consistency issues is where the real work begins. In our next deep dive, we build a transactional outbox publisher and idempotent consumer in Go using Kafka. Read the Event-Driven Microservices Guide.
Sources
- 1.Apache Kafka Documentation — ASF, 2026
- 2.RabbitMQ Documentation — rabbitmq.com, 2026
- 3.AMQP 0-9-1 Protocol Specification — amqp.org, 2008
- 4.NATS Documentation — docs.nats.io, 2026
- 5.Amazon Simple Queue Service Developer Guide — AWS Documentation, 2026
- 6.RabbitMQ Documentation — Queues — rabbitmq.com, 2026
- 7.kgo package — github.com/twmb/franz-go/pkg/kgo — pkg.go.dev, 2026
- 8.Apache Kafka 4.3 Documentation — Geo-Replication (Cross-Cluster Data Mirroring) — ASF, 2026
- 9.Apache Pulsar 4.0 (LTS) — Pulsar geo-replication (Administration) — ASF, 2026
- 10.Apache Pulsar 4.0 (LTS) — Geo Replication — ASF, 2026
- 11.Apache Pulsar 4.0 (LTS) — Architecture Overview — ASF, 2026
- 12.Redpanda Documentation — Redpanda Data, 2026
- 13.Redpanda Documentation — How Redpanda Works — Redpanda Data, 2026
- 14.Redpanda Documentation — Kafka Compatibility — Redpanda Data, 2026
Engineering Team
An independent engineering publication covering distributed systems, databases, and production infrastructure. Every factual claim is cited to a primary source or removed.
Read Next
Consistent Hashing: The Algorithm Behind Scalable Distributed Systems
Adding one cache server shouldn't invalidate every key. Consistent hashing with virtual nodes and bounded loads — full Go and Java implementations.
Microservices Architecture: From Monolith to Production-Ready Services
When to decompose a monolith, how to define boundaries, and the patterns that work: API gateways, sagas, and event-driven comms.
Distributed Rate Limiting at Scale: The Probabilistic Drop Architecture
Probabilistic drop rate limiting: uncoordinated enforcement that takes Redis off the request path, with no coordination per request.