Skip to content

Custom Brokers

Learn how to create custom broker adapters for Protean to integrate with any message broker or streaming platform.

Overview

Custom brokers allow you to:

  • Integrate new messaging technologies (Kafka, RabbitMQ, AWS SQS, etc.)
  • Add company-specific messaging systems
  • Create specialized brokers for testing or development

Architecture

All brokers must inherit from BaseBroker and implement its abstract methods. They are these, plus the capabilities property:

  • Publishing and reading: _publish, _get_next, _read
  • Acknowledgment: _ack, _nack
  • Consumer groups: _ensure_group
  • Connection and health: _ping, _health_stats, _ensure_connection, _info, _data_reset
  • Dead letter queue: _dlq_list, _dlq_inspect, _dlq_replay, _dlq_replay_all, _dlq_purge

A broker that does not support a feature still implements the method. It returns an empty result, as the Redis PubSub broker does for the dead letter queue methods.

The broker below implements every method by handing it to an in-memory InlineBroker. Replace those calls with calls to your messaging client:

from typing import Any

from protean import Domain
from protean.adapters.broker.inline import InlineBroker
from protean.port.broker import BaseBroker, BrokerCapabilities, DLQEntry, registry


class DelegatingBroker(BaseBroker):
    """A broker that hands every operation to an in-memory InlineBroker.

    Replace the delegate calls with calls to your messaging client.
    """

    def __init__(self, name: str, domain: "Domain", conn_info: dict[str, Any]) -> None:
        super().__init__(name, domain, conn_info)
        # Initialize your broker connection here
        self._delegate = InlineBroker(name, domain, conn_info)

    @property
    def capabilities(self) -> BrokerCapabilities:
        """Declare only the capabilities the broker implements."""
        return self._delegate.capabilities

    # Publishing and reading
    def _publish(self, stream: str, message: dict[str, Any]) -> str:
        return self._delegate._publish(stream, message)

    def _get_next(
        self, stream: str, consumer_group: str
    ) -> tuple[str, dict[str, Any]] | None:
        return self._delegate._get_next(stream, consumer_group)

    def _read(
        self, stream: str, consumer_group: str, no_of_messages: int
    ) -> list[tuple[str, dict[str, Any]]]:
        return self._delegate._read(stream, consumer_group, no_of_messages)

    # Acknowledgment
    def _ack(self, stream: str, identifier: str, consumer_group: str) -> bool:
        return self._delegate._ack(stream, identifier, consumer_group)

    def _nack(self, stream: str, identifier: str, consumer_group: str) -> bool:
        return self._delegate._nack(stream, identifier, consumer_group)

    # Consumer groups
    def _ensure_group(self, group_name: str, stream: str | None = None) -> None:
        self._delegate._ensure_group(group_name, stream)

    # Connection and health
    def _ping(self) -> bool:
        return self._delegate._ping()

    def _health_stats(self) -> dict[str, Any]:
        return self._delegate._health_stats()

    def _ensure_connection(self) -> bool:
        return self._delegate._ensure_connection()

    def _info(self) -> dict[str, Any]:
        return self._delegate._info()

    def _data_reset(self) -> None:
        self._delegate._data_reset()

    # Dead letter queue
    def _dlq_list(self, dlq_streams: list[str], limit: int) -> list[DLQEntry]:
        return self._delegate._dlq_list(dlq_streams, limit)

    def _dlq_inspect(self, dlq_stream: str, dlq_id: str) -> DLQEntry | None:
        return self._delegate._dlq_inspect(dlq_stream, dlq_id)

    def _dlq_replay(self, dlq_stream: str, dlq_id: str, target_stream: str) -> bool:
        return self._delegate._dlq_replay(dlq_stream, dlq_id, target_stream)

    def _dlq_replay_all(self, dlq_stream: str, target_stream: str) -> int:
        return self._delegate._dlq_replay_all(dlq_stream, target_stream)

    def _dlq_purge(self, dlq_stream: str) -> int:
        return self._delegate._dlq_purge(dlq_stream)


# Make the broker available as `provider = "delegating_inline"`
registry.register("delegating_inline", f"{__name__}.DelegatingBroker")

Once registered, the broker is configured by its name like any other:

domain = Domain(name="CustomBroker")
domain.config["brokers"]["default"] = {"provider": "delegating_inline"}
domain.init(traverse=False)

