# 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.

Rendered page: https://iggy.apache.org/docs/tutorials/event-driven-service/

Source: https://github.com/apache/iggy-website/blob/main/content/docs/tutorials/event-driven-service.mdx

Docs version: server 0.9.0, Rust SDK 0.11.0, CLI 0.14.0, Node.js SDK 0.10.0, and the Python, Java, Go and C# SDKs at 0.9.0

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](https://iggy.apache.org/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](https://iggy.apache.org/docs/introduction/concepts#consumer-groups) and the [FAQ](https://iggy.apache.org/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](https://iggy.apache.org/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](https://iggy.apache.org/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?](https://iggy.apache.org/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.
