This tutorial ends with an Apache Kafka broker running in Docker on your Mac, a topic named orders, and two Python programs that talk to it: a producer that publishes JSON order events and a consumer that reads them. You will then use those programs to show the property that makes Kafka different from a job queue. Reading a message does not remove it, so a consumer that restarts picks up where it stopped, and a second, independent consumer can replay the whole history from the beginning. You will see both happen, and measure how far behind a consumer is, with Kafka’s own tools.

Kafka is a distributed log. Programs called producers append records (Kafka calls them events or messages; this article says message) to a named stream, and programs called consumers read them back in order. The server that stores the messages is a broker. Production clusters run several brokers; one is enough to learn every idea in this article, and the client code you write does not change when you add more. A named stream is a topic, and each topic is split into ordered partitions; a message’s position within its partition is its offset. Consumers that share a name form a consumer group, which is how Kafka tracks what each reader has seen. Steps 3 and 6 show each of these at work.

Two version notes before you start. Kafka 4.0 removed ZooKeeper, the separate coordination service older tutorials start alongside Kafka, so the setup here runs Kafka on its own, in what Kafka calls KRaft mode. If a guide asks you to start a zookeeper container, it was written for Kafka 3 or earlier. And the Python client is confluent-kafka, Confluent’s widely used wrapper around the C library librdkafka.

What you end up with

  • A compose.yaml that runs a single Kafka 4.3.1 broker in KRaft mode on localhost:9092.
  • A topic named orders with three partitions, created with Kafka’s bundled command-line tools.
  • producer.py, which publishes orders keyed by customer and reports the partition and offset each one landed at.
  • consumer.py, which reads orders as a member of a named consumer group and records its progress in Kafka.
  • A Makefile whose default target prints a help screen, and a pytest suite that checks the round trip against the running broker.

Prerequisites

  • macOS 13 or later.
  • Docker Desktop 4.x (docker.com), running. You need Compose v2 (docker compose version), which Docker Desktop bundles. This tutorial was validated with Docker Engine 29.7.2.
  • uv 0.5 or later (docs.astral.sh/uv), for example brew install uv. uv installs Python 3.13 for the project if you do not have it. Validated with uv 0.11.26.
  • make, from the Xcode Command Line Tools (xcode-select --install).
  • Familiarity with running containers from a compose.yaml file is assumed; the Docker Compose quickstart covers it. No prior Kafka knowledge is needed.

Step 1: Create the project and its .gitignore

The project is one directory holding the broker’s Compose file and the Python code. The .gitignore comes first so that the virtual environment uv creates in the next steps, and the caches pytest writes later, never get committed.

Create the file

mkdir -p ~/projects/kafka/getting-started-kafka-python-macos
cd ~/projects/kafka/getting-started-kafka-python-macos
touch .gitignore

Add the code: .gitignore

# Python
__pycache__/
*.py[cod]
.venv/
.pytest_cache/

# macOS and scratch
.DS_Store
*.log
tmp/

Detailed breakdown

  • .venv/ is where uv puts the project’s virtual environment. It is large, machine-specific, and rebuilt from uv.lock on any machine with uv sync.
  • __pycache__/, *.py[cod], and .pytest_cache/ are bytecode and test caches that Python and pytest regenerate on every run.
  • Nothing Kafka writes needs ignoring. The broker keeps its data inside its container, not in the project directory, which matters in Step 7.

Step 2: Run a single-node broker with Docker Compose

Everything else in this tutorial talks to this broker, so it has to be running and reachable from your Mac before any Python exists. The official apache/kafka image will start with no configuration at all, but the settings below are written out on purpose. One of them, the advertised listener, is the setting behind a very common “my client can’t connect” failure, and it is much easier to understand when you can see it.

Create the file

touch compose.yaml

Add the code: compose.yaml

