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.yamlthat runs a single Kafka 4.3.1 broker in KRaft mode onlocalhost:9092. - A topic named
orderswith 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
Makefilewhose default target prints a help screen, and apytestsuite 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.yamlfile 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 fromuv.lockon any machine withuv sync.__pycache__/,*.py[cod], and.pytest_cache/are bytecode and test caches that Python andpytestregenerate 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 floatinglatesttag 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), andKAFKA_CONTROLLER_QUORUM_VOTERSlists the quorum as just itself: node1, reachable on port9093. - Listeners versus advertised listeners. A listener is a named network
endpoint the broker accepts connections on.
KAFKA_LISTENERSopens two: thePLAINTEXTlistener on port 9092 for clients, and theCONTROLLERlistener 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. Advertisinglocalhost:9092works here because theportsmapping 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_MAPmaps each listener name to a protocol. Both are unencryptedPLAINTEXT, 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
1is 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 --waitreturn 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_SERVERSis 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.TOPICdefaults to theorderstopic 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_orderbuilds ordernfrom 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, andcarolall hash to partition 2.ana,ben, andgiawere 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.librdkafkabatches buffered messages per partition and sends them from a background thread. That is why the result arrives later, through theon_deliverycallback.- 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 afterflushrather than inside the callback, so a failure surfaces as one clear exception fromproduce_orders.- Caveat for long-running producers: callbacks only run during
flush()orpoll(). A producer that sends continuously should callproducer.poll(0)inside its loop, or its callbacks pile up unprocessed. --startand--countlet 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.idnames the consumer group. Thebillingdefault 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: earliestonly 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 firstpoll()can take longer than later ones.poll(idle_seconds)returns one message, orNoneif 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 frommakeand 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.commitis on by default), andconsumer.close()commits the final position before leaving the group. Leaving outclose()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 := helpmakes a baremakeprint the help screen. Thehelprecipe finds every line of the formtarget: ## descriptionand prints the pairs, so adding a target with a##comment documents it.KAFKA_BINandBOOTSTRAPhold the repeated prefix for the bundled tools. They run inside the container, wherelocalhost:9092is the broker itself.GROUP ?= billingsets a default that you can override per command, as inmake consume GROUP=analytics.downdeletes data. The broker writes its log inside the container, anddocker compose downremoves the container. The-valso removes the anonymous volumes the image declares, so repeated runs do not leave unused volumes behind. Everymake upafter amake downstarts with no topics and no messages.--if-not-existsmakesmake topicsafe to repeat on a broker that already has the topic.- Makefile recipes must be indented with a tab, not spaces. If
makereportsmissing 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
adminfixture uses Kafka’sAdminClient, the programmatic equivalent ofkafka-topics.sh.list_topics(timeout=5)is a cheap request that raisesKafkaExceptionif 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
topicfixture creates a uniquely named topic per test and deletes it afterward.create_topicsreturns 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 makesauto.offset.reset: earliestapply.test_every_order_comes_backcompares sorted order IDs, because Step 8 showed that the interleaving across partitions is not fixed.test_same_key_lands_on_one_partition_in_ordersends 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_replaysis 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_orderscall 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(oripv6#[::1]), repeated, thenRuntimeError: delivery failed: 6 still queued. No broker is listening. Runmake upand wait forHealthy. The%3|...|FAIL|lines come fromlibrdkafka, which logs connection attempts to stderr and keeps retrying; theRuntimeErrorisproduce_ordersgiving up afterflush(10), with the number of orders it could not send. A finalTERMINATE ... use flush()warning follows because the program exits with those orders still buffered.- All six orders land on
partition 0. Theorderstopic 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 runmake produceaftermake downandmake upbut beforemake topic. Runningmake topicafterwards does not help:--if-not-existssees the one-partition topic and leaves it alone. Runmake down,make up,make topic, then produce again. - The consumer prints
0 order(s). Either nothing has been produced since the lastmake 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_LISTENERSisPLAINTEXT://localhost:9092and that theportsmapping is9092: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 psto find a container,brew services listfor 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.yamlso the broker’s data survivesdocker compose down. PointKAFKA_LOG_DIRSat the mount too: the broker writes to a directory under/tmpby default, not to the image’s declared data volume.