Advanced Messaging Patterns¶
Overview¶
This tutorial covers advanced messaging patterns available in the FCC event bus system, including priority queues, dead-letter queues, circuit breakers, and event replay strategies.
Difficulty: Advanced Prerequisites: Familiarity with the FCC event bus (see Event Bus tutorial). Time: 45 minutes
Step 1: Basic Event Bus Review¶
from fcc.messaging.bus import EventBus
from fcc.messaging.events import Event, EventType
bus = EventBus()
received = []
def handler(event: Event) -> None:
received.append(event)
bus.subscribe(EventType.SIMULATION_STARTED, handler)
bus.publish(Event(event_type=EventType.SIMULATION_STARTED))
print(f"Events received: {len(received)}")
Step 2: Priority Queue Pattern¶
Implement priority-based event processing:
import heapq
from dataclasses import dataclass, field
@dataclass(order=True)
class PrioritizedEvent:
priority: int
event: Event = field(compare=False)
queue = []
heapq.heappush(queue, PrioritizedEvent(2, Event(event_type=EventType.SIMULATION_STARTED)))
heapq.heappush(queue, PrioritizedEvent(0, Event(event_type=EventType.COMPLIANCE_CHECK_COMPLETED)))
heapq.heappush(queue, PrioritizedEvent(1, Event(event_type=EventType.WORKFLOW_COMPLETED)))
while queue:
item = heapq.heappop(queue)
print(f"Priority {item.priority}: {item.event.event_type.value}")
Step 3: Dead-Letter Queue Pattern¶
Capture failed events for retry:
class SimpleDeadLetterQueue:
def __init__(self):
self._failed = []
def add(self, event, error):
self._failed.append((event, str(error)))
def retry(self, bus):
for event, _ in self._failed:
bus.publish(event)
count = len(self._failed)
self._failed.clear()
return count
Step 4: Circuit Breaker Pattern¶
Prevent cascading failures:
class SimpleCircuitBreaker:
def __init__(self, threshold=3):
self._failures = 0
self._threshold = threshold
self._open = False
def call(self, func, *args):
if self._open:
raise RuntimeError("Circuit open")
try:
return func(*args)
except Exception:
self._failures += 1
if self._failures >= self._threshold:
self._open = True
raise
def reset(self):
self._failures = 0
self._open = False
Step 5: Event Replay¶
Replay serialized events for debugging:
from fcc.messaging.serialization import EventSerializer
serializer = EventSerializer()
event = Event(event_type=EventType.SIMULATION_STARTED, data={"scenario": "test"})
# Serialize
json_str = serializer.serialize(event)
print(f"Serialized: {json_str[:80]}...")
# Deserialize and replay
restored = serializer.deserialize(json_str)
bus.publish(restored)
print(f"Replayed event type: {restored.event_type.value}")
Next Steps¶
- Read Guidebook Chapter 23 for in-depth coverage.
- Explore the
ComplianceSubscriberfor a real-world subscriber example. - Build a custom subscriber with DLQ and circuit breaker integration.