services:
  kafka:
    image: apache/kafka:4.3.1
    ports:
      - "9092:9092"
    environment:
      # One process plays both KRaft roles: it stores data and runs the metadata quorum.
      KAFKA_NODE_ID: 1
      KAFKA_PROCESS_ROLES: broker,controller
      KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093
      KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
      # Where the broker listens, and the address it tells clients to use.
      KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      # A single broker cannot keep three copies of Kafka's internal topics.
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      # Start a new consumer group immediately instead of waiting 3 seconds.
      KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
    healthcheck:
      test: ["CMD", "/opt/kafka/bin/kafka-broker-api-versions.sh", "--bootstrap-server", "localhost:9092"]
      interval: 5s
      timeout: 10s
      retries: 20

Detailed breakdown

  • The image is pinned to apache/kafka:4.3.1, the current stable release when this was written. A floating latest tag would give a reader next year a different broker than the one this article’s output came from.
  • The KRaft block. A Kafka cluster needs a controller: the component that tracks which topics exist, which broker leads each partition, and which brokers are alive. In Kafka 4 the controllers form a small voting group (the quorum) using the Raft consensus algorithm, hence “KRaft”. Here one process is both broker and controller (KAFKA_PROCESS_ROLES: broker,controller), and KAFKA_CONTROLLER_QUORUM_VOTERS lists the quorum as just itself: node 1, reachable on port 9093.
  • Listeners versus advertised listeners. A listener is a named network endpoint the broker accepts connections on. KAFKA_LISTENERS opens two: the PLAINTEXT listener on port 9092 for clients, and the CONTROLLER listener on 9093 for quorum traffic, which stays inside the container. The advertised listener is different: it is the address the broker tells clients to use. A Kafka client connects to whatever address you give it (the bootstrap server), asks for cluster metadata, and then reconnects to each broker at the advertised address in that metadata. Advertising localhost:9092 works here because the ports mapping makes that address reach the container from your Mac. Advertise the container’s own hostname instead and the first connection succeeds while every one after it fails, because your Mac cannot resolve a Docker-internal name.
  • KAFKA_LISTENER_SECURITY_PROTOCOL_MAP maps each listener name to a protocol. Both are unencrypted PLAINTEXT, which is fine on your own machine and never acceptable across a network.
  • The replication-factor settings. Kafka stores consumer progress and transaction state in internal topics, and by default wants three copies of each on three different brokers. With one broker it cannot create them, and consumers typically hang waiting for them. Setting the factor to 1 is what makes a single-node broker usable.
  • The health check runs a bundled tool that only succeeds once the broker answers requests. It lets docker compose up --wait return when Kafka is actually ready, not merely when the container has started.

Start the broker. The --wait flag blocks until the health check passes:

docker compose up -d --wait

The last lines of output are:

 Container getting-started-kafka-python-macos-kafka-1 Started
 Container getting-started-kafka-python-macos-kafka-1 Waiting
 Container getting-started-kafka-python-macos-kafka-1 Healthy

Healthy is the line to look for. The first run also pulls the image, which takes a minute or so; later starts take a few seconds.

Step 3: Create a topic with three partitions

Producers and consumers both refer to a stream of messages by name, and that name is a topic. You could let the broker create topics on first use, but it would give orders a single partition, and the partition count shapes everything you will observe later. Creating it deliberately, with Kafka’s own command-line tools, lets you choose.

A partition is one ordered, append-only log inside a topic. Each message written to a partition gets the next sequence number in that partition, its offset, starting at 0. Kafka guarantees order within a partition and makes no promise across partitions. Splitting a topic into partitions is how Kafka spreads load: different partitions can live on different brokers and be read by different consumers in parallel.

The CLI tools ship inside the image under /opt/kafka/bin, so run them with docker compose exec:

docker compose exec kafka /opt/kafka/bin/kafka-topics.sh \
  --bootstrap-server localhost:9092 \
  --create --topic orders --partitions 3 --replication-factor 1
docker compose exec kafka /opt/kafka/bin/kafka-topics.sh \
  --bootstrap-server localhost:9092 --describe --topic orders
Created topic orders.
Topic: orders	TopicId: KFSxc7YjRQO58_O7BL3vsA	PartitionCount: 3	ReplicationFactor: 1	Configs: min.insync.replicas=1
	Topic: orders	Partition: 0	Leader: 1	Replicas: 1	Isr: 1	Elr: 	LastKnownElr: 
	Topic: orders	Partition: 1	Leader: 1	Replicas: 1	Isr: 1	Elr: 	LastKnownElr: 
	Topic: orders	Partition: 2	Leader: 1	Replicas: 1	Isr: 1	Elr: 	LastKnownElr: 

