Skip to content
  • 8 min read

Kafka — Testing Scenarios, Load & CI

The Kafka-specific behaviours that fakes cannot show: ordering per key, redelivery, poison messages, consumer lag under load. The tests reuse the fixtures and helpers from 05 Testing Kafka-Based Systems: bootstrap_servers, topic_factory, group_id, producer, the service fixture that runs the consumer loop in a thread, produce_json and consume_until. The test functions without a file name below live in tests/test_invoicing_integration.py. Versions used: pytest 9.1.1, confluent-kafka 2.15.1, testcontainers 4.15.0, Kafka 4.3.1.

Ordering by Key

def test_order_per_key_is_preserved(bootstrap_servers, topics, producer, service):
    for order_id in ("ord-1", "ord-2", "ord-3"):
        for amount in (100, 200, 300):
            produce_json(producer, topics["in"], order_id, new_event(order_id, amount))

    messages = consume_until(bootstrap_servers, topics["out"], lambda ms: len(ms) >= 9)

    per_key = defaultdict(list)
    for m in messages:
        per_key[m.key()].append(json.loads(m.value())["total_cents"])
    assert all(amounts == [100, 200, 300] for amounts in per_key.values()), per_key
    assert all(len({m.partition() for m in messages if m.key() == k}) == 1 for k in per_key)
  • Assert order per key, never globally: records of different keys in different partitions interleave in any order.
  • Use a topic with more than one partition — on a single partition, bugs caused by key handling and partitioning cannot show up.
  • The second assertion checks that each key stayed on one partition of the output topic.

Duplicates and Idempotency

At-least-once delivery means the consumer will see some records twice: after a crash before commit, after a rebalance, after a producer retry without idempotence. Test both sides.

def test_redelivered_event_is_not_invoiced_twice(bootstrap_servers, topics, producer, service):
    event = new_event("ord-7")
    produce_json(producer, topics["in"], "ord-7", event)
    produce_json(producer, topics["in"], "ord-7", event)               # same event_id: simulated redelivery
    produce_json(producer, topics["in"], "ord-7", new_event("ord-7"))  # marker: a later, different event

    messages = consume_until(bootstrap_servers, topics["out"], lambda ms: len(ms) >= 2)

    sources = [json.loads(m.value())["source_event"] for m in messages]
    assert sources.count(event["event_id"]) == 1

And that uncommitted work really is redelivered — the property the idempotency above relies on:

# tests/test_redelivery.py
import pytest
from confluent_kafka import Consumer

from tests.kafka_helpers import new_event, produce_json

pytestmark = pytest.mark.integration


def read_one(bootstrap_servers: str, topic: str, group_id: str, commit: bool) -> int:
    consumer = Consumer({
        "bootstrap.servers": bootstrap_servers,
        "group.id": group_id,
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
    })
    consumer.subscribe([topic])
    try:
        for _ in range(50):
            msg = consumer.poll(0.2)
            if msg is not None and not msg.error():
                if commit:
                    consumer.commit(message=msg, asynchronous=False)
                return msg.offset()
        raise AssertionError("no message")
    finally:
        consumer.close()


def test_uncommitted_message_is_redelivered(bootstrap_servers, topic_factory, producer, group_id):
    topic = topic_factory("orders", partitions=1)
    produce_json(producer, topic, "ord-1", new_event("ord-1"))

    first = read_one(bootstrap_servers, topic, group_id, commit=False)    # "crash" before commit
    second = read_one(bootstrap_servers, topic, group_id, commit=True)
    assert first == second == 0                                          # at-least-once: same record again

    produce_json(producer, topic, "ord-1", new_event("ord-1"))
    assert read_one(bootstrap_servers, topic, group_id, commit=True) == 1  # committed -> not redelivered

Poison Messages and the DLQ