with domain.domain_context():
    broker = domain.brokers["default"]
    broker.publish("orders", {"order_id": "1"})

    identifier, message = broker.get_next("orders", "order-processor")
    acknowledged = broker.ack("orders", identifier, "order-processor")

Example: Kafka Broker

Here's a complete example of creating a Kafka broker as an external package:

Project Structure

protean-kafka-broker/
├── pyproject.toml
├── src/
│   └── protean_kafka/
│       ├── __init__.py
│       └── broker.py
└── tests/

pyproject.toml

[project]
name = "protean-kafka-broker"
version = "0.1.0"
description = "Kafka broker for Protean"
requires-python = ">=3.11"
dependencies = [
    "protean>=0.14",
    "kafka-python>=2.0",
]

[project.entry-points."protean.brokers"]
kafka = "protean_kafka:register"

Registration Function

# fragment
# src/protean_kafka/__init__.py
"""Kafka broker plugin for Protean."""

def register():
    """Register Kafka broker with Protean."""
    try:
        # Only register if kafka is available
        import kafka
        from protean.port.broker import registry

        registry.register(
            "kafka",
            "protean_kafka.broker.KafkaBroker"
        )
    except ImportError:
        # Kafka not available, skip registration
        pass

Broker Implementation

# fragment
# src/protean_kafka/broker.py
"""Kafka broker implementation."""

import json
from typing import TYPE_CHECKING, Dict, List, Tuple
import kafka
from protean.port.broker import BaseBroker, BrokerCapabilities

if TYPE_CHECKING:
    from protean.domain import Domain


class KafkaBroker(BaseBroker):
    """Kafka broker implementation for Protean."""

    __broker__ = "kafka"

    def __init__(self, name: str, domain: "Domain", conn_info: Dict) -> None:
        super().__init__(name, domain, conn_info)

        # Initialize Kafka connection
        self.producer = kafka.KafkaProducer(
            bootstrap_servers=conn_info.get("BOOTSTRAP_SERVERS", ["localhost:9092"]),
            value_serializer=lambda v: json.dumps(v).encode('utf-8')
        )

        self.consumer = kafka.KafkaConsumer(
            bootstrap_servers=conn_info.get("BOOTSTRAP_SERVERS", ["localhost:9092"]),
            value_deserializer=lambda m: json.loads(m.decode('utf-8'))
        )

    @property
    def capabilities(self) -> BrokerCapabilities:
        return BrokerCapabilities.ORDERED_MESSAGING

    def _publish(self, stream: str, message: dict) -> str:
        """Publish message to Kafka topic."""
        future = self.producer.send(stream, message)
        record_metadata = future.get(timeout=10)
        return f"{record_metadata.partition}:{record_metadata.offset}"

    def _read(self, stream: str, consumer_group: str, no_of_messages: int) -> List[Tuple[str, dict]]:
        """Read messages from Kafka topic."""
        # Subscribe to topic with consumer group
        self.consumer.subscribe([stream])
        messages = []

        # Poll for messages
        records = self.consumer.poll(timeout_ms=1000, max_records=no_of_messages)

        for topic_partition, msgs in records.items():
            for msg in msgs:
                msg_id = f"{msg.partition}:{msg.offset}"
                messages.append((msg_id, msg.value))

        return messages

    def _ack(self, stream: str, identifier: str, consumer_group: str) -> bool:
        """Acknowledge message (commit offset in Kafka)."""
        try:
            self.consumer.commit()
            return True
        except Exception:
            return False

    def _nack(self, stream: str, identifier: str, consumer_group: str) -> bool:
        """Negative acknowledgment - seek back for reprocessing."""
        # In Kafka, NACK typically means not committing the offset
        # Message will be redelivered on next poll
        return True

    def _ping(self) -> bool:
        """Test Kafka connectivity."""
        try:
            # Check if we can list topics
            self.consumer.list_topics(timeout=5)
            return True
        except:
            return False

    def _health_stats(self) -> dict:
        """Get Kafka broker health statistics."""
        try:
            metrics = self.producer.metrics()
            return {
                "healthy": True,
                "connection_count": metrics.get('connection-count', 0),
                "request_rate": metrics.get('request-rate', 0)
            }
        except Exception as e:
            return {"healthy": False, "error": str(e)}

    def _ensure_connection(self) -> bool:
        """Ensure Kafka connection is healthy."""
        return self._ping()

