diff --git a/examples/python/README.md b/examples/python/README.md index 7e5da180d4..45acc31693 100644 --- a/examples/python/README.md +++ b/examples/python/README.md @@ -71,6 +71,19 @@ python basic/consumer.py Demonstrates fundamental client connection, authentication, batch message sending, and polling with support for TCP/QUIC/HTTP protocols. +### Message Partitioning + +Shows how to route message batches to a fixed partition, use server-side +balanced routing, or keep the same message key on the same partition. + +```bash +# Using uv +uv run partitioning/producer.py + +# Without using uv +python partitioning/producer.py +``` + ### Message Headers Shows how to attach and read Python SDK user headers with `str`, `bytes`, `bool`, `int`, and `float` values. Two variants share their logic through `message-headers/common.py`: diff --git a/examples/python/partitioning/producer.py b/examples/python/partitioning/producer.py new file mode 100644 index 0000000000..ffc71355d7 --- /dev/null +++ b/examples/python/partitioning/producer.py @@ -0,0 +1,74 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +import argparse +import asyncio + +from apache_iggy import IggyClient, Partitioning, SendMessage +from loguru import logger + +STREAM_NAME = "partitioning-stream" +TOPIC_NAME = "partitioning-topic" +PARTITIONS_COUNT = 3 + + +async def init_system(client: IggyClient) -> None: + if await client.get_stream(STREAM_NAME) is None: + await client.create_stream(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=PARTITIONS_COUNT, + ) + + +async def send(client: IggyClient, label: str, partitioning: Partitioning) -> None: + response = await client.send_messages( + stream=STREAM_NAME, + topic=TOPIC_NAME, + partitioning=partitioning, + messages=[SendMessage(label)], + ) + for confirmation in response.confirmations: + logger.info( + "{} was written to partition {} at offset {}", + label, + confirmation.partition_id, + confirmation.base_offset, + ) + + +async def main(connection_string: str) -> None: + client = IggyClient.from_connection_string(connection_string) + await client.connect() + await init_system(client) + + await send(client, "fixed", Partitioning.partition_id(0)) + await send(client, "balanced", Partitioning.balanced()) + await send(client, "keyed", Partitioning.messages_key(b"customer-42")) + + +if __name__ == "__main__": + parser = argparse.ArgumentParser() + parser.add_argument( + "connection_string", + nargs="?", + default="iggy+tcp://iggy:iggy@127.0.0.1:8090", + ) + args = parser.parse_args() + asyncio.run(main(args.connection_string)) diff --git a/foreign/python/apache_iggy.pyi b/foreign/python/apache_iggy.pyi index 50b65786db..4dcca9a608 100644 --- a/foreign/python/apache_iggy.pyi +++ b/foreign/python/apache_iggy.pyi @@ -1301,7 +1301,7 @@ class IggyClient: self, stream: builtins.str | builtins.int, topic: builtins.str | builtins.int, - partitioning: builtins.int, + partitioning: Partitioning | builtins.int, messages: list[SendMessage], ) -> collections.abc.Awaitable[SendMessagesResponse]: r""" @@ -1310,6 +1310,10 @@ class IggyClient: confirmations, or a PyRuntimeError on failure. The confirmation list is empty when the server reports no offsets, and the legacy server never reports any. + + `partitioning` is required. Pass `Partitioning.balanced()`, + `Partitioning.partition_id(id)`, or `Partitioning.messages_key(key)`. + An integer remains supported as shorthand for `partition_id`. """ def poll_messages( self, @@ -1601,6 +1605,30 @@ class Partition: The number of messages in the partition. """ +@typing.final +class Partitioning: + r""" + Defines how a batch of messages is assigned to a topic partition. + """ + @staticmethod + def balanced() -> Partitioning: + r""" + Routes the batch to partitions using server-side round-robin selection. + """ + @staticmethod + def partition_id(partition_id: builtins.int) -> Partitioning: + r""" + Routes the batch to the specified partition. + """ + @staticmethod + def messages_key(key: builtins.str | builtins.bytes) -> Partitioning: + r""" + Routes the batch using a binary key hashed by the server. + + String keys are encoded as UTF-8. The encoded key must contain between + 1 and 255 bytes. + """ + @typing.final class Permissions: r""" diff --git a/foreign/python/src/client.rs b/foreign/python/src/client.rs index 816ba4e61f..4e9cbf624c 100644 --- a/foreign/python/src/client.rs +++ b/foreign/python/src/client.rs @@ -39,6 +39,7 @@ use crate::consumer::{ use crate::duration::{py_delta_to_iggy_duration, reject_zero}; use crate::identifier::PyIdentifier; use crate::options::OptionSpec as PyOptionSpec; +use crate::partitioning::PyPartitioning; use crate::permissions::Permissions as PyPermissions; use crate::receive_message::{PollingStrategy, ReceiveMessage}; use crate::send_message::{SendMessage, SendMessagesResponse as PySendMessagesResponse}; @@ -1000,13 +1001,18 @@ impl IggyClient { /// confirmations, or a PyRuntimeError on failure. The confirmation list is /// empty when the server reports no offsets, and the legacy server never /// reports any. + /// + /// `partitioning` is required. Pass `Partitioning.balanced()`, + /// `Partitioning.partition_id(id)`, or `Partitioning.messages_key(key)`. + /// An integer remains supported as shorthand for `partition_id`. #[gen_stub(override_return_type(type_repr="collections.abc.Awaitable[SendMessagesResponse]", imports=("collections.abc")))] fn send_messages<'a>( &self, py: Python<'a>, stream: PyIdentifier, topic: PyIdentifier, - partitioning: u32, + #[gen_stub(override_type(type_repr = "Partitioning | builtins.int"))] + partitioning: PyPartitioning, #[gen_stub(override_type(type_repr = "list[SendMessage]"))] messages: &Bound<'_, PyList>, ) -> PyResult> { let messages: Vec = messages @@ -1023,7 +1029,7 @@ impl IggyClient { let stream = Identifier::try_from(stream)?; let topic = Identifier::try_from(topic)?; - let partitioning = Partitioning::partition_id(partitioning); + let partitioning = partitioning.into(); let inner = self.inner.clone(); future_into_py(py, async move { diff --git a/foreign/python/src/lib.rs b/foreign/python/src/lib.rs index 7910af77b9..bba4ab791f 100644 --- a/foreign/python/src/lib.rs +++ b/foreign/python/src/lib.rs @@ -21,6 +21,7 @@ mod consumer; mod duration; mod identifier; mod options; +mod partitioning; mod permissions; mod receive_message; mod send_message; @@ -36,6 +37,7 @@ use consumer::{ ConsumerGroupMember, IggyConsumer, ReceiveMessageIterator, }; use options::OptionSpec; +use partitioning::Partitioning; use permissions::{GlobalPermissions, Permissions, StreamPermissions, TopicPermissions}; use pyo3::prelude::*; use receive_message::{PollingStrategy, ReceiveMessage}; @@ -62,6 +64,7 @@ fn apache_iggy(_py: Python, m: &Bound<'_, PyModule>) -> PyResult<()> { m.add_class::()?; m.add_class::()?; m.add_class::()?; + m.add_class::()?; m.add_class::()?; m.add_class::()?; m.add_class::()?; diff --git a/foreign/python/src/partitioning.rs b/foreign/python/src/partitioning.rs new file mode 100644 index 0000000000..d403b7dd5f --- /dev/null +++ b/foreign/python/src/partitioning.rs @@ -0,0 +1,93 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use iggy::prelude::Partitioning as RustPartitioning; +use pyo3::{exceptions::PyValueError, prelude::*, types::PyBytes}; +use pyo3_stub_gen::{ + derive::{gen_stub_pyclass, gen_stub_pymethods}, + impl_stub_type, +}; + +/// Defines how a batch of messages is assigned to a topic partition. +#[derive(Clone)] +#[pyclass(from_py_object)] +#[gen_stub_pyclass] +pub struct Partitioning { + pub(crate) inner: RustPartitioning, +} + +#[gen_stub_pymethods] +#[pymethods] +impl Partitioning { + /// Routes the batch to partitions using server-side round-robin selection. + #[staticmethod] + pub fn balanced() -> Self { + Self { + inner: RustPartitioning::balanced(), + } + } + + /// Routes the batch to the specified partition. + #[staticmethod] + pub fn partition_id(partition_id: u32) -> Self { + Self { + inner: RustPartitioning::partition_id(partition_id), + } + } + + /// Routes the batch using a binary key hashed by the server. + /// + /// String keys are encoded as UTF-8. The encoded key must contain between + /// 1 and 255 bytes. + #[staticmethod] + pub fn messages_key(py: Python<'_>, key: PyMessagesKey) -> PyResult { + let key = match key { + PyMessagesKey::String(key) => key.into_bytes(), + PyMessagesKey::Bytes(key) => key.extract::>(py)?, + }; + let inner = RustPartitioning::messages_key(&key) + .map_err(|error| PyValueError::new_err(error.to_string()))?; + Ok(Self { inner }) + } +} + +#[derive(FromPyObject)] +pub enum PyMessagesKey { + #[pyo3(transparent, annotation = "str")] + String(String), + #[pyo3(transparent, annotation = "bytes")] + Bytes(Py), +} +impl_stub_type!(PyMessagesKey = String | PyBytes); + +#[derive(FromPyObject)] +pub(crate) enum PyPartitioning { + #[pyo3(transparent)] + Strategy(Partitioning), + #[pyo3(transparent, annotation = "int")] + PartitionId(u32), +} +impl_stub_type!(PyPartitioning = Partitioning | isize); + +impl From for RustPartitioning { + fn from(partitioning: PyPartitioning) -> Self { + match partitioning { + PyPartitioning::Strategy(partitioning) => partitioning.inner, + PyPartitioning::PartitionId(partition_id) => Self::partition_id(partition_id), + } + } +} diff --git a/foreign/python/tests/test_message_operations.py b/foreign/python/tests/test_message_operations.py index d9128a0e11..ca85f0686b 100644 --- a/foreign/python/tests/test_message_operations.py +++ b/foreign/python/tests/test_message_operations.py @@ -24,12 +24,53 @@ HeaderKey, HeaderValue, IggyClient, + Partitioning, PollingStrategy, UserHeaders, ) from apache_iggy import SendMessage as Message +class TestPartitioning: + """Test message partitioning strategy construction.""" + + @pytest.mark.unit + def test_balanced_and_partition_id_strategies(self): + assert isinstance(Partitioning.balanced(), Partitioning) + assert isinstance(Partitioning.partition_id(1), Partitioning) + + @pytest.mark.unit + @pytest.mark.parametrize("key", [b"customer-42", "customer-42", "客户-42"]) + def test_messages_key_accepts_bytes_and_strings(self, key): + assert isinstance(Partitioning.messages_key(key), Partitioning) + + @pytest.mark.unit + @pytest.mark.parametrize("key", [b"a" * 255, "a" * 255]) + def test_messages_key_accepts_255_bytes(self, key): + assert isinstance(Partitioning.messages_key(key), Partitioning) + + @pytest.mark.unit + @pytest.mark.parametrize("key", [b"", "", b"a" * 256, "a" * 256, "界" * 86]) + def test_messages_key_rejects_invalid_encoded_length(self, key): + with pytest.raises(ValueError): + Partitioning.messages_key(key) + + @pytest.mark.unit + @pytest.mark.parametrize("partition_id", [-1, 2**32]) + def test_partition_id_rejects_values_outside_u32(self, partition_id): + with pytest.raises(OverflowError): + Partitioning.partition_id(partition_id) + + @pytest.mark.unit + def test_partitioning_rejects_invalid_types(self): + with pytest.raises(TypeError): + # pyrefly: ignore # bad-argument-type + Partitioning.partition_id("0") + with pytest.raises(TypeError): + # pyrefly: ignore # bad-argument-type + Partitioning.messages_key(1) + + class TestMessageOperations: """Test message sending, polling, and processing.""" @@ -109,6 +150,87 @@ async def test_send_messages_reports_committed_confirmation( assert confirmation.partition_id == partition_id assert confirmation.base_offset == polled_messages[0].offset() + @pytest.mark.asyncio + async def test_send_messages_with_partition_id_strategy( + self, iggy_client: IggyClient, unique_name + ): + stream_name = unique_name() + topic_name = unique_name() + partition_id = 2 + + await iggy_client.create_stream(stream_name) + await iggy_client.create_topic( + stream=stream_name, name=topic_name, partitions_count=3 + ) + + response = await iggy_client.send_messages( + stream=stream_name, + topic=topic_name, + partitioning=Partitioning.partition_id(partition_id), + messages=[Message("fixed partition")], + ) + + assert len(response.confirmations) == 1 + assert response.confirmations[0].partition_id == partition_id + + @pytest.mark.asyncio + async def test_send_messages_with_balanced_strategy( + self, iggy_client: IggyClient, unique_name + ): + stream_name = unique_name() + topic_name = unique_name() + partitions_count = 3 + + await iggy_client.create_stream(stream_name) + await iggy_client.create_topic( + stream=stream_name, + name=topic_name, + partitions_count=partitions_count, + ) + + response = await iggy_client.send_messages( + stream=stream_name, + topic=topic_name, + partitioning=Partitioning.balanced(), + messages=[Message("balanced")], + ) + + assert len(response.confirmations) == 1 + assert response.confirmations[0].partition_id < partitions_count + + @pytest.mark.asyncio + @pytest.mark.parametrize("key", [b"customer-42", "customer-42"]) + async def test_send_messages_with_same_key_uses_same_partition( + self, iggy_client: IggyClient, unique_name, key + ): + stream_name = unique_name() + topic_name = unique_name() + + await iggy_client.create_stream(stream_name) + await iggy_client.create_topic( + stream=stream_name, name=topic_name, partitions_count=3 + ) + partitioning = Partitioning.messages_key(key) + + first = await iggy_client.send_messages( + stream=stream_name, + topic=topic_name, + partitioning=partitioning, + messages=[Message("first")], + ) + second = await iggy_client.send_messages( + stream=stream_name, + topic=topic_name, + partitioning=partitioning, + messages=[Message("second")], + ) + + assert len(first.confirmations) == 1 + assert len(second.confirmations) == 1 + assert ( + first.confirmations[0].partition_id == second.confirmations[0].partition_id + ) + @pytest.mark.asyncio async def test_send_and_poll_messages_as_bytes( self, iggy_client: IggyClient, unique_name