Skip to content

Scaling Kafka Consumers with Partitions in asyncio

A Kafka consumer group scales by giving each consumer a share of the topic's partitions, and a partition is consumed by at most one member of the group. That makes the partition count a hard ceiling on parallelism for one-message-at-a-time consumers — and the asyncio answer is to process several keys at once inside each consumer. Measured on Python 3.14 with aiokafka 0.14.0 against Kafka 4.1.0 in Docker, with a 6-partition topic of 6,000 messages and a handler that awaits 5 ms per message (simulated I/O): one consumer handled 196 messages per second, three 559, six 1,053 — and eight 1,050, with two consumers holding no partition. When half the messages shared one key, six consumers managed only 322 per second, because that key's partition held 3,615 of the messages. Processing different keys concurrently within each batch, while keeping each key in order, raised a single consumer to 17,158 per second with 0 ordering violations. Six such consumers, joining one after another, reached 26,029 unique messages per second but processed 10,995 duplicates as rebalances moved partitions mid-batch. This guide measures each lever.

Prerequisites

1. Measure the partition ceiling

A typical consumer handles one message at a time in a getmany loop:

consumer = AIOKafkaConsumer("orders", bootstrap_servers=BS, group_id="order-workers",
                            auto_offset_reset="earliest")
await consumer.start()
while True:
    batch = await consumer.getmany(timeout_ms=200, max_records=100)
    for tp, messages in batch.items():
        for message in messages:
            await handle(message)                      # 5 ms of awaited I/O

Measured with 1, 3, 6 and 8 consumers in one group on a 6-partition topic: 196, 559, 1,053 and 1,050 messages per second. Up to six, each consumer added about one partition's worth of throughput — a 5 ms handler allows about 200 messages per second per partition. At eight, Kafka assigned one partition each to six consumers and none to the other two, which sat idle. Adding consumers beyond the partition count adds standby capacity for failover, not throughput.

Verify: the number of busy consumers equals min(consumers, partitions), read from each consumer's assignment().

6 partitions, 6,000 messages, 5 ms handler A grid of 8 rows by 3 columns. 6 partitions, 6,000 messages, 5 ms handler setup msg/s detail 1 consumer, sequential 196 6 partitions 3 consumers, sequential 559 2 partitions each 6 consumers, sequential 1,053 1 partition each 8 consumers, sequential 1,050 2 consumers idle 6 consumers, half the messages on one key 322 hot partition held 3,615 1 consumer, keys concurrent 17,158 0 ordering violations 6 consumers, keys concurrent 26,029 unique 10,995 duplicates from rebalances 6 consumers, keys concurrent, skewed 370 unique the hot key is still sequential aiokafka 0.14.0, Kafka 4.1.0; the handler awaits 5 ms, standing in for I/O.

2. Choose the partition count up front

Since partitions bound a group's parallelism, size them for the consumer count you expect to need, with room to grow. Increasing partitions later is possible, but it changes which partition each key hashes to, so messages for one key can be split across partitions around the change, briefly breaking per-key ordering:

await admin.create_topics([NewTopic("orders", num_partitions=24, replication_factor=3)])

A rule of thumb from the measurements: required throughput divided by one consumer's throughput per partition — about 200 messages per second here — gives the partition count, and then add headroom. Each partition costs broker resources, so very large counts have their own price; a few dozen per topic is common for services of this size.

Verify: the partition count covers the planned consumer count at peak, and changes to it are scheduled with ordering in mind.

3. Watch for hot keys

The producer chooses a partition by hashing the message key. If one key dominates, its partition dominates. Measured with half of the 6,000 messages using the same key: six consumers handled 322 messages per second, a third of the even-key rate, and the consumer holding the hot partition processed 3,615 messages while the others processed 462–502 each. The other five finished early and waited.

from collections import Counter

def partition_skew(messages_per_partition: dict[int, int]) -> float:
    counts = list(messages_per_partition.values())
    return max(counts) / (sum(counts) / len(counts))     # 1.0 = even

A skew ratio well above 1 — here about 3.6 — means one partition sets the group's speed. Fixes are on the producer side: a finer key, such as order ID instead of customer ID, when per-customer ordering is not required; or a composite key that splits a hot entity into several sub-streams when it is.

Verify: per-partition lag or message counts are monitored, and the skew ratio stays close to 1.

Group throughput by approach 4 horizontal bars comparing 6 consumers, sequential, hot key with the others. Group throughput by approach 6 consumers, sequential, hot key 322/s 6 consumers, sequential, even keys 1,053/s 1 consumer, keys concurrent 17,158/s 6 consumers, keys concurrent 26,029/s unique Concurrency inside a consumer beat adding consumers.