This listing shows the core methods only. A complete Kafka broker also implements _get_next, _ensure_group, _info, _data_reset and the five _dlq_* methods from the Architecture list. Without them, Python refuses to create an instance of the class.

Installation & Usage

For External Packages

Users install your broker package:

pip install protean-kafka-broker

Then configure it in their domain:

# domain.toml
[brokers.default]
provider = "kafka"
BOOTSTRAP_SERVERS = ["localhost:9092"]

For Internal Brokers

If adding a broker to Protean itself:

  1. Add the broker implementation in src/protean/adapters/broker/
  2. Create a register() function in your broker module
  3. Add entry point in pyproject.toml:
[project.entry-points."protean.brokers"]
mybroker = "protean.adapters.broker.mybroker:register"

Declaring Capabilities

Choose the appropriate capability tier for your broker:

# fragment
@property
def capabilities(self) -> BrokerCapabilities:
    # Basic pub/sub only
    return BrokerCapabilities.BASIC_PUBSUB

    # With consumer groups
    return BrokerCapabilities.SIMPLE_QUEUING

    # With acknowledgments
    return BrokerCapabilities.RELIABLE_MESSAGING

    # With ordering guarantees
    return BrokerCapabilities.ORDERED_MESSAGING

    # Full enterprise features
    return BrokerCapabilities.ENTERPRISE_STREAMING

Blocking reads

A broker that declares BLOCKING_READ implements _read_blocking. The public read_blocking method calls it, and falls back to _read on brokers without the capability.

read_blocking_streams reads several streams in one call and returns a dict with one key per stream, in the order given. Stream subscriptions with priority lanes use it to read the primary and the backfill stream together. The default _read_blocking_streams on BaseBroker builds on _read_blocking. It reads every stream but the last without waiting, and waits on the last stream only if the others are empty. While it waits, it does not see the earlier streams. Override _read_blocking_streams if your backend can wait on several streams in one call, as the Redis broker does with one blocking XREADGROUP:

# fragment
def _read_blocking_streams(
    self,
    streams: Sequence[str],
    consumer_group: str,
    consumer_name: str,
    timeout_ms: int = 5000,
    count: int = 1,  # Applies to each stream
) -> dict[str, list[tuple[str, dict[str, Any]]]]:
    ...

Testing Your Broker

Unit Tests

# fragment
import pytest
from unittest.mock import Mock, patch

def test_broker_capabilities():
    """Test broker declares correct capabilities."""
    broker = KafkaBroker("test", Mock(), {"BOOTSTRAP_SERVERS": ["localhost:9092"]})
    assert broker.has_capability(BrokerCapabilities.MESSAGE_ORDERING)

def test_publish():
    """Test message publishing."""
    with patch("kafka.KafkaProducer"):
        broker = KafkaBroker("test", Mock(), {})
        result = broker.publish("test-stream", {"data": "test"})
        assert result is not None

Integration Tests

# fragment
@pytest.mark.integration
def test_end_to_end():
    """Test full message flow."""
    domain = Domain(__name__)
    domain.config['brokers'] = {
        'default': {
            'provider': 'kafka',
            'BOOTSTRAP_SERVERS': ['localhost:9092']
        }
    }
    domain.init()

    # Publish message
    msg_id = domain.brokers.publish("test", {"data": "test"})

    # Read message
    messages = domain.brokers['default'].read("test", "group", 1)
    assert len(messages) == 1

Guidance for adapter authors

  1. Handle Connection Failures: Implement reconnection logic in _ensure_connection()
  2. Declare Accurate Capabilities: Only declare capabilities you actually implement
  3. Follow Protean Conventions: Use consistent naming and error handling
  4. Provide Health Checks: Implement meaningful _ping() and _health_stats()
  5. Document Configuration: Clearly document all configuration options

Common Patterns

Connection Pooling

# fragment
def __init__(self, name: str, domain: "Domain", conn_info: Dict) -> None:
    super().__init__(name, domain, conn_info)

    # Create connection pool
    self.pool = []
    pool_size = conn_info.get("POOL_SIZE", 10)

    for _ in range(pool_size):
        conn = self._create_connection()
        self.pool.append(conn)

Retry Logic

# fragment
import time

def _publish(self, stream: str, message: dict) -> str:
    max_retries = 3

    for attempt in range(max_retries):
        try:
            return self._do_publish(stream, message)
        except Exception as e:
            if attempt == max_retries - 1:
                raise
            time.sleep(2 ** attempt)  # Exponential backoff