An event-driven service with a consumer group
Build a small order-processing service in Python: key messages by customer, share a topic between two service instances with a consumer group, and see them take over from each other and resume.
You will build a small order-processing service. A producer writes order events to a topic, keyed by customer. Two copies of the service share the topic through a consumer group. You stop one and watch the other take over, then start it again and see it resume where it left off.
Before we start
You need Docker, Python 3.10 to 3.13, and the Python SDK in a virtual environment. On Linux the server needs kernel 5.19 or newer, but macOS has no kernel requirement (see System requirements).
python -m venv .venv
source .venv/bin/activate
pip install apache-iggyConsumer groups need a binary transport, so this tutorial uses TCP. HTTP cannot join a group (see Consumer groups and the FAQ).
Start the server
docker run --rm --name iggy \
--cap-add=SYS_NICE --security-opt seccomp=unconfined --ulimit memlock=-1:-1 \
-p 8090:8090 \
-e IGGY_TCP_ADDRESS=0.0.0.0:8090 \
-e IGGY_NODE_ADVERTISED_ADDRESS=localhost \
-e IGGY_ROOT_USERNAME=iggy -e IGGY_ROOT_PASSWORD=iggy \
-e IGGY_SHARDING_CPU_ALLOCATION=all \
apache/iggy:0.9.0The flags are explained on Docker & Helm. IGGY_SHARDING_CPU_ALLOCATION=all turns off NUMA binding, which Docker Desktop does not support.
The producer
Save as orders_producer.py.
import asyncio
import json
import os
import random
from apache_iggy import IggyClient, Partitioning, SendMessage
CONNECTION_STRING = os.environ.get(
"IGGY_CONNECTION_STRING", "iggy+tcp://iggy:iggy@127.0.0.1:8090"
)
STREAM = "shop"
TOPIC = "orders"
PARTITIONS = 3
CUSTOMERS = ["ada", "bob", "cy", "dee", "eve"]
async def main():
client = IggyClient.from_connection_string(CONNECTION_STRING)
await client.connect()
# Create only what is missing, so the producer can be run again.
if await client.get_stream(STREAM) is None:
await client.create_stream(name=STREAM)
if await client.get_topic(STREAM, TOPIC) is None:
await client.create_topic(
stream=STREAM, name=TOPIC, partitions_count=PARTITIONS
)
for order_id in range(1, 16):
customer = random.choice(CUSTOMERS)
event = {
"order_id": order_id,
"customer": customer,
"total": round(random.uniform(5, 200), 2),
}
# One partition per customer: the key is hashed, so every order
# from the same customer lands on the same partition, in order.
response = await client.send_messages(
stream=STREAM,
topic=TOPIC,
partitioning=Partitioning.messages_key(customer),
messages=[SendMessage(json.dumps(event))],
)
confirmation = response.confirmations[0]
print(
f"order {order_id:>2} customer={customer:<3} "
f"-> partition {confirmation.partition_id} offset {confirmation.base_offset}"
)
asyncio.run(main())Partitioning.messages_key hashes the key to pick the partition, so orders for one customer share a partition and stay in order. The connection string comes from the environment, so the same script works against a server on another host (see Connection strings).
Run it:
python orders_producer.pyThis is the output from one run. Yours will differ because customers are picked at random.
order 1 customer=eve -> partition 2 offset 0
order 2 customer=dee -> partition 1 offset 0
order 3 customer=dee -> partition 1 offset 1
order 4 customer=eve -> partition 2 offset 1
order 5 customer=eve -> partition 2 offset 2
order 6 customer=bob -> partition 2 offset 3
order 7 customer=ada -> partition 2 offset 4
order 8 customer=dee -> partition 1 offset 2
order 9 customer=cy -> partition 0 offset 0
order 10 customer=bob -> partition 2 offset 5
order 11 customer=eve -> partition 2 offset 6
order 12 customer=bob -> partition 2 offset 7
order 13 customer=cy -> partition 0 offset 1
order 14 customer=dee -> partition 1 offset 3
order 15 customer=dee -> partition 1 offset 4The service
Save as order_service.py.
import asyncio
import json
import os
import signal
import sys
from apache_iggy import (
AutoCommit,
AutoCommitAfter,
IggyClient,
PollingStrategy,
ReceiveMessage,
)
CONNECTION_STRING = os.environ.get(
"IGGY_CONNECTION_STRING", "iggy+tcp://iggy:iggy@127.0.0.1:8090"
)
STREAM = "shop"
TOPIC = "orders"
GROUP = "order-processors"
INSTANCE = sys.argv[1] if len(sys.argv) > 1 else "service"
async def main():
client = IggyClient.from_connection_string(CONNECTION_STRING)
await client.connect()
# Joins the group (creating it if needed). The server assigns this
# member a share of the partitions; Next() resumes after the group's
# stored offset, and the offset is stored after each message is handled.
consumer = await client.consumer_group(
GROUP,
STREAM,
TOPIC,
polling_strategy=PollingStrategy.Next(),
auto_commit=AutoCommit.After(AutoCommitAfter.ConsumingEachMessage()),
)
shutdown = asyncio.Event()
loop = asyncio.get_running_loop()
for sig in (signal.SIGINT, signal.SIGTERM):
loop.add_signal_handler(sig, shutdown.set)
async def handle(message: ReceiveMessage) -> None:
event = json.loads(message.payload())
print(
f"[{INSTANCE}] partition {message.partition_id()} offset {message.offset()}: "
f"order {event['order_id']} for {event['customer']} total {event['total']}",
flush=True,
)
print(f"[{INSTANCE}] joined group {GROUP}, waiting for orders", flush=True)
await consumer.consume_messages(handle, shutdown)
print(f"[{INSTANCE}] shutting down", flush=True)
asyncio.run(main())consume_messages runs until the shutdown event is set. Ctrl-C, or a SIGTERM from a supervisor, sets it, and the current message finishes before the call returns. The instance name is only used in the log.
AutoCommit.After stores the offset once the handler returns, even if it raised an exception. Handle errors inside handle, or a failed message is skipped rather than retried.
Run it
One instance. Start the service in a second terminal:
python order_service.py svc-aIt is the only member, so it gets all three partitions and drains what the producer wrote:
[svc-a] joined group order-processors, waiting for orders
[svc-a] partition 0 offset 0: order 9 for cy total 24.91
[svc-a] partition 0 offset 1: order 13 for cy total 41.86
[svc-a] partition 1 offset 0: order 2 for dee total 56.72
[svc-a] partition 1 offset 1: order 3 for dee total 103.83
[svc-a] partition 1 offset 2: order 8 for dee total 99.38
[svc-a] partition 1 offset 3: order 14 for dee total 13.7
[svc-a] partition 1 offset 4: order 15 for dee total 42.16
[svc-a] partition 2 offset 0: order 1 for eve total 45.66
[svc-a] partition 2 offset 1: order 4 for eve total 105.5
[svc-a] partition 2 offset 2: order 5 for eve total 83.05
[svc-a] partition 2 offset 3: order 6 for bob total 42.6
[svc-a] partition 2 offset 4: order 7 for ada total 130.19
[svc-a] partition 2 offset 5: order 10 for bob total 73.41
[svc-a] partition 2 offset 6: order 11 for eve total 70.27
[svc-a] partition 2 offset 7: order 12 for bob total 13.56Scale out. In a third terminal start a second instance, then run the producer again:
python order_service.py svc-bpython orders_producer.pyThe server rebalances the group. In this run svc-b took partition 2 and svc-a kept 0 and 1, and no order was printed by both.
[svc-b] joined group order-processors, waiting for orders
[svc-b] partition 2 offset 8: order 1 for eve total 6.29
[svc-b] partition 2 offset 9: order 2 for eve total 30.1
[svc-b] partition 2 offset 10: order 3 for ada total 181.11
[svc-b] partition 2 offset 11: order 4 for ada total 82.09
[svc-b] partition 2 offset 12: order 8 for eve total 24.11
[svc-b] partition 2 offset 13: order 9 for eve total 65.43
[svc-b] partition 2 offset 14: order 10 for bob total 60.97
[svc-b] partition 2 offset 15: order 11 for eve total 41.1
[svc-b] partition 2 offset 16: order 12 for ada total 47.29
[svc-b] partition 2 offset 17: order 13 for eve total 11.01[svc-a] partition 0 offset 2: order 5 for cy total 171.77
[svc-a] partition 1 offset 5: order 7 for dee total 178.08
[svc-a] partition 0 offset 3: order 6 for cy total 152.96
[svc-a] partition 0 offset 4: order 15 for cy total 194.45
[svc-a] partition 1 offset 6: order 14 for dee total 81.75Scale in. Press Ctrl-C in the svc-a terminal. It prints [svc-a] shutting down and exits. Run the producer again. svc-b now owns every partition and receives all 15 orders (output trimmed):
[svc-b] partition 0 offset 5: order 1 for cy total 141.92
[svc-b] partition 2 offset 18: order 2 for bob total 91.52
[svc-b] partition 2 offset 19: order 3 for bob total 113.24
...
[svc-b] partition 1 offset 9: order 14 for dee total 159.5
[svc-b] partition 2 offset 25: order 12 for ada total 149.35
[svc-b] partition 2 offset 26: order 15 for ada total 81.49Resume. Start svc-a again. It prints only its joined line, because the group's stored offsets mean nothing is replayed. Run the producer once more and svc-a gets partition 2 back:
[svc-a] joined group order-processors, waiting for orders
[svc-a] partition 2 offset 27: order 1 for eve total 50.28
[svc-a] partition 2 offset 28: order 2 for eve total 188.94
...
[svc-a] partition 2 offset 37: order 15 for bob total 158.09When a partition moves, the server stops the old owner polling it before handing it to the new one, and waits up to 30 seconds by default (see How does consumer group rebalancing work?).
Running it again
The producer creates the stream and topic only if they are missing, so it can be run any number of times. A restarted service resumes from the group's stored offset. To start from nothing, stop the container, and --rm removes it along with its data.
The service is the same consumer shape as the high-level consumer example in the repository, with a JSON handler and signal handling added.
