Skip to content
Packages Examples Agents Blog Get started

Message bus and queue with Named DI multi-backend support for the Oridecon Framework.


oridecon-queue provides async message queue and bus functionality with Redis, RabbitMQ, Kafka, SQS, and in-memory backends. It includes MessageConsumer workers, a dead-letter queue utility, transactional outbox for atomic DB+message publishing, and a composable message pipeline — all wired through the DI container.

Full documentation: oridecon.dev

Terminal window
uv add oridecon oridecon-queue
# With Redis support
uv add "oridecon-queue[redis]"
# With RabbitMQ support
uv add "oridecon-queue[rabbitmq]"
# With Kafka support
uv add "oridecon-queue[kafka]"
# With AWS SQS support
uv add "oridecon-queue[sqs]"
# With Azure Service Bus support
uv add "oridecon-queue[azure]"
# With GCP Pub/Sub support
uv add "oridecon-queue[gcp]"
from oridecon import Application
from oridecon.queue import BusMessage, MessageConsumer, QueueModule
from oridecon.queue.config import KafkaDriverConfig, NamedQueueConfig, QueueConfig
from oridecon.contracts.queue.protocols import QueueProtocol
class OrderConsumer(MessageConsumer):
topic = "orders"
async def handle(self, message: BusMessage) -> None:
print(f"Processing order: {message.payload}")
async def main() -> None:
async with Application.boot(
modules=[
QueueModule.configure(
QueueConfig(
backends=[
NamedQueueConfig(
name="primary",
primary=True,
driver="kafka",
kafka=KafkaDriverConfig(
bootstrap_servers="localhost:9092",
),
)
]
)
)
]
) as app:
queue = await app.container.resolve(QueueProtocol)
consumer = OrderConsumer(queue)
await consumer.start()
await queue.publish(
"orders", BusMessage(payload={"order_id": "12345", "total": 99.99})
)
if __name__ == "__main__":
import asyncio
asyncio.run(main())

Note: consumers are constructed with the resolved queue and started explicitly via consumer.start() (which subscribes to the topic).

Note: QueueModule.configure() with empty/absent backends registers no queue backend — always declare at least one backend. For tests, QueueModule.stub() uses an in-memory backend.

application.yaml
queue:
backends:
- name: default
primary: true
driver: kafka
max_retries: 3
kafka:
bootstrap_servers: "localhost:9092"
group_id: "oridecon-consumers"
Section titled “Option 2 — Profiles + Environment Variables (recommended)”

Note: backends is a list and cannot be set via environment variables — configure backends in YAML or Python instead.

from oridecon.queue import QueueModule
from oridecon.queue.config import QueueConfig, NamedQueueConfig, KafkaDriverConfig
QueueModule.configure(
QueueConfig(
backends=[
NamedQueueConfig(
name="default",
primary=True,
driver="kafka",
kafka=KafkaDriverConfig(
bootstrap_servers="localhost:9092",
),
),
]
)
)
FieldDefaultEnv varDescription
backends[]ORI_QUEUE__BACKENDSList of named queue backend configurations
backends[n].name(required)ORI_QUEUE__BACKENDS__N__NAMEUnique identifier used for Named() injection
backends[n].driver"memory"ORI_QUEUE__BACKENDS__N__DRIVERDriver: memory, redis, rabbitmq, kafka, sqs, azure_servicebus, gcp_pubsub
backends[n].primaryfalseORI_QUEUE__BACKENDS__N__PRIMARYAlso register as unnamed QueueProtocol binding
backends[n].max_retries3ORI_QUEUE__BACKENDS__N__MAX_RETRIESRetry budget stamped on published messages (BusMessage.max_retries)
backends[n].redis.urlnullORI_QUEUE__BACKENDS__N__REDIS__URLRedis connection URL
backends[n].kafka.bootstrap_serversnullORI_QUEUE__BACKENDS__N__KAFKA__BOOTSTRAP_SERVERSKafka broker addresses (comma-separated)
backends[n].kafka.group_id"oridecon-consumers"ORI_QUEUE__BACKENDS__N__KAFKA__GROUP_IDKafka consumer group ID
backends[n].rabbitmq.urlnullORI_QUEUE__BACKENDS__N__RABBITMQ__URLRabbitMQ connection URL
backends[n].sqs.queue_urlnullORI_QUEUE__BACKENDS__N__SQS__QUEUE_URLSQS queue URL

