Message Queuing (RabbitMQ)
Okustera Message Queuing provides enterprise-ready, distributed message brokering and asynchronous task coordination powered by the RabbitMQ Cluster Operator.
Engineered for mission-critical microservices, event-driven architectures, and distributed pipelines, it delivers high-throughput AMQP queuing, Raft-based Quorum Queues, native web management consoles, and automated TLS certificate lifecycle management.
Architectural Highlights
- Raft Quorum Queues: Replicated FIFO queues providing consensus-based data safety and automatic recovery from network partitions.
- Protocol Versatility: Native support for AMQP 0-9-1, AMQP 1.0, MQTT 3.1.1, STOMP, and WebSockets.
- Integrated Observability: Built-in Prometheus metrics exporter exposing queue depth, publish/delivery rates, and consumer acknowledgments to Grafana.
- Kubernetes Native Autoscaling: Dynamically scale worker pods based on queue depth metrics using KEDA (Kubernetes Event-driven Autoscaling).
Example Cluster Manifest
apiVersion: rabbitmq.com/v1beta1
kind: RabbitmqCluster
metadata:
name: production-broker
namespace: default
spec:
replicas: 3
image: rabbitmq:3.13-management
persistence:
storageClassName: ceph-rbd
storage: 20Gi
resources:
requests:
cpu: 500m
memory: 2Gi
limits:
cpu: 2000m
memory: 4Gi
Python (pika) Publishing Example
import pika
credentials = pika.PlainCredentials('app_user', 'app_password')
parameters = pika.ConnectionParameters(
'production-broker.default.svc.cluster.local',
5672,
'/',
credentials
)
connection = pika.BlockingConnection(parameters)
channel = connection.channel()
# Declare durable quorum queue
channel.queue_declare(
queue='orders-queue',
durable=True,
arguments={'x-queue-type': 'quorum'}
)
channel.basic_publish(
exchange='',
routing_key='orders-queue',
body='{"order_id": 9821, "status": "pending"}',
properties=pika.BasicProperties(delivery_mode=2)
)
print("Order published successfully")
connection.close()