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:
- Add the broker implementation in
src/protean/adapters/broker/ - Create a
register()function in your broker module - 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
- Handle Connection Failures: Implement reconnection logic in
_ensure_connection() - Declare Accurate Capabilities: Only declare capabilities you actually implement
- Follow Protean Conventions: Use consistent naming and error handling
- Provide Health Checks: Implement meaningful
_ping()and_health_stats() - 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
Related pages
- Review existing broker implementations for examples
- Understand broker capabilities in detail
- Test with the broker test suite
- Share your broker with the Protean community