Your TopicId will differ; Kafka generates it. --replication-factor 1 asks for one copy of each partition, the only option with one broker. In the per-partition lines, Leader is the broker that serves reads and writes for that partition, and Replicas and Isr (in-sync replicas) list the brokers holding copies. All three name broker 1, the node ID from compose.yaml. Elr (eligible leader replicas) and LastKnownElr list replicas that have dropped out of the ISR but could still safely become leader; they stay empty on a one-broker cluster. min.insync.replicas=1 in the header means a write succeeds once one copy has it, the only workable value with one broker.

Step 4: Add the Kafka client and shared settings

Both Python programs and the tests need the same two facts: where the broker is, and which topic to use. Putting them in one module means the tests can point the same code at a throwaway topic, and you can aim everything at a different broker later with an environment variable instead of an edit.

Initialize the project with uv, delete the placeholder script it generates, and add the client library plus pytest:

uv init --app --name orders-stream --python 3.13 .
rm main.py
uv add 'confluent-kafka==2.15.1'
uv add --dev 'pytest==9.1.1'

uv init keeps your existing .gitignore. It creates pyproject.toml, .python-version, and an empty README.md, and runs git init if the directory is not already inside a Git repository; uv add creates uv.lock and the .venv/ environment. Confirm the client imports and see which librdkafka it wraps:

uv run python -c 'import confluent_kafka as c; print(c.version(), c.libversion())'
2.15.1 ('2.15.1', 34537983)

The integer is the same librdkafka version packed into a number; you can ignore it.

Create the file

touch config.py

Add the code: config.py

"""Connection settings shared by the producer, the consumer, and the tests."""

import os

BOOTSTRAP_SERVERS = os.environ.get("KAFKA_BOOTSTRAP_SERVERS", "localhost:9092")
TOPIC = os.environ.get("KAFKA_TOPIC", "orders")

Detailed breakdown

  • BOOTSTRAP_SERVERS is the address a client dials first. As Step 2 explained, the client only uses it to fetch cluster metadata; after that it connects to the advertised address of each broker. With several brokers you would list two or three here, comma-separated, so startup survives one being down.
  • TOPIC defaults to the orders topic from Step 3. The producer and consumer import this module, so neither hard-codes the name.
  • Environment variables are the override path. Nothing in this tutorial sets them, but they are how you would point the same code at another cluster.

Step 5: Write the producer

The producer turns Python dictionaries into Kafka messages. Each message has a key and a value, both bytes as far as Kafka is concerned. The value here is the order as JSON. The key is the customer name, and it decides where the message goes: the client hashes the key and uses the result to pick a partition, so every order for one customer lands in the same partition and therefore stays in order relative to that customer’s other orders. (Retries after a network error can reorder messages unless the producer sets enable.idempotence; on a local broker that does not come up.)

Create the file

touch producer.py

Add the code: producer.py

"""Publish order events to Kafka, keyed by customer."""

import argparse
import json

from confluent_kafka import Producer

from config import BOOTSTRAP_SERVERS, TOPIC

CUSTOMERS = ["ana", "ben", "gia"]


def make_order(order_id: int) -> dict:
    """Build a deterministic order so every run of the tutorial prints the same data."""
    return {
        "order_id": order_id,
        "customer": CUSTOMERS[(order_id - 1) % len(CUSTOMERS)],
        "amount_cents": 1000 + order_id * 250,
    }


def produce_orders(orders: list[dict], topic: str = TOPIC) -> list[tuple[int, int]]:
    """Send each order and return the (partition, offset) Kafka assigned to it."""
    producer = Producer({"bootstrap.servers": BOOTSTRAP_SERVERS})
    placements: list[tuple[int, int]] = []
    errors: list[str] = []

    def on_delivery(err, msg):
        if err is not None:
            errors.append(str(err))
            return
        placements.append((msg.partition(), msg.offset()))
        print(
            f"sent order {json.loads(msg.value())['order_id']} "
            f"key={msg.key().decode()} -> partition {msg.partition()} offset {msg.offset()}"
        )

    for order in orders:
        producer.produce(
            topic,
            key=order["customer"],
            value=json.dumps(order),
            on_delivery=on_delivery,
        )
    remaining = producer.flush(10)
    if errors or remaining:
        raise RuntimeError(f"delivery failed: {errors or f'{remaining} still queued'}")
    return placements


