Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions content/docs/meta.json
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
"pages": [
"index",
"introduction",
"tutorials",
"server",
"clustering",
"binary-protocol",
Expand Down
280 changes: 280 additions & 0 deletions content/docs/tutorials/event-driven-service.mdx
Original file line number Diff line number Diff line change
@@ -0,0 +1,280 @@
---
title: An event-driven service with a consumer group
description: "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](/docs/server/introduction#system-requirements)).

```bash
python -m venv .venv
source .venv/bin/activate
pip install apache-iggy
```

Consumer groups need a binary transport, so this tutorial uses TCP. HTTP cannot join a group (see [Consumer groups](/docs/introduction/concepts#consumer-groups) and the [FAQ](/docs/faq/faq)).

## Start the server

```bash
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.0
```

The flags are explained on [Docker & Helm](/docs/server/docker). `IGGY_SHARDING_CPU_ALLOCATION=all` turns off NUMA binding, which Docker Desktop does not support.

## The producer

Save as `orders_producer.py`.

```python
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](/docs/sdk/connection-strings)).

Run it:

```bash
python orders_producer.py
```

This 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 4
```

## The service

Save as `order_service.py`.

```python
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:

```bash
python order_service.py svc-a
```

It 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.56
```

**Scale out.** In a third terminal start a second instance, then run the producer again:

```bash
python order_service.py svc-b
```

```bash
python orders_producer.py
```

The 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.75
```

**Scale 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.49
```

**Resume.** 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.09
```

When 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?](/docs/faq/faq#q-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](https://github.com/apache/iggy/blob/master/examples/python/high-level/consumer.py) in the repository, with a JSON handler and signal handling added.
4 changes: 4 additions & 0 deletions content/docs/tutorials/meta.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
{
"title": "Tutorials",
"pages": ["event-driven-service"]
}
Loading