The in-memory backend runs one asyncio task per subscribed handler per published message, with no bound on how many handlers run concurrently by default. A producer that publishes faster than its handlers can process will pile up unbounded resource usage (DB connections, file handles, …) in a single-process deployment.

Set max_concurrency on InMemoryQueue for any non-trivial single-process deployment:

from oridecon.queue.backends.memory import InMemoryQueue
queue = InMemoryQueue(max_concurrency=16)
await queue.connect()

With a bound set, handler tasks queue behind an internal semaphore once the cap is reached — publish() still returns immediately, but at most max_concurrency handlers execute at any instant. Leave the default (None, unbounded) only when you are certain handler throughput will keep up with publish throughput; the parameter is not yet plumbed through QueueConfig backends, so set it where the backend is constructed.

MethodDescription
QueueModule.configure(config=None)Register queue backends; exports QueueProtocol
QueueModule.scope(*consumers)Exists for feature scoping; consumers are still constructed manually (OrderConsumer(queue)) and started via start()
QueueModule.stub(config=None)In-memory backend for testing
  • Multi-backend messaging — Redis Pub/Sub, RabbitMQ, Kafka, AWS SQS, Azure Service Bus, GCP Pub/Sub, and in-memory
  • Message consumersMessageConsumer subclasses with per-topic handle(), started via consumer.start()
  • Dead-letter queue utilityDeadLetterQueue collects failed messages for inspection and replay
  • In-process publish batchingBatchedPublisher stages messages for an atomic in-process flush() (in-memory only; pair with the durable SQL outbox — OutboxStoreProtocol/SQLOutboxStore/OutboxPublisher — for crash-safe delivery)
  • Message pipelineMessagePipeline with pluggable MiddlewareBase middleware
  • Named DI multi-backendAnnotated[QueueProtocol, Named("events")] for multiple backends
  • Retry metadataBusMessage carries retry_count / max_retries with should_retry() / is_expired()
  • Consumer groups — Kafka consumer groups for load balancing
from oridecon import Application
from oridecon.queue import BusMessage, QueueModule
from oridecon.contracts.queue.protocols import QueueProtocol
async def test_message_consumer():
async with Application.boot(modules=[QueueModule.stub()]) as app:
queue = await app.container.resolve(QueueProtocol)
await queue.publish("test-topic", BusMessage(payload={"key": "value"}))
# Test with in-memory backend
FileWhat it contains
src/oridecon/queue/module.pyQueueModule.configure(), .scope(), .stub()
src/oridecon/queue/config.pyQueueConfig, NamedQueueConfig, backend configs
src/oridecon/queue/di/provider.pyQueueProvider boot and registration
src/oridecon/queue/consumers/consumer.pyMessageConsumer base class
src/oridecon/queue/core/dlq.pyDeadLetterQueue implementation
src/oridecon/queue/core/batch_publisher.pyBatchedPublisher in-memory publish batching
src/oridecon/queue/core/pipeline.pyMessagePipeline and MiddlewareBase
src/oridecon/queue/backends/kafka.pyKafka backend implementation
src/oridecon/queue/backends/rabbitmq.pyRabbitMQ backend implementation
src/oridecon/queue/backends/redis.pyRedis backend implementation
BackendDurabilityOrderingThroughputUse Case
MemoryNoneFIFOVery HighDevelopment, testing
Redis Pub/SubAt-most-onceNo guaranteeVery HighReal-time events, ephemeral messages
RabbitMQAt-least-oncePer-queueHighTask queues, work distribution
KafkaAt-least-oncePer-partitionVery HighEvent streams, audit logs
SQSAt-least-onceBest-effort (FIFO available)HighAWS-native, decoupled systems