def main() -> None:
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("--start", type=int, default=1, help="first order id")
    parser.add_argument("--count", type=int, default=6, help="number of orders")
    args = parser.parse_args()
    orders = [make_order(i) for i in range(args.start, args.start + args.count)]
    produce_orders(orders)


if __name__ == "__main__":
    main()

Detailed breakdown

  • make_order builds order n from its number alone, cycling through three customers, so order 1 is always Ana’s for 1,250 cents and order 4 is Ana’s again for 2,000. Deterministic data is what lets you compare your terminal to the output in Step 8.
  • Why these three names. The client’s default partitioner hashes the key with CRC32 and takes the result modulo the partition count. With only three keys and three partitions, keys often collide: alice, bob, and carol all hash to partition 2. ana, ben, and gia were chosen because they land on three different partitions, which makes the per-partition behavior visible. Real systems have thousands of keys and the spread evens out. Java clients use a different hash (murmur2), so the same key can map to a different partition from a Java producer; mixing client languages on one topic needs the partitioner configured to match.
  • produce() does not send anything by itself. It adds the message to an in-memory buffer and returns immediately. librdkafka batches buffered messages per partition and sends them from a background thread. That is why the result arrives later, through the on_delivery callback.
  • The delivery callback (Kafka calls its result a delivery report) receives either an error or the delivered message, now carrying the partition and offset the broker assigned. Offsets exist only after the broker has written the message, which is why the program prints them here and not at produce() time.
  • flush(10) waits up to 10 seconds for every buffered message to be delivered and runs the callbacks. It returns the number of messages still undelivered. Skipping it is the classic producer bug: the script exits, the buffer is discarded, and nothing reaches Kafka. Errors are collected and raised after flush rather than inside the callback, so a failure surfaces as one clear exception from produce_orders.
  • Caveat for long-running producers: callbacks only run during flush() or poll(). A producer that sends continuously should call producer.poll(0) inside its loop, or its callbacks pile up unprocessed.
  • --start and --count let you send a specific range of orders. Step 8 uses them to add orders 7 and 8 without repeating 1 through 6.

Step 6: Write the consumer

The consumer reads the topic and remembers how far it got. That memory lives in Kafka, not in the program, and it is attached to a consumer group: a name that one or more consumer processes share. Kafka divides a topic’s partitions among the members of a group, and stores, per group and per partition, the offset of the next message the group should read. That stored position is the group’s committed offset. A different group has its own committed offsets and reads the same messages independently.

Create the file

touch consumer.py

Add the code: consumer.py

"""Read order events from Kafka as a member of a consumer group."""

import argparse
import json

from confluent_kafka import Consumer

from config import BOOTSTRAP_SERVERS, TOPIC


def consume_orders(group: str, topic: str = TOPIC, idle_seconds: float = 5.0) -> list[dict]:
    """Read until no message arrives for idle_seconds, then commit and return."""
    consumer = Consumer(
        {
            "bootstrap.servers": BOOTSTRAP_SERVERS,
            "group.id": group,
            "auto.offset.reset": "earliest",
        }
    )
    consumer.subscribe([topic])
    received: list[dict] = []
    try:
        while True:
            msg = consumer.poll(idle_seconds)
            if msg is None:
                break
            if msg.error():
                raise RuntimeError(msg.error())
            order = json.loads(msg.value())
            received.append(order)
            print(
                f"got order {order['order_id']} customer={order['customer']} "
                f"amount_cents={order['amount_cents']} "
                f"(partition {msg.partition()} offset {msg.offset()})"
            )
    finally:
        consumer.close()
    print(f"{len(received)} order(s) read by group {group!r}")
    return received


def main() -> None:
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("--group", default="billing", help="consumer group id")
    args = parser.parse_args()
    consume_orders(args.group)


