Skip to content

Brokers

Brokers enable asynchronous message passing between different parts of your system and external services. They decouple message producers from consumers, allowing for scalable, resilient architectures.

Overview

The Broker port in Protean provides a unified interface for different message broker implementations. Each broker adapter implements this interface while providing access to the unique features of the underlying technology.

Note

Protean internally uses an Event Store for domain events and commands within a bounded context. Brokers are primarily used for integration between different systems and for publishing messages to external consumers.

Available Brokers

Protean includes several broker adapters:

Inline Broker

The inline broker processes messages synchronously within the same process. It's ideal for development, testing, and simple applications that don't require distributed messaging.

  • Use cases: Development, testing, small applications
  • Capabilities: Basic pub/sub, simple queuing, reliable messaging
  • No external dependencies required

Redis Stream Broker

The redis broker uses Redis Streams for durable message streaming with consumer groups support.

  • Use cases: Production environments requiring reliable message delivery
  • Capabilities: Consumer groups, message acknowledgment, ordered delivery
  • Requires: Redis 5.0+

Redis PubSub Broker

The redis_pubsub broker uses Redis Lists for simple queuing with basic consumer group support.

  • Use cases: Simple message distribution, development environments
  • Capabilities: Simple queuing with position tracking
  • Requires: Redis 2.0+

Configuration

Brokers are configured in your domain configuration file (domain.toml or .domain.toml):

# Default broker configuration (required)
[brokers.default]
provider = "inline"

# Additional named brokers
[brokers.notifications]
provider = "redis_pubsub"
URI = "redis://localhost:6379/0"

[brokers.analytics]
provider = "redis"
URI = "redis://localhost:6379/1"

Each broker configuration must specify:

  • provider: The broker adapter to use (inline, redis, redis_pubsub, or custom)
  • Additional provider-specific options (like URI for Redis brokers)

Important

You must define a default broker in your configuration. This broker will be used unless a specific broker is requested.

Broker Capabilities

Brokers in Protean declare their capabilities through a capability-based system. This allows you to understand what features each broker supports and write code that adapts to available capabilities.

Capability Tiers

Brokers are organized into capability tiers, each building upon the previous:

  1. BASIC_PUBSUB: Fire-and-forget message publishing
  2. SIMPLE_QUEUING: Basic pub/sub + consumer groups
  3. RELIABLE_MESSAGING: Simple queuing + acknowledgment/rejection
  4. ORDERED_MESSAGING: Reliable messaging + message ordering
  5. ENTERPRISE_STREAMING: Full features including DLQ, replay, partitioning

Checking Capabilities

You can check broker capabilities at runtime:

from protean import Domain
from protean.port.broker import BrokerCapabilities

domain = Domain(name="Capabilities")
domain.init(traverse=False)

with domain.domain_context():
    # Get a broker instance
    broker = domain.brokers["default"]
    broker.publish("orders", {"order_id": "A1"})

    # Check for one capability
    if broker.has_capability(BrokerCapabilities.CONSUMER_GROUPS):
        # Read up to 10 messages as the "order-processor" group
        messages = broker.read(
            stream="orders",
            consumer_group="order-processor",
            no_of_messages=10,
        )

    # Check for all of several capabilities
    if broker.has_all_capabilities(
        BrokerCapabilities.CONSUMER_GROUPS | BrokerCapabilities.ACK_NACK
    ):
        # Acknowledge each message once it is handled
        for identifier, _message in messages:
            broker.ack("orders", identifier, "order-processor")

Basic Usage

Consuming Messages

Messages are typically consumed through Subscribers. A subscriber is a class with a __call__ method that receives the message as a dict. Register subscribers before you initialize the domain.

By default, a published message waits for the message processing engine (see below). With message_processing = "sync", the inline broker delivers each message to its subscribers as soon as it is published. The examples on this page use that setting:

from protean import Domain
from protean.exceptions import ValidationError

domain = Domain(name="Messaging")
domain.config["brokers"]["notifications"] = {"provider": "inline"}

