If Apache Iggy is useful to you, give it a star on GitHubStar apache/iggy
Apache Iggy
Introduction

Getting started (Python)

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.

The producer and consumer below match the getting-started examples in the repository.

Starting the server

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

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

Producer

Save this as producer.py:

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:iggy@127.0.0.1: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:

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

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:iggy@127.0.0.1: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:

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:

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:iggy@127.0.0.1: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:

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.

On this page