if __name__ == "__main__":
    main()

Detailed breakdown

  • group.id names the consumer group. The billing default stands for one downstream system; Step 8 adds a second, analytics, to show two groups reading the same topic without affecting each other.
  • auto.offset.reset: earliest only matters for a group with no committed offset yet, which is every group on its first run. It says “start at the oldest message still in the partition”. The default, latest, would start at the end of each partition and skip everything already written, so a brand-new consumer run after the producer would print nothing.
  • subscribe() joins the group. The broker then assigns this consumer its share of partitions; as the only member, it gets all three. Joining takes a moment, which is one reason the first poll() can take longer than later ones.
  • poll(idle_seconds) returns one message, or None if nothing arrived within the timeout. A real consumer loops forever. This one stops after 5 quiet seconds so that each run ends on its own, which makes it easy to run from make and from tests. The cost is that every run takes at least 5 seconds.
  • msg.error() covers errors reported through the message stream, such as a partition that is temporarily unavailable. Raising on them makes failures visible; production code would retry or log the recoverable ones.
  • Commits. The client commits offsets automatically in the background every 5 seconds (enable.auto.commit is on by default), and consumer.close() commits the final position before leaving the group. Leaving out close() means the last few messages you processed are read again next time. Automatic commits also mean a crash between processing a message and the next commit replays that message, so consumers that must not double-process need more care than this one takes.

Step 7: Wrap the commands in a Makefile

Step 8 runs the same handful of commands several times, some of them long docker compose exec lines. A Makefile gives each one a short name and keeps the flags consistent between runs. Running make on its own prints the list of targets, so you never have to open the file to remember them.

Create the file

touch Makefile

Add the code: Makefile

KAFKA_BIN := docker compose exec kafka /opt/kafka/bin
BOOTSTRAP := --bootstrap-server localhost:9092
GROUP ?= billing

.DEFAULT_GOAL := help
.PHONY: help up down topic produce consume lag test

help: ## Show this help
	@grep -E '^[a-z-]+:.*## ' $(MAKEFILE_LIST) | awk 'BEGIN {FS = ":.*## "} {printf "  %-9s %s\n", $$1, $$2}'

up: ## Start the Kafka broker and wait until it is healthy
	docker compose up -d --wait

down: ## Stop the broker and delete its data
	docker compose down -v

topic: ## Create the orders topic with three partitions
	$(KAFKA_BIN)/kafka-topics.sh $(BOOTSTRAP) --create --if-not-exists --topic orders --partitions 3 --replication-factor 1
	$(KAFKA_BIN)/kafka-topics.sh $(BOOTSTRAP) --describe --topic orders

produce: ## Send six orders to the orders topic
	uv run producer.py

consume: ## Read new orders as GROUP (default: billing)
	uv run consumer.py --group $(GROUP)

lag: ## Show committed offsets and lag for GROUP
	$(KAFKA_BIN)/kafka-consumer-groups.sh $(BOOTSTRAP) --describe --group $(GROUP)

test: ## Run the integration tests against the running broker
	uv run pytest -v

Detailed breakdown

  • .DEFAULT_GOAL := help makes a bare make print the help screen. The help recipe finds every line of the form target: ## description and prints the pairs, so adding a target with a ## comment documents it.
  • KAFKA_BIN and BOOTSTRAP hold the repeated prefix for the bundled tools. They run inside the container, where localhost:9092 is the broker itself.
  • GROUP ?= billing sets a default that you can override per command, as in make consume GROUP=analytics.
  • down deletes data. The broker writes its log inside the container, and docker compose down removes the container. The -v also removes the anonymous volumes the image declares, so repeated runs do not leave unused volumes behind. Every make up after a make down starts with no topics and no messages.
  • --if-not-exists makes make topic safe to repeat on a broker that already has the topic.
  • Makefile recipes must be indented with a tab, not spaces. If make reports missing separator, your editor converted them.

Check the help screen:

make
  help      Show this help
  up        Start the Kafka broker and wait until it is healthy
  down      Stop the broker and delete its data
  topic     Create the orders topic with three partitions
  produce   Send six orders to the orders topic
  consume   Read new orders as GROUP (default: billing)
  lag       Show committed offsets and lag for GROUP
  test      Run the integration tests against the running broker