# Deliver to subscribers as soon as a message is published
domain.config["message_processing"] = "sync"

welcomed = []


@domain.subscriber(stream="user-events")
class UserEventSubscriber:
    def __call__(self, payload: dict) -> None:
        # A subscriber receives the raw message dict
        if payload["event_type"] == "user.registered":
            welcomed.append(payload["email"])
            # Send welcome email...


domain.init(traverse=False)

Publishing Messages

with domain.domain_context():
    # Publish to the default broker
    domain.brokers.publish(
        stream="user-events",
        message={
            "event_type": "user.registered",
            "user_id": "123",
            "email": "user@example.com",
        },
    )

    # Publish to a specific broker
    domain.brokers["notifications"].publish(
        stream="notifications",
        message={
            "type": "email",
            "to": "user@example.com",
            "subject": "Welcome!",
        },
    )

After the first publish, the subscriber above has recorded user@example.com.

Message Processing Engine

Protean includes a built-in message processing engine that handles message consumption from brokers:

# Start the message processing engine
protean server

# With specific domain
protean --domain path.to.domain server

The engine automatically:

  • Discovers all registered subscribers
  • Manages consumer groups
  • Handles message acknowledgment
  • Implements retry logic based on broker capabilities
  • Provides graceful shutdown

Error Handling

publish rejects an empty message with a ValidationError before it reaches the broker:

with domain.domain_context():
    try:
        domain.brokers.publish(stream="events", message={})
    except ValidationError as exc:
        # An empty message is rejected before it reaches the broker
        error = exc.messages

Here error is {"message": ["Message cannot be empty"]}.

When the broker's client library raises a connection error, Protean tries to reconnect. If it reconnects, it retries the operation once. If it cannot reconnect, or the retry fails, the client library's error is raised (for Redis, a redis.exceptions.ConnectionError).

This covers only the errors that reach Protean. An adapter can handle errors itself: the Redis broker's get_next, ack and nack log any error and return None or False. Inside a Unit of Work, publish does not contact the broker at all. It records the message, which is sent when the Unit of Work commits.

Health Checks

Monitor broker health and connectivity. health_stats() returns a dict with status ("healthy", "degraded" or "unhealthy"), connected, last_ping_ms, uptime_seconds and a broker-specific details dict:

with domain.domain_context():
    # Check broker connection
    broker = domain.brokers["default"]
    if broker.ping():
        print("Broker is healthy")

    # Get detailed health statistics
    health_stats = broker.health_stats()
    print(f"Status: {health_stats['status']}")  # healthy, degraded or unhealthy
    print(f"Connected: {health_stats['connected']}")
    print(f"Details: {health_stats['details']}")  # broker-specific

Configuring a broker

  1. A default broker is required: Even if it's just the inline broker for development

  2. Check capabilities before using features: Not all brokers support all features

  3. Handle broker failures gracefully: Implement retry logic and circuit breakers

  4. Use appropriate brokers for different concerns:

    • Inline for tests
    • Redis PubSub for notifications
    • Redis Streams for reliable event processing
  5. Monitor broker health: Set up alerts for connection failures and high queue depths

  6. Consider message size limits: Different brokers have different message size constraints

Broker Registry

Brokers register themselves through Python entry points under the protean.brokers group. This means third-party broker packages can be pip-installed and automatically discovered by Protean.

Protean's built-in brokers are registered in pyproject.toml:

[project.entry-points."protean.brokers"]
inline = "protean.adapters.broker.inline:register"
redis = "protean.adapters.broker.redis:register"
redis_pubsub = "protean.adapters.broker.redis_pubsub:register"

Each entry point maps a broker name to a register() function. The function is called on first access and registers the broker class path with the BrokerRegistry. Dependencies are wrapped in try/except so that brokers with uninstalled optional dependencies are silently skipped.

External packages can register their own brokers the same way. See Custom Brokers for a complete guide including a Kafka example.