Series overview
Part 18 of 2864% complete
2026-07-02•15 min read

Competing consumers and work queues

Chapter 7’s groupId configuration was introduced almost in passing, distinguishing publish-subscribe (each service its own consumer group) from something else it deferred. This chapter is that something else: using multiple instances of one service, sharing one consumer group, to distribute — not duplicate — a stream of work.

1. Problem the Pattern Solves

notification-service (added in Chapter 7’s exercise) sends an order-confirmation email for every OrderPlaced event. During a normal day, one instance keeps up easily. During a flash sale, order volume spikes tenfold, and a single instance’s email-sending throughput — bounded by the sending provider’s API rate limit per connection and the CPU cost of rendering each email template — can’t keep up. A growing backlog of unsent confirmation emails means customers wait twenty, then forty minutes for a receipt that should arrive in seconds.

Scaling notification-service to three replicas doesn’t help by itself if each replica is in its own Kafka consumer group (as inventory-service and payment-service correctly are, per Chapter 7’s design for independent, parallel reactions to the same event) — that configuration would mean each of the three replicas processes every email independently, sending each confirmation three times, which is a worse outcome than the original backlog.

Forces in tension:

  • Throughput vs. correctness of “exactly one worker per task.” Distributing work across multiple instances scales throughput, but only if the distribution mechanism guarantees each unit of work goes to exactly one instance — not zero, and not several.
  • Ordering vs. parallelism. Kafka’s ordering guarantee is per-partition; distributing work across multiple consumer instances (which Kafka does by assigning different partitions to different instances within a group) means work from different partitions can be processed out of relative order — acceptable for independent emails, unacceptable if the work items have a required processing sequence.
  • Elasticity vs. partition count as a hard ceiling. Kafka can’t have more active consumers in a group than partitions on the topic — scaling beyond the partition count adds idle instances, not more throughput, a specific and easy-to-miss capacity-planning trap.
  • Failure handling per unit of work. A worker crashing mid-task shouldn’t lose that task — it should become available for another worker to pick up, which depends on correctly configured offset-commit behavior (commit only after successful processing, not before).

2. Core Idea

Competing consumers means multiple instances of the same service share one Kafka consumer group, so that each message on the topic is delivered to exactly one instance within that group — the instances “compete” for messages, each processing a disjoint subset, giving horizontal scaling of throughput for a stream of independent work items.

same consumer group:

notification-service

Kafka topic: order-confirmations

4 partitions

notification-service

replica 1

(partitions 0,1)

notification-service

replica 2

(partition 2)

notification-service

replica 3

(partition 3)

same consumer group:

notification-service

Kafka topic: order-confirmations

4 partitions

notification-service

replica 1

(partitions 0,1)

notification-service

replica 2

(partition 2)

notification-service

replica 3

(partition 3)

Compare this directly to Chapter 7’s diagram: there, inventory-service and payment-service were in different consumer groups, each receiving its own full copy of every OrderPlaced event — publish-subscribe. Here, all three notification-service replicas share one consumer group — competing consumers. Same broker, same topic mechanics, a completely different delivery semantic, chosen by how consumer groups are configured, not by anything different about Kafka itself.

Participants:

  • Work queue (topic) — order-confirmations, partitioned to allow parallel consumption; partition count sets the hard ceiling on how many active consumer instances can do useful work simultaneously.
  • Competing consumers — multiple notification-service replicas, all in the notification-service consumer group, each assigned a disjoint subset of partitions by Kafka’s group coordinator.
  • Partition assignment — Kafka’s own rebalancing mechanism, which reassigns partitions among the group’s current members whenever an instance joins or leaves (a scale-up, a scale-down, or a crash).

Commonly confused with:

  • Publish-subscribe (Chapter 7 and 8). The distinguishing factor is entirely the consumer group configuration: same group, one instance per message (competing consumers, this chapter); different groups, every instance its own copy (pub-sub, Chapters 7–8). The same Kafka topic can technically be consumed both ways simultaneously by different services, as Northwind’s orders.events topic already effectively demonstrates across inventory-service, payment-service, and (via its own topic) notification-service.
  • A saga’s choreography (Chapter 11). A saga’s event handlers react to specific business events with specific compensating logic; competing consumers is a generic scaling mechanism for any stream of independent, parallelizable work — a saga step could itself be implemented by a competing-consumers group if that step’s processing needs to scale, but the two concepts operate at different levels.
  • Kubernetes horizontal pod autoscaling alone. Scaling notification-service’s replica count via HPA (Chapter 1’s Kubernetes discussion) only produces more throughput if those replicas are actually distributing work correctly via a shared consumer group — adding replicas without this pattern’s consumer-group configuration, as Section 1 shows, produces duplicated processing instead of parallelized processing.