Step 8: Stream orders and watch Kafka keep them

Every piece is in place, and a short sequence of commands shows the three behaviors the introduction promised. Start from an empty broker so your partitions and offsets match the output below. make down deletes the topic you created in Step 3, and make topic creates it again:

make down
make up
make topic

Apart from make echoing each command before running it, the make topic output matches Step 3 with a new TopicId.

Publish six orders

Send orders 1 through 6:

make produce
uv run producer.py
sent order 2 key=ben -> partition 2 offset 0
sent order 5 key=ben -> partition 2 offset 1
sent order 1 key=ana -> partition 0 offset 0
sent order 4 key=ana -> partition 0 offset 1
sent order 3 key=gia -> partition 1 offset 0
sent order 6 key=gia -> partition 1 offset 1

Each customer’s orders went to one partition (Ana to 0, Gia to 1, Ben to 2), and each partition numbered its own messages from offset 0. The lines are not in order-ID order, and on your machine the three pairs may appear in a different sequence. Delivery reports come back per partition batch, and nothing orders one partition’s batch relative to another’s. Within each customer, order 1 is still before order 4, and order 2 before order 5.

Read them as the billing group

Run the consumer. It uses the billing group by default:

make consume
uv run consumer.py --group billing
got order 2 customer=ben amount_cents=1500 (partition 2 offset 0)
got order 5 customer=ben amount_cents=2250 (partition 2 offset 1)
got order 3 customer=gia amount_cents=1750 (partition 1 offset 0)
got order 6 customer=gia amount_cents=2500 (partition 1 offset 1)
got order 1 customer=ana amount_cents=1250 (partition 0 offset 0)
got order 4 customer=ana amount_cents=2000 (partition 0 offset 1)
6 order(s) read by group 'billing'

All six arrived, each customer’s orders in the order they were sent. As with the producer, the interleaving of partitions can differ on your run; the order within a partition cannot. The command takes about eight seconds: joining the group, reading, then the 5-second idle wait before it stops.

Run it again: the group remembers

Run exactly the same command a second time:

make consume
uv run consumer.py --group billing
0 order(s) read by group 'billing'

Nothing new. The previous run committed offset 2 on each partition when it closed, so this run started after the last message. The messages have not gone anywhere, as the analytics group will show in a moment. Kafka records that billing has read them.

Add two orders and measure the lag

Now publish orders 7 and 8 while no consumer is running:

uv run producer.py --start 7 --count 2
sent order 8 key=ben -> partition 2 offset 2
sent order 7 key=ana -> partition 0 offset 2

Before reading them, ask Kafka where billing stands:

make lag
docker compose exec kafka /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group billing

Consumer group 'billing' has no active members.

GROUP           TOPIC           PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG             CONSUMER-ID     HOST            CLIENT-ID
billing         orders          0          2               3               1               -               -               -
billing         orders          1          2               2               0               -               -               -
billing         orders          2          2               3               1               -               -               -

This table is the group’s saved position, straight from the broker. CURRENT-OFFSET is the committed offset, the next message billing will read. LOG-END-OFFSET is the offset the next new message will get. The difference is LAG: how many messages are waiting for this group. Partitions 0 and 2 each received one new order; partition 1 received none. “No active members” means no billing consumer is connected right now, and the dashes in the last three columns are empty for the same reason. The committed offsets persist with no member connected.

Resume where it stopped

Read as billing again:

make consume
uv run consumer.py --group billing
got order 7 customer=ana amount_cents=2750 (partition 0 offset 2)
got order 8 customer=ben amount_cents=3000 (partition 2 offset 2)
2 order(s) read by group 'billing'

Exactly the two new orders, starting at offset 2 in each partition that had lag. A restarted consumer continues from its committed offsets rather than from the beginning or the end.

A second group replays everything

Finally, read as a group that has never connected:

make consume GROUP=analytics
uv run consumer.py --group analytics
got order 2 customer=ben amount_cents=1500 (partition 2 offset 0)
got order 5 customer=ben amount_cents=2250 (partition 2 offset 1)
got order 8 customer=ben amount_cents=3000 (partition 2 offset 2)
got order 3 customer=gia amount_cents=1750 (partition 1 offset 0)
got order 6 customer=gia amount_cents=2500 (partition 1 offset 1)
got order 1 customer=ana amount_cents=1250 (partition 0 offset 0)
got order 4 customer=ana amount_cents=2000 (partition 0 offset 1)
got order 7 customer=ana amount_cents=2750 (partition 0 offset 2)
8 order(s) read by group 'analytics'

All eight orders, every one of which billing had already read. With no committed offset, auto.offset.reset: earliest started analytics at offset 0 of every partition. In a job queue, the first reader would have removed these messages. In Kafka they stay in the log until the topic’s retention period expires (seven days by default), and any number of groups can read them at their own pace.

Run make lag once more and billing shows LAG 0 on all three partitions (offsets 3, 2, 3). The analytics run moved only its own group’s offsets.

Step 9: Test the round trip with pytest

Step 8 showed the behavior by eye. The tests check it on every run: that everything produced comes back, that one key stays on one partition in order, and that a group resumes while a new group replays. They run against the real broker, so each test creates its own throwaway topic and uses fresh group names. That keeps them from reading or disturbing the orders data you just made, and from interfering with each other.

Create the file

mkdir -p tests
touch tests/test_orders_stream.py

Add the code: tests/test_orders_stream.py

"""Integration tests: these talk to the broker started by `make up`."""

import uuid

import pytest
from confluent_kafka import KafkaException
from confluent_kafka.admin import AdminClient, NewTopic

from config import BOOTSTRAP_SERVERS
from consumer import consume_orders
from producer import make_order, produce_orders


@pytest.fixture(scope="session")
def admin():
    """Connect once per test run, and skip every test if the broker is down."""
    client = AdminClient({"bootstrap.servers": BOOTSTRAP_SERVERS})
    try:
        client.list_topics(timeout=5)
    except KafkaException:
        pytest.skip(f"no Kafka broker at {BOOTSTRAP_SERVERS}; run `make up` first")
    return client


@pytest.fixture
def topic(admin):
    """Create a throwaway three-partition topic so tests never touch `orders`."""
    name = f"test-orders-{uuid.uuid4().hex[:8]}"
    admin.create_topics([NewTopic(name, num_partitions=3, replication_factor=1)])[name].result()
    yield name
    admin.delete_topics([name])[name].result()


def new_group() -> str:
    return f"test-group-{uuid.uuid4().hex[:8]}"


def test_every_order_comes_back(topic):
    orders = [make_order(i) for i in range(1, 7)]
    produce_orders(orders, topic=topic)
    received = consume_orders(new_group(), topic=topic)
    assert sorted(o["order_id"] for o in received) == [1, 2, 3, 4, 5, 6]


def test_same_key_lands_on_one_partition_in_order(topic):
    placements = produce_orders([make_order(i) for i in (1, 4, 7)], topic=topic)
    assert len({partition for partition, _ in placements}) == 1
    received = consume_orders(new_group(), topic=topic)
    assert [o["order_id"] for o in received] == [1, 4, 7]


def test_group_resumes_and_new_group_replays(topic):
    produce_orders([make_order(i) for i in range(1, 4)], topic=topic)
    group = new_group()
    assert len(consume_orders(group, topic=topic)) == 3
    assert consume_orders(group, topic=topic) == []
    assert len(consume_orders(new_group(), topic=topic)) == 3