def test_poison_message_goes_to_dlq_and_consumer_keeps_going(bootstrap_servers, topics, producer, service):
    produce_json(producer, topics["in"], "ord-9", b"{not json")
    produce_json(producer, topics["in"], "ord-9", new_event("ord-9"))

    [dead] = consume_until(bootstrap_servers, topics["dlq"], lambda ms: len(ms) >= 1)
    [invoice] = consume_until(bootstrap_servers, topics["out"], lambda ms: len(ms) >= 1)

    headers = dict(dead.headers())
    assert dead.value() == b"{not json"
    assert headers["error"].startswith(b"JSONDecodeError")
    assert json.loads(invoice.value())["order_id"] == "ord-9"

What this proves: the bad record is kept with its original bytes and the reason, and the next record on the same key is still processed — a consumer that crashes or retries forever on a poison message blocks its whole partition.

More cases worth a test:

Case Expected behaviour
Invalid JSON / unknown schema ID DLQ, no retry
Valid payload, business rule violated DLQ or a domain "rejected" event — whatever the spec says
Transient failure (DB down, HTTP 503) Retry with backoff, no DLQ, no commit until success
DLQ produce itself fails Consumer does not commit the original offset
Header-only difference (schema-version: 2) Routed to the right parser

Consumer Lag in Load Tests

For load tests, "requests per second" is the wrong main number for a consumer. Measure how fast the group drains a backlog and whether lag returns to zero after the load stops.

# tests/lag.py
import time

from confluent_kafka import ConsumerGroupTopicPartitions, TopicPartition
from confluent_kafka.admin import AdminClient, OffsetSpec


def consumer_lag(admin: AdminClient, group: str, topic: str) -> dict[int, int]:
    """Lag per partition = log-end offset - committed offset of the group."""
    partitions = [TopicPartition(topic, p) for p in admin.list_topics(topic, timeout=5).topics[topic].partitions]
    committed = admin.list_consumer_group_offsets(
        [ConsumerGroupTopicPartitions(group, partitions)]
    )[group].result().topic_partitions
    ends = {tp.partition: f.result().offset for tp, f in admin.list_offsets(
        {tp: OffsetSpec.latest() for tp in partitions}
    ).items()}
    # offset < 0: the group never committed this partition -> count all of it (fine for fresh test topics)
    return {tp.partition: ends[tp.partition] - max(tp.offset, 0) for tp in committed}


def wait_for_zero_lag(admin: AdminClient, group: str, topic: str, timeout: float = 30.0) -> float:
    """Return how long the group needed to catch up; fail if it did not within `timeout`."""
    start = time.monotonic()
    while True:
        lag = consumer_lag(admin, group, topic)
        if sum(lag.values()) == 0:
            return time.monotonic() - start
        if time.monotonic() - start > timeout:
            raise AssertionError(f"group {group} still lags on {topic}: {lag}")
        time.sleep(0.5)
def test_consumer_drains_a_burst(admin, bootstrap_servers, topics, producer, service, group_id):
    for i in range(500):
        producer.produce(topics["in"], key=f"ord-{i % 20}", value=json.dumps(new_event(f"ord-{i % 20}")).encode())
    assert producer.flush(10) == 0

    seconds = wait_for_zero_lag(admin, group_id, topics["in"], timeout=60)

    print(f"500 events drained in {seconds:.1f}s")
    assert seconds < 30

For raw broker throughput, use the bundled perf tools before blaming the service:

docker exec kafka /opt/kafka/bin/kafka-producer-perf-test.sh --bootstrap-server localhost:9092 \
  --topic perf --num-records 100000 --record-size 512 --throughput -1 --command-property acks=all
# 100000 records sent, 41893.6 records/sec (20.46 MB/sec), 602.63 ms avg latency, ...

docker exec kafka /opt/kafka/bin/kafka-consumer-perf-test.sh --bootstrap-server localhost:9092 \
  --topic perf --num-records 100000 --group perf-check

Load-test checklist for a consumer service:

  • Pre-fill a backlog (N records), start the consumers, measure time to zero lag → drain rate per consumer.
  • Scale consumers 1 → partition count and check the drain rate grows; beyond the partition count it must stay flat.
  • Kill one consumer during the run: lag spikes during the rebalance and then drains; no records are lost (compare counts or IDs).
  • Use realistic keys: a single hot key keeps one partition busy while the others idle.
  • Drive HTTP load that produces events with Locust and track lag as a separate metric next to response times.