4. Process keys concurrently inside each consumer

Kafka guarantees order within a partition, but most applications only need order per key. Group each batch by key, process each key's messages in order, and run the groups concurrently:

async def process_batch(batch):
    groups = []
    for tp, messages in batch.items():
        by_key = defaultdict(list)
        for m in messages:
            by_key[m.key].append(m)
        groups.extend(by_key.values())

    async def in_order(messages):
        for m in messages:
            await handle(m)                      # one key's messages, one at a time

    await asyncio.gather(*(in_order(g) for g in groups))

while True:
    batch = await consumer.getmany(timeout_ms=200, max_records=500)
    if batch:
        await process_batch(batch)
        await consumer.commit()                  # only after every message in the batch

Measured with one consumer and enable_auto_commit=False: 17,158 messages per second against 196 sequentially, and a check that recorded the last offset seen per key found 0 cases of a key's messages being handled out of order. The batch is committed only after all of its groups finish, so a crash re-delivers the batch rather than skipping it. The gain depends on the handler waiting on I/O; for CPU-bound handlers, concurrency inside one event loop does not help. A hot key stays sequential: on the skewed topic this approach reached only 370 messages per second.

Verify: a test that records per-key offsets as messages are handled finds no key processed out of order, and commits happen only after a batch completes.

5. Plan for rebalances and duplicates

When consumers join or leave, the group rebalances and partitions move. Work in progress on a partition that moves is not committed, and the new owner reads from the last committed offset. Measured with six concurrent-key consumers starting one after another: 10,995 messages were processed more than once, as early consumers had partitions revoked part-way through 500-message batches. On the skewed topic, where one batch could take seconds, there were 24,450 duplicates. Larger batches finish more work per commit but lose more on a rebalance.

consumer = AIOKafkaConsumer(
    "orders", bootstrap_servers=BS, group_id="order-workers",
    enable_auto_commit=False,
    max_poll_records=200,                        # bounded work between commits
    group_instance_id=os.environ["POD_NAME"],    # static membership: restarts don't rebalance
)

Keep handlers idempotent, so duplicates are harmless; bound the work between commits; and consider static membership (group_instance_id), which is designed to let a consumer that restarts quickly reclaim its partitions without a full rebalance — not measured here. For commit strategies in detail, see committing Kafka offsets safely in async consumers.

Verify: a test that adds and removes consumers during consumption produces no lost messages, and duplicates are absorbed by idempotent handling.

Scaling a consumer group A flow of 5 stages. Scaling a consumer group Partitions sized for peak, ~200/s each here Keys no hot partition, skew near 1 Inside a consumer keys concurrent, each key in order Commit after the whole batch Rebalances bounded batches, idempotent handlers More consumers help only up to the partition count.

Verification

A Kafka consumer group scales when:

  • Partitions cover the needed parallelism, and no consumer sits without an assignment by design.
  • Partition skew is near 1, with hot keys split at the producer.
  • Each consumer processes keys concurrently, preserving order per key, when handlers are I/O-bound.
  • Commits follow completed batches, and handlers are idempotent for rebalance duplicates.

Diagnostic Hook: when adding consumers to a group stops increasing throughput, compare the consumer count with the partition count and check each consumer's assignment(). Eight consumers on six partitions left two with no partitions and ran at the same 1,050 messages per second as six.

Pitfalls & edge cases

  • More consumers than partitions. Measured: two idle consumers, no gain.
  • A dominant key. Measured: 322 msg/s against 1,053 with even keys.
  • Concurrent processing with auto-commit. Offsets may be committed before work finishes.
  • Long batches during rebalances. Measured: 10,995 to 24,450 duplicates.

Frequently Asked Questions

How many Kafka consumers can a group use?

At most one per partition. On a 6-partition topic, 6 consumers gave 1,053 msg/s and 8 gave 1,050, with two consumers assigned nothing.

How do I speed up an aiokafka consumer without more partitions?

Group each batch by key and process keys concurrently, each key in order, committing after the batch. One consumer went from 196 to 17,158 msg/s with I/O-bound handlers.

Why is one Kafka consumer much busier than the others?

One key dominates its partition. With half the messages on one key, that consumer handled 3,615 messages and the others about 480 each.

Why does my Kafka consumer process messages twice?

A rebalance moved a partition before its batch was committed, so the new owner re-read it. Bound batch size, use static membership and make handlers idempotent.