Detailed breakdown

  • The admin fixture uses Kafka’s AdminClient, the programmatic equivalent of kafka-topics.sh. list_topics(timeout=5) is a cheap request that raises KafkaException if no broker answers, and the fixture turns that into a skip with a hint. It is session-scoped, so a stopped broker costs one 5-second wait instead of one per test.
  • The topic fixture creates a uniquely named topic per test and deletes it afterward. create_topics returns a future per topic, and .result() blocks until the broker confirms, so the test never races topic creation.
  • new_group() gives each consumer run a group with no committed offsets, which is what makes auto.offset.reset: earliest apply.
  • test_every_order_comes_back compares sorted order IDs, because Step 8 showed that the interleaving across partitions is not fixed.
  • test_same_key_lands_on_one_partition_in_order sends three of Ana’s orders, checks they all landed on one partition, then reads them back and checks they arrive as orders 1, 4, 7. Comparing offsets from the delivery reports would not do: reports for one partition always arrive in offset order, so that check could never fail. It does not assert which partition, since that depends on the hash and the partition count, not on anything the test controls.
  • test_group_resumes_and_new_group_replays is Step 8 in miniature: a group reads three messages, reads nothing the second time, and a new group reads all three again.
  • Speed. Every consume_orders call ends with the 5-second idle wait, so the five consumer runs across the suite put a floor of about 25 seconds under it. Keeping the idle timeout in one place makes that trade-off easy to change.

The tests import config, producer, and consumer from the project root, so tell pytest to put that directory on the import path by appending a section to pyproject.toml:

cat >> pyproject.toml <<'EOF'

[tool.pytest.ini_options]
pythonpath = ["."]
testpaths = ["tests"]
EOF

Run the suite with the broker still up:

make test
uv run pytest -v
============================= test session starts ==============================
platform darwin -- Python 3.13.14, pytest-9.1.1, pluggy-1.6.0
configfile: pyproject.toml
testpaths: tests
collecting ... collected 3 items

tests/test_orders_stream.py::test_every_order_comes_back PASSED          [ 33%]
tests/test_orders_stream.py::test_same_key_lands_on_one_partition_in_order PASSED [ 66%]
tests/test_orders_stream.py::test_group_resumes_and_new_group_replays PASSED [100%]

============================== 3 passed in 30.11s ==============================

The header lines also print your virtual-environment path, cache directory, and project root, trimmed above. When you are done, make down stops the broker and deletes everything it stored.

Troubleshooting

  • Connect to ipv4#127.0.0.1:9092 failed: Connection refused (or ipv6#[::1]), repeated, then RuntimeError: delivery failed: 6 still queued. No broker is listening. Run make up and wait for Healthy. The %3|...|FAIL| lines come from librdkafka, which logs connection attempts to stderr and keeps retrying; the RuntimeError is produce_orders giving up after flush(10), with the number of orders it could not send. A final TERMINATE ... use flush() warning follows because the program exits with those orders still buffered.
  • All six orders land on partition 0. The orders topic was created with one partition. When a producer writes to a topic that does not exist, the broker creates it automatically with the default of one partition. This happens if you run make produce after make down and make up but before make topic. Running make topic afterwards does not help: --if-not-exists sees the one-partition topic and leaves it alone. Run make down, make up, make topic, then produce again.
  • The consumer prints 0 order(s). Either nothing has been produced since the last make down, or that group already committed its offsets. Try a new group name: make consume GROUP=scratch.
  • Clients connect once, then hang or fail with a hostname error. The advertised listener points somewhere your Mac cannot reach. Check that KAFKA_ADVERTISED_LISTENERS is PLAINTEXT://localhost:9092 and that the ports mapping is 9092:9092.
  • Bind for 0.0.0.0:9092 failed: port is already allocated. Another Kafka, often from an earlier tutorial or a Homebrew service, holds port 9092. Stop it (docker ps to find a container, brew services list for Homebrew) or change both the port mapping and the advertised listener to another port.

Recap

You ran a single Kafka 4.3.1 broker in KRaft mode with no ZooKeeper, and learned why its advertised listener has to be an address your Mac can reach. You created a topic with three partitions, published orders whose customer key kept each customer’s orders together and in sequence, and read them as a consumer group that committed its position back to Kafka. Then you watched the difference between a log and a queue directly: the same group resumed from its committed offsets, kafka-consumer-groups.sh showed its lag, and a second group replayed every message from offset 0.

Next steps:

  • Run several consumers in the same group and watch Kafka divide the three partitions between them, and reassign them when one stops.
  • Replace automatic commits with manual consumer.commit() calls after processing, and work out what happens when a consumer crashes between the two.
  • Add a named volume to compose.yaml so the broker’s data survives docker compose down. Point KAFKA_LOG_DIRS at the mount too: the broker writes to a directory under /tmp by default, not to the image’s declared data volume.