From ba1ddb14c09c24353506aef9a595b24194a326a7 Mon Sep 17 00:00:00 2001 From: Justin Mclean Date: Tue, 29 Sep 2026 14:06:18 +1000 Subject: [PATCH 1/2] docs: add a Tutorials section with an event-driven service tutorial A Python order-processing service that shares a topic through a consumer group, scales out and in, and resumes from stored offsets. Run end to end against 0.9.0. --- content/docs/meta.json | 1 + .../docs/tutorials/event-driven-service.mdx | 280 ++++++++++++++++++ content/docs/tutorials/meta.json | 4 + 3 files changed, 285 insertions(+) create mode 100644 content/docs/tutorials/event-driven-service.mdx create mode 100644 content/docs/tutorials/meta.json diff --git a/content/docs/meta.json b/content/docs/meta.json index 3c2bc42b88..f70313d5e8 100644 --- a/content/docs/meta.json +++ b/content/docs/meta.json @@ -3,6 +3,7 @@ "pages": [ "index", "introduction", + "tutorials", "server", "clustering", "binary-protocol", diff --git a/content/docs/tutorials/event-driven-service.mdx b/content/docs/tutorials/event-driven-service.mdx new file mode 100644 index 0000000000..5dd030a32c --- /dev/null +++ b/content/docs/tutorials/event-driven-service.mdx @@ -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; 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. TCP is used here. 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. Orders for one customer stay in order because they share a partition. 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 +``` + +Output from one run. Yours will differ because customers are 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; the current message finishes and the call returns. The instance name is only for the log. + +`AutoCommit.After` stores the offset once the handler returns, whether or not it raised, so handle errors inside `handle`. + +## 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. In this run `svc-b` took partition 2 and `svc-a` kept 0 and 1. 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: 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 +``` + +Rebalancing is cooperative. When a partition moves, the server stops the previous owner polling it, then hands it over, with a timeout that defaults to 30 seconds. 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; `--rm` removes it and 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. diff --git a/content/docs/tutorials/meta.json b/content/docs/tutorials/meta.json new file mode 100644 index 0000000000..e5fb4e6d6f --- /dev/null +++ b/content/docs/tutorials/meta.json @@ -0,0 +1,4 @@ +{ + "title": "Tutorials", + "pages": ["event-driven-service"] +} From 2ec5aa70c1eb5766db518d015f5b71b1dc4fadcb Mon Sep 17 00:00:00 2001 From: Justin Mclean Date: Wed, 30 Sep 2026 15:32:20 +1000 Subject: [PATCH 2/2] docs: tidy the event-driven tutorial wording Replace semicolons and short fragments, and say what happens when the handler raises. --- .../docs/tutorials/event-driven-service.mdx | 20 +++++++++---------- 1 file changed, 10 insertions(+), 10 deletions(-) diff --git a/content/docs/tutorials/event-driven-service.mdx b/content/docs/tutorials/event-driven-service.mdx index 5dd030a32c..85db2d3c20 100644 --- a/content/docs/tutorials/event-driven-service.mdx +++ b/content/docs/tutorials/event-driven-service.mdx @@ -7,7 +7,7 @@ You will build a small order-processing service. A producer writes order events ## 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; macOS has no kernel requirement. See [System requirements](/docs/server/introduction#system-requirements). +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 @@ -15,7 +15,7 @@ source .venv/bin/activate pip install apache-iggy ``` -Consumer groups need a binary transport. TCP is used here. HTTP cannot join a group; see [Consumer groups](/docs/introduction/concepts#consumer-groups) and the [FAQ](/docs/faq/faq). +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 @@ -90,7 +90,7 @@ async def main(): asyncio.run(main()) ``` -`Partitioning.messages_key` hashes the key to pick the partition. Orders for one customer stay in order because they share a partition. 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). +`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: @@ -98,7 +98,7 @@ Run it: python orders_producer.py ``` -Output from one run. Yours will differ because customers are random. +This is the output from one run. Yours will differ because customers are picked at random. ``` order 1 customer=eve -> partition 2 offset 0 @@ -182,9 +182,9 @@ async def main(): asyncio.run(main()) ``` -`consume_messages` runs until the shutdown event is set. Ctrl-C, or a SIGTERM from a supervisor, sets it; the current message finishes and the call returns. The instance name is only for the log. +`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, whether or not it raised, so handle errors inside `handle`. +`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 @@ -225,7 +225,7 @@ python order_service.py svc-b python orders_producer.py ``` -The server rebalances. In this run `svc-b` took partition 2 and `svc-a` kept 0 and 1. No order was printed by both. +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 @@ -261,7 +261,7 @@ The server rebalances. In this run `svc-b` took partition 2 and `svc-a` kept 0 a [svc-b] partition 2 offset 26: order 15 for ada total 81.49 ``` -**Resume.** Start `svc-a` again. It prints only its `joined` line: the group's stored offsets mean nothing is replayed. Run the producer once more and `svc-a` gets partition 2 back: +**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 @@ -271,10 +271,10 @@ The server rebalances. In this run `svc-b` took partition 2 and `svc-a` kept 0 a [svc-a] partition 2 offset 37: order 15 for bob total 158.09 ``` -Rebalancing is cooperative. When a partition moves, the server stops the previous owner polling it, then hands it over, with a timeout that defaults to 30 seconds. See [How does consumer group rebalancing work?](/docs/faq/faq#q-how-does-consumer-group-rebalancing-work). +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; `--rm` removes it and its data. +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.