Flakiness Pitfalls

Symptom Cause Fix
Test sometimes sees no messages Consumer started with auto.offset.reset=latest (the default) after the event was produced earliest for assertion consumers and test groups
Test sees messages from earlier runs Reused topic or group.id with committed offsets Unique topic and group per test (uuid)
UNKNOWN_TOPIC_OR_PART right after creating a topic Topic metadata not propagated yet Wait until list_topics shows all partitions with leaders
First test is slow or logs Coordinator load in progress __consumer_offsets is created on first group use Accept it, or warm up with a throwaway consumer in the session fixture
Every test waits about 3 s before consuming group.initial.rebalance.delay.ms default 3000 KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS=0 in test brokers
Test passes locally, times out in CI time.sleep(2) instead of a condition Poll with a deadline (consume_until)
Last events missing Producer not flushed before assertion or exit produce_json waits for the delivery report; flush() in teardown
Next test waits for a rebalance Consumer from the previous test not closed close() in finally / fixture teardown
Ordering test is always green Topic has one partition, or all test keys hash to one partition 3+ partitions and several keys
Assert on offsets or counts is off by one Transaction markers take offsets; compaction leaves gaps Assert on record contents, not on offset arithmetic
New topic created by a typo Broker auto-creates topics KAFKA_AUTO_CREATE_TOPICS_ENABLE=false
Tests interfere under pytest-xdist Shared topic names or group IDs Unique names per test; each worker gets its own session broker (or share one via KAFKA_BOOTSTRAP_SERVERS)
Works on host, fails from a container Advertised listener is localhost Second listener for the Docker network (02 Local Setup)

General techniques against flaky tests: Reliability & Flakiness.

CI

Option 1: Testcontainers — nothing to configure, the GitHub-hosted Ubuntu runners have Docker. Option 2: a service container, shared by all tests of the job:

# .github/workflows/kafka-tests.yml (fragment)
jobs:
  tests:
    runs-on: ubuntu-latest
    services:
      kafka:
        image: apache/kafka:4.3.1
        ports:
          - 9092:9092
        env:
          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_LISTENER_NAMES: CONTROLLER
          KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093
          KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: "1"
          KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: "1"
          KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: "1"
          KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: "0"
          KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false"
        options: >-
          --health-cmd "/opt/kafka/bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092"
          --health-interval 5s --health-timeout 10s --health-retries 20
    env:
      KAFKA_BOOTSTRAP_SERVERS: localhost:9092
    steps:
      - uses: actions/checkout@v4
      - uses: astral-sh/setup-uv@v6
      - name: Unit and contract tests
        run: uv run pytest -m "not integration"
      - name: Integration tests
        run: uv run pytest -m integration --junitxml=report.xml
      - name: Broker logs on failure
        if: failure()
        run: docker logs "${{ job.services.kafka.id }}" | tail -200
  • Run unit and contract tests first: they fail in seconds and need no broker.
  • The advertised localhost:9092 works because steps run on the runner host; a job that runs inside a container needs the service name (kafka:9092) as the advertised address.
  • Pin the broker image and the client versions; upgrade both in a dedicated PR.

QA Checklist

  • Handler logic is unit-tested without Kafka: valid, invalid, duplicate, poison
  • Producer sample events are validated against the consumer contract, or schema compatibility is checked against the registry in CI
  • Integration tests use a real broker (Testcontainers or a service container), with pinned image versions
  • Every test creates its own topics and consumer group and deletes the topics afterwards
  • No fixed sleeps: every wait is a condition with a deadline
  • Producers in tests wait for delivery reports; consumers are always closed
  • Ordering is asserted per key on a multi-partition topic
  • Redelivery of uncommitted records and idempotent handling of duplicates are tested
  • Poison messages land in the DLQ with the original payload and error details, and the partition keeps flowing
  • Load tests report drain rate and time-to-zero lag, including a consumer restart during the run
  • CI prints broker logs on failure

See also