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..85db2d3c20 --- /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, 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. 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"] +}