Skip to content
Merged
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
213 changes: 213 additions & 0 deletions content/docs/introduction/getting-started-python.mdx
Original file line number Diff line number Diff line change
@@ -0,0 +1,213 @@
---
title: Getting started (Python)
description: "A first program with the Python SDK: create a stream and topic, send messages, then read them back on their own and as a consumer group."
---

## Before we start

This tutorial uses the Python SDK's low-level API (`IggyClient`) directly, so you can see each step. It needs Python 3.10 or newer. On Linux the server needs kernel 5.19 or newer; macOS has no kernel requirement. See [System requirements](/docs/server/introduction#system-requirements).

The producer and consumer below match the [getting-started examples](https://github.com/apache/iggy/tree/master/examples/python/getting-started) in the repository.

## Starting the server

```bash
docker run --rm \
--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 \
apache/iggy:0.9.0
```

`IGGY_TCP_ADDRESS` and `IGGY_NODE_ADVERTISED_ADDRESS` are needed because the server binds `127.0.0.1` inside the container by default, which the published port can't reach. The Rust tutorial's [Starting the Iggy server](/docs/introduction/getting-started#starting-the-iggy-server) section explains each flag.

On Docker Desktop, 0.9.0 can fail to start with `failed to bind shard 0 memory to its NUMA node`. If you hit this, add `-e IGGY_SHARDING_CPU_ALLOCATION=all` to the command above.

Install the SDK in a virtual environment:

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

## Producer

Save this as `producer.py`:

```python
import asyncio

from apache_iggy import IggyClient
from apache_iggy import SendMessage as Message

STREAM_NAME = "sample-stream"
TOPIC_NAME = "sample-topic"
PARTITION_ID = 0


async def main():
client = IggyClient.from_connection_string(
"iggy+tcp://iggy:[email protected]:8090"
)
await client.connect()

# Re-running this example is fine: only create what is missing.
if await client.get_stream(STREAM_NAME) is None:
await client.create_stream(name=STREAM_NAME)

if await client.get_topic(STREAM_NAME, TOPIC_NAME) is None:
await client.create_topic(
stream=STREAM_NAME,
name=TOPIC_NAME,
partitions_count=1,
)

messages = [Message(f"message-{i}") for i in range(10)]

await client.send_messages(
stream=STREAM_NAME,
topic=TOPIC_NAME,
partitioning=PARTITION_ID,
messages=messages,
)
print(f"Sent {len(messages)} message(s)")


asyncio.run(main())
```

Run it with `python producer.py`:

```text
Sent 10 message(s)
```

`from_connection_string` takes the transport, host and credentials in one string, `iggy+tcp://user:pass@host:port`. See [Connection strings](/docs/sdk/connection-strings) for the other options. Creating a stream or topic that already exists is an error, so the producer checks first and only creates what's missing. That's also what makes it safe to run again.

## Consumer

Save this as `consumer.py`:

```python
import asyncio

from apache_iggy import Consumer, IggyClient, PollingStrategy

STREAM_NAME = "sample-stream"
TOPIC_NAME = "sample-topic"
PARTITION_ID = 0


async def main():
client = IggyClient.from_connection_string(
"iggy+tcp://iggy:[email protected]:8090"
)
await client.connect()

# Next() resumes this named consumer after its stored offset. Auto-commit
# can store the offset before application processing finishes.
polled_messages = await client.poll_messages(
stream=STREAM_NAME,
topic=TOPIC_NAME,
consumer=Consumer.Single("sample-consumer"),
partition_id=PARTITION_ID,
polling_strategy=PollingStrategy.Next(),
count=10,
auto_commit=True,
)

for message in polled_messages:
payload = message.payload().decode("utf-8")
print(f"Offset: {message.offset()}, Payload: {payload}")


asyncio.run(main())
```

Run it with `python consumer.py`, after the producer:

```text
Offset: 0, Payload: message-0
Offset: 1, Payload: message-1
Offset: 2, Payload: message-2
Offset: 3, Payload: message-3
Offset: 4, Payload: message-4
Offset: 5, Payload: message-5
Offset: 6, Payload: message-6
Offset: 7, Payload: message-7
Offset: 8, Payload: message-8
Offset: 9, Payload: message-9
```

`Consumer.Single(name)` is a consumer that polls on its own, and the server stores its offset under that name. With `auto_commit=True` the offset is stored after the poll returns, without waiting for the store to finish, so a crash straight afterwards can deliver the last batch again.

## Consumer group

A consumer group lets several copies of a consumer share a topic's partitions, with each message going to only one of them. Groups need a binary transport (TCP, QUIC or WebSocket); HTTP doesn't support them.

Save this as `consumer_group.py`:

```python
import asyncio

from apache_iggy import IggyClient, PollingStrategy

STREAM_NAME = "sample-stream"
TOPIC_NAME = "sample-topic"
GROUP_NAME = "sample-consumer-group"


async def main():
client = IggyClient.from_connection_string(
"iggy+tcp://iggy:[email protected]:8090"
)
await client.connect()

# Creates the group if it doesn't exist and joins it. A group member
# reads whichever partitions the server assigns it, so there's no
# partition_id to set, unlike Consumer.Single.
group = await client.consumer_group(
name=GROUP_NAME,
stream=STREAM_NAME,
topic=TOPIC_NAME,
polling_strategy=PollingStrategy.First(),
)

# iter_messages() keeps waiting for new messages once it catches up,
# like a live subscription, so this example stops after 20.
count = 0
async for message in group.iter_messages():
payload = message.payload().decode("utf-8")
print(f"Offset: {message.offset()}, Payload: {payload}")
count += 1
if count == 20:
break
print(f"Consumed {count} message(s), stopping.")


asyncio.run(main())
```

Run it with `python consumer_group.py`. With the producer run twice, so 20 messages on the topic, it prints every offset from 0 to 19 and then stops:

```text
Offset: 0, Payload: message-0
...
Offset: 19, Payload: message-9
Consumed 20 message(s), stopping.
```

`iter_messages()` is the simplest way to read from a group. For a long-running service with its own shutdown signal, use `consume_messages(callback, shutdown_event)` instead.

## Running it again

Run the producer again and it finds the stream and topic already there, and adds another 10 messages at offsets 10 to 19.

Run `consumer.py` again and it carries on from its stored offset, printing only the new messages rather than the first ten again.

Run `consumer_group.py` again and it prints every message from offset 0, because it asks for `PollingStrategy.First()`. That keeps the example repeatable. A long-running group consumer would leave the polling strategy at its default and keep `auto_commit` on, so it carries on from the group's stored offset instead.

To start again with an empty server, stop the container, which `--rm` removes, and run `docker run` again.
4 changes: 2 additions & 2 deletions content/docs/introduction/meta.json
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
{
"title": "Introduction",
"pages": ["about", "concepts", "coming-from-kafka", "glossary", "serialising-messages", "architecture", "quickstart", "getting-started"]
}
"pages": ["about", "concepts", "coming-from-kafka", "glossary", "serialising-messages", "architecture", "quickstart", "getting-started", "getting-started-python"]
}
Loading