3. When to Use It

Strong indicators:

  • A stream of independent, order-insensitive (or order-tolerant-within-acceptable-bounds) units of work whose processing throughput needs to scale with volume — Northwind’s confirmation emails are exactly this: each is independent of every other.
  • Demonstrated or anticipated volume spikes (flash sales) that a single consumer instance can’t handle within an acceptable processing-latency target.
  • Work whose individual processing cost (CPU, external API calls) is high enough that parallelizing across instances meaningfully helps, as opposed to work so cheap that a single instance already processes it faster than it arrives.

Concrete use cases:

  • E-commerce, as here: order-confirmation emails, receipt PDF generation, and post-purchase analytics event processing are all naturally parallelizable, independent work streams.
  • Image/video processing pipelines: resizing uploaded images or transcoding video are CPU-intensive, embarrassingly parallel tasks that scale directly with more competing consumer instances.
  • Batch data processing: nightly ETL jobs split into per-partition work units, processed by a pool of worker instances that scale up during the batch window and down afterward.
  • Any background-job system: sending SMS notifications, generating reports, processing webhook deliveries to external systems — all are classic competing-consumer workloads once volume exceeds a single worker’s capacity.

Prerequisites:

  • A topic partitioned appropriately for the target parallelism — partition count must be planned with headroom for future scaling, since increasing partition count later is possible but can disrupt existing key-based ordering guarantees (a consideration specific to Kafka’s partition-increase semantics).
  • Idempotent or safely-retriable work items — a rebalance can cause a message to be reprocessed by a different instance if the original didn’t commit its offset in time, so “sent this email exactly once” needs the same idempotency discipline this series has required since Chapter 1 (Section 5 shows the specific mechanism for this workload).
  • An understanding of what ordering guarantee, if any, the workload actually needs — Northwind’s confirmation emails need none; a workload that does needs careful partition-key design (Chapter 7’s per-order keying) rather than this pattern’s default parallelism.

4. When Not to Use It

  • Work with a strict, cross-item ordering requirement that spans partitions. If processing order must be strictly global (rare, but real for some ledger-style workloads), competing consumers’ partition-parallel processing directly conflicts with that requirement — a single-consumer, single-partition design, accepting the throughput ceiling, may be the correct trade-off there.
  • Low, steady volume that a single instance handles comfortably. Provisioning multiple notification-service replicas and worrying about partition counts for a workload that never exceeds a single instance’s capacity is unneeded operational complexity — evaluate against measured or realistically projected peak volume.
  • Work that isn’t safely retriable/idempotent, with no way to make it so. If a work item’s processing has an irreversible side effect that can’t be made idempotent (rare, but worth checking explicitly), a rebalance-triggered reprocessing becomes a real correctness problem this pattern can’t fully protect against — that work item needs a different design, or a compensating mechanism (Chapter 11’s saga thinking, applied at the level of a single work item).
  • Overengineering signal: massively over-partitioning a topic “for future scale” well beyond any realistic projected need, since each partition is itself a small resource cost (open file handles, replication overhead) on the broker — size partition count to a reasonable multiple of anticipated peak parallelism, not an arbitrarily large number.

5. Implementation Example

Topic provisioning, sized for Northwind’s anticipated flash-sale peak parallelism:

Terminal window
kafka-topics.sh --create --topic order-confirmations --partitions 12 --replication-factor 3 \
--bootstrap-server kafka.default.svc.cluster.local:9092

Twelve partitions gives headroom for up to twelve actively-consuming notification-service replicas — chosen based on Northwind’s projected flash-sale peak of roughly ten replicas needed, with a small margin, not an arbitrarily large number (Section 4’s overengineering warning, made concrete).

The competing consumer, deployed as multiple replicas sharing one consumer group:

notification-service/src/main/kotlin/in/o612/eng/northwind/notification/OrderConfirmationListener.kt
package `in`.o612.eng.northwind.notification
import `in`.o612.eng.northwind.order.api.OrderPlacedNotification
import org.springframework.kafka.annotation.KafkaListener
import org.springframework.stereotype.Component
@Component
class OrderConfirmationListener(
private val emailSender: EmailSender,
private val sentEmailLog: SentEmailLog, // idempotency guard, see below
) {
@KafkaListener(
topics = ["order-confirmations"],
groupId = "notification-service", // shared across every replica — this is what makes them compete, not duplicate
concurrency = "3", // 3 consumer threads per pod instance, in addition to horizontal replica scaling
)
fun onOrderPlaced(notification: OrderPlacedNotification) {
if (sentEmailLog.alreadySent(notification.orderId)) return // idempotency check — see Section 7
val email = emailSender.sendConfirmation(notification.orderId)
sentEmailLog.recordSent(notification.orderId)
}
}
notification-service/src/main/resources/application.yml
spring:
kafka:
consumer:
group-id: notification-service
enable-auto-commit: false # commit only after successful processing — see Section 7
max-poll-records: 10 # bound how much work one poll claims, avoiding one slow batch blocking rebalance
listener:
ack-mode: manual_immediate # explicit ack after successful send, not before
k8s/notification-service-hpa.yaml
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: notification-service
spec:
scaleTargetRef: { apiVersion: apps/v1, kind: Deployment, name: notification-service }
minReplicas: 2
maxReplicas: 10
metrics:
- type: External
external:
metric:
name: kafka_consumergroup_lag
selector: { matchLabels: { topic: order-confirmations, consumergroup: notification-service } }
target: { type: AverageValue, averageValue: "50" }

Scaling on consumer lag, not CPU, is the deliberate, correct choice here — CPU utilization doesn’t directly reflect whether the backlog is growing; lag does, directly answering “is this service keeping up with incoming work,” which is precisely the question autoscaling needs answered for a queue-consuming workload.

Idempotency guard, since a rebalance mid-processing (a replica scaling down, or crashing) can cause the same message to be redelivered to a different replica after the original committed no offset:

notification-service/src/main/kotlin/in/o612/eng/northwind/notification/SentEmailLog.kt
package `in`.o612.eng.northwind.notification
import org.springframework.jdbc.core.JdbcTemplate
import org.springframework.dao.DuplicateKeyException
import org.springframework.stereotype.Component
import java.util.UUID
@Component
class SentEmailLog(private val jdbc: JdbcTemplate) {
fun alreadySent(orderId: UUID): Boolean =
jdbc.queryForObject("SELECT COUNT(*) FROM sent_emails WHERE order_id = ?", Int::class.java, orderId)!! > 0
fun recordSent(orderId: UUID) {
try {
jdbc.update("INSERT INTO sent_emails (order_id, sent_at) VALUES (?, now())", orderId)
} catch (e: DuplicateKeyException) {
// Two competing consumers raced on the same redelivered message
// — the email was (or is about to be) sent exactly once either
// way; this just prevents a duplicate log entry, not a duplicate send.
}
}
}

This is the same inbox-style idempotency discipline from Chapter 12, applied to a competing-consumers workload rather than a pub-sub one — the mechanism generalizes directly.

6. Step-by-Step Flow

sent_emails tablenotification-service replica 2notification-service replica 1Kafka (order-confirmations, 12 partitions)sent_emails tablenotification-service replica 2notification-service replica 1Kafka (order-confirmations, 12 partitions)Flash sale begins — HPA scales notification-service 2 -> 6 replicaspar[competing consumption]No overlap — each message processed by exactly one replicarebalance: reassign partitions across 6 replicasmessages from partitions 0-1messages from partitions 2-3check + record order A (new)send email for order Acheck + record order B (new)send email for order B
sent_emails tablenotification-service replica 2notification-service replica 1Kafka (order-confirmations, 12 partitions)sent_emails tablenotification-service replica 2notification-service replica 1Kafka (order-confirmations, 12 partitions)Flash sale begins — HPA scales notification-service 2 -> 6 replicaspar[competing consumption]No overlap — each message processed by exactly one replicarebalance: reassign partitions across 6 replicasmessages from partitions 0-1messages from partitions 2-3check + record order A (new)send email for order Acheck + record order B (new)send email for order B
  1. Client action. A flash sale drives a sharp spike in OrderPlaced events, each producing a corresponding order-confirmations message.
  2. API request equivalent. Kafka’s consumer-group coordinator rebalances partition assignment as the HPA scales notification-service from 2 to 6 replicas in response to growing lag.
  3. Service behavior. Each replica processes only the messages from its assigned partitions — no replica sees a message another replica is also handling.
  4. Database interaction. Each replica checks and records into the shared sent_emails table before/after sending, guarding against the specific redelivery risk a rebalance introduces.
  5. Inter-service communication. None beyond the Kafka consumption itself — notification-service doesn’t call any other Northwind service to do its job.
  6. Error or failure handling. If a replica crashes mid-processing before committing its offset, that message is redelivered to whichever replica picks up the orphaned partition after rebalancing — the idempotency guard ensures this doesn’t produce a duplicate email.
  7. Observability signals. Consumer lag per partition, and the rate of DuplicateKeyException catches in SentEmailLog (a direct measure of how often rebalance-triggered redelivery is actually happening) are the two most important metrics for this pattern specifically.
  8. Final response/outcome. Confirmation emails go out within the target latency even during a tenfold volume spike, because throughput scaled with the number of active competing consumers, up to the partition-count ceiling.

7. Production Concerns

  • Timeouts, retries, idempotency. enable-auto-commit: false plus ack-mode: manual_immediate (Section 5) ensures an offset is only committed after successful processing — the specific configuration that makes “crash mid-processing → redelivered, not lost” true; getting this backward (committing before processing) would silently lose work on a crash instead.
  • Data consistency. The idempotency guard’s sent_emails table check-and-insert should itself be as atomic as the database allows (a unique constraint on order_id, as the DuplicateKeyException handling assumes) — a check-then-insert with a race window wide enough for two competing consumers to both pass the check is a real risk under high concurrency, worth verifying with a targeted concurrency test.
  • Partition count as a scaling ceiling. Twelve partitions means twelve is the hard limit on useful parallelism, regardless of how many replicas the HPA spins up — monitor for “replicas without partitions assigned” (idle, wasted replicas) as a sign the partition count needs revisiting, not just the replica count.
  • API versioning. Not directly relevant — this pattern is about consumption topology, not message schema, though the message schema itself follows the same versioning discipline as any other event (Chapters 7–8).
  • Authentication and service-to-service trust. No new surface beyond standard Kafka ACLs (Chapter 7) — all replicas in a consumer group share the same identity/credential for consuming this topic.
  • Logging, metrics, tracing, correlation IDs. Tag every log line with which replica (pod name) and which partition processed a given message — essential for diagnosing an uneven load distribution or a specific partition falling behind.
  • Kubernetes deployment, autoscaling. Scale on consumer lag (Section 5), not CPU or memory — a queue-consuming workload’s “am I keeping up” question is answered directly by lag and only indirectly, unreliably, by resource utilization.
  • Testing strategy. Test the idempotency guard under simulated concurrent redelivery (two threads racing to process the same orderId) explicitly — this is the one behavior in this pattern that’s genuinely hard to catch with a single-threaded test but easy to get wrong under real rebalance conditions.
  • Migration strategy. If notification-service currently runs as a single instance with no consumer-group scaling concerns, introduce competing consumers only once a real volume-driven backlog (as Section 1 describes) justifies the added partition-planning and idempotency work — not preemptively for a workload that hasn’t shown the need.

8. Common Mistakes

  1. Scaling replicas without a shared consumer group. Running three notification-service replicas each in its own consumer group (a copy-paste of Chapter 7’s pub-sub configuration) triples email sends instead of parallelizing them — exactly Section 1’s opening failure mode. Fix: verify every replica intended to compete for work shares the identical groupId.
  2. Partition count lower than the target maximum parallelism. Provisioning a topic with 3 partitions while planning to scale to 10 replicas wastes 7 replicas that receive no partition assignment at all. Fix: provision partition count with headroom for realistic peak parallelism, as Section 5 does.
  3. Committing offsets before processing completes. Using auto-commit (or manually acknowledging before the email actually sends) means a crash after the commit but before the send silently loses that confirmation email — the opposite failure mode from the intended “redelivered, not lost” guarantee. Fix: commit (acknowledge) only after processing succeeds, as Section 5’s manual_immediate configuration does.
  4. No idempotency guard, assuming rebalances are rare enough to ignore. Rebalances happen routinely during normal autoscaling, not just crashes — treating redelivery as a rare edge case rather than a routine occurrence under this pattern leads to real, periodic duplicate sends. Fix: build the idempotency guard as a standard, always-on part of the consumer, not a defensive afterthought.
  5. Scaling on CPU instead of consumer lag. A competing-consumers workload can have low CPU utilization per instance while still falling badly behind (e.g., each task spends most of its time waiting on a slow external email-provider API) — CPU-based autoscaling won’t react to this at all. Fix: scale on consumer lag directly, as shown in Section 5’s HPA configuration.
  6. Assuming any ordering guarantee across the whole topic. Relying on messages from different partitions being processed in the order they were produced (they generally won’t be, under competing consumers) breaks any workload with cross-item sequencing needs. Fix: if ordering matters at all, design the partition key deliberately (as Chapter 7 did with orderId) so related items land on the same partition and are processed in relative order by whichever single consumer holds that partition — never assume global ordering across partitions.

9. Decision Guide

Problem signalUse this pattern?WhyAlternative
A stream of independent, parallelizable work whose volume exceeds one instance’s throughputYesScales throughput by distributing work across instances that don’t duplicate each other’s effort—
Low, steady volume within a single instance’s comfortable capacityNoAdded partition-planning and idempotency complexity for no throughput benefitSingle consumer instance
Work has a strict, cross-item, whole-topic ordering requirementNo, or only with careful partition-key designPartition-parallel processing conflicts with strict global orderingSingle partition (accepting the throughput ceiling), or careful per-key partitioning
Multiple independent services need their own full copy of every eventNo (this is pub-sub, not competing consumers)Separate consumer groups deliver independent copies; shared groups distribute, not duplicatePublish-subscribe (Chapters 7-8), one consumer group per service
Work items aren’t safely retriable and can’t be made idempotentNo, or redesign the work item firstRebalance-triggered redelivery becomes a genuine correctness risk without idempotencyRedesign the work to be idempotent, or use a different delivery mechanism entirely

10. Hands-On Exercise

Extend it: provision a second competing-consumers topic, receipt-pdf-generation, for a new, more CPU-intensive work type (rendering a PDF receipt), and choose a partition count independently justified by that workload’s expected volume and per-item processing cost — explain why it might reasonably differ from order-confirmations’ twelve.

Simulate a failure: force a rebalance mid-processing (scale notification-service down by one replica while it’s actively consuming a backlog) and confirm, using the sent_emails table, that no order’s confirmation email is either lost or sent more than once.

Decision question, with justification required: Northwind wants confirmation emails sent in the exact order orders were placed, per customer (so a customer never receives their second order’s confirmation before their first’s, if both were placed within the same second during a flash sale). Does the current partition-by-nothing-in-particular scheme support this, and if not, what partition-key change would you make, and what does that change cost in terms of this pattern’s parallelism ceiling?

11. Key Takeaways

  • Competing consumers means multiple instances of the same service share one Kafka consumer group, so each message is processed by exactly one instance — the direct opposite delivery semantic from publish-subscribe’s “every group gets its own copy,” configured purely by whether instances share a groupId.
  • Partition count sets a hard ceiling on useful parallelism — provision it with headroom for realistic peak scale, and monitor for idle, unassigned replicas as a sign it needs revisiting.
  • Commit offsets only after successful processing, never before — this is what makes a crash mid-processing result in safe redelivery rather than silent data loss.
  • Idempotency is not optional under this pattern — rebalances (including routine autoscaling events, not just crashes) will cause redelivery, and every competing-consumer workload needs a guard against processing the same item twice.
  • Scale on consumer lag, not CPU or memory — lag directly answers whether the consumer group is keeping up with incoming work, which is the actual question autoscaling needs to answer here.
  • Kafka provides ordering only within a partition — competing consumers’ parallelism across partitions means no default guarantee of cross-item ordering; design the partition key deliberately if any ordering requirement exists.
  • Apply this pattern to a demonstrated, volume-driven throughput problem, exactly as every other pattern in this series — not preemptively to a workload that hasn’t shown the need.
Spring BootKotlinMicroservicesKafka

Type to search the site.

↑↓ navigate⏎ openPowered by Pagefind