Building State Machines That Survive Missing Events
Handle out-of-order and lost events in distributed systems without strict ordering. Practical patterns for gap detection, recovery, and avoiding zombie states.
The Problem Nobody Talks About Until It Costs Them
You're running a subscription billing system. A customer's payment processes successfully, triggering an OrderPaid event. A few milliseconds later, the system should transition to OrderFulfilling. But the event broker hiccups. The fulfillment service never sees OrderPaid. Three hours later, when the queue recovers, that event arrives—but by then, the order has already timed out and moved to a Cancelled state.
Now you have a zombie order: paid but never fulfilled. Your reconciliation job catches it at 2 AM. Your on-call engineer manually updates the database. A customer complains on Twitter.
This scenario plays out differently across industries—payment processing, order fulfillment, subscription management, IoT sensor networks—but the root cause is always the same: event-driven systems assume events arrive in order and never disappear. They usually don't. And when they don't, your state machine breaks.
Most teams discover this the hard way. The fix isn't to demand perfect ordering (that's expensive and often impossible). It's to design state machines that acknowledge gaps, tolerate reordering, and recover gracefully.
Designing State Machines for an Unreliable World
1. Explicit State Versioning and Gap Detection
Instead of assuming the next event is always the one you expect, build state machines that track what events they've seen and what they're waiting for.
class OrderStateMachine:
def __init__(self, order_id: str):
self.order_id = order_id
self.current_state = "Created"
self.seen_events = set()
self.expected_event_sequence = [
"OrderCreated",
"PaymentProcessed",
"InventoryReserved",
"OrderFulfilling",
"OrderShipped"
]
self.last_seen_sequence_number = 0
def handle_event(self, event: dict) -> bool:
"""
Returns True if transition succeeded, False if event is unexpected
or out of order. Logs gap detection.
"""
event_type = event["type"]
sequence_number = event["sequence_number"]
# Detect gaps
if sequence_number > self.last_seen_sequence_number + 1:
gap_size = sequence_number - self.last_seen_sequence_number - 1
self.log_gap(gap_size, self.last_seen_sequence_number, sequence_number)
# Don't fail—queue this event for later replay
return self.defer_event(event)
if event_type in self.seen_events:
# Idempotent: we've already processed this
return True
# Validate state transition
if not self._is_valid_transition(self.current_state, event_type):
self.log_invalid_transition(self.current_state, event_type)
return False
self.seen_events.add(event_type)
self.last_seen_sequence_number = sequence_number
self.current_state = self._next_state(self.current_state, event_type)
return True
def _is_valid_transition(self, current: str, event: str) -> bool:
"""Explicit mapping of allowed transitions."""
transitions = {
"Created": ["OrderCreated"],
"Pending": ["PaymentProcessed"],
"PendingInventory": ["InventoryReserved"],
"Fulfilling": ["OrderFulfilling"],
"Shipped": ["OrderShipped"]
}
return event in transitions.get(current, [])
def defer_event(self, event: dict):
"""Store out-of-order event for later replay."""
# Write to a deferred event queue, keyed by order_id
# Include timestamp so we can detect truly lost events
pass
def log_gap(self, gap_size: int, last_seen: int, current: int):
"""Alert monitoring system to potential message loss."""
passThis approach gives you three immediate benefits:
- Idempotency by design: Replaying the same event twice won't corrupt state.
- Visible gaps: You know when events are missing, not just when your system breaks.
- Deferred processing: Out-of-order events don't fail—they wait in a queue until the gap closes.
2. Handling Reordering Without Strict Ordering Guarantees
Real systems rarely guarantee strict ordering across partitions. Instead, use logical timestamps and event dependencies.
class Event:
def __init__(self, event_type: str, order_id: str, timestamp: int,
depends_on: list = None):
self.type = event_type
self.order_id = order_id
self.timestamp = timestamp # Logical clock or wall-clock
self.depends_on = depends_on or [] # Event types this requires
self.idempotency_key = f"{order_id}_{event_type}_{timestamp}"
class OrderEventProcessor:
def __init__(self):
self.deferred_queue = {} # order_id -> [events]
self.processed = {} # order_id -> set of idempotency_keys
def process_event(self, event: Event, state_machine: OrderStateMachine):
"""
Process event only if all dependencies have been seen.
Otherwise, defer.
"""
order_id = event.order_id
# Idempotency check
if event.idempotency_key in self.processed.get(order_id, set()):
return
# Check dependencies
unmet_deps = [dep for dep in event.depends_on
if not self._dependency_satisfied(order_id, dep)]
if unmet_deps:
self._defer(event, f"Waiting for: {unmet_deps}")
return
# Safe to process
success = state_machine.handle_event(event.__dict__)
if success:
if order_id not in self.processed:
self.processed[order_id] = set()
self.processed[order_id].add(event.idempotency_key)
# Attempt to flush deferred events
self._flush_deferred(order_id, state_machine)
def _dependency_satisfied(self, order_id: str, event_type: str) -> bool:
"""Check if we've seen this event type for this order."""
if order_id not in self.processed:
return False
matching = [key for key in self.processed[order_id]
if event_type in key]
return len(matching) > 0
def _defer(self, event: Event, reason: str):
order_id = event.order_id
if order_id not in self.deferred_queue:
self.deferred_queue[order_id] = []
self.deferred_queue[order_id].append((event, reason))
def _flush_deferred(self, order_id: str, state_machine: OrderStateMachine):
"""Retry deferred events in dependency order."""
if order_id not in self.deferred_queue:
return
queue = self.deferred_queue[order_id]
# Sort by dependency satisfaction, not arrival time
retry_count = 0
while queue and retry_count < len(queue):
event, reason = queue.pop(0)
try:
self.process_event(event, state_machine)
retry_count = 0 # Reset on success
except Exception:
queue.append((event, reason))
retry_count += 1The key insight: dependencies are explicit. An OrderShipped event depends on OrderFulfilling, not on arrival order. This lets you handle reordering naturally.
Tradeoffs and Real Failure Modes
Strict Ordering vs. Availability
Demanding strict ordering (via a single-partition event log or distributed lock) guarantees correctness but costs throughput and availability. If your ordering broker goes down, so does your entire system.
Practical middle ground: Use logical clocks and dependency graphs instead. Trade some complexity in your state machine for higher availability and lower latency.
Replay Complexity and Zombie States
Here's a real scenario I encountered:
The Incident: A subscription renewal service stored state in Postgres and published events to Kafka. During a database failover, a few hundred renewal transactions were replayed—but the Kafka events had already been published hours earlier. Downstream systems saw duplicate SubscriptionRenewed events for the same subscription in the same billing cycle.
The billing system's state machine looked like this (naive version):
# ❌ BROKEN: No idempotency, no gap detection
if event.type == "SubscriptionRenewed":
subscription.balance -= event.amount
subscription.next_renewal_date = calculate_next_renewal()
subscription.save()Running this twice doubled the charge. We caught it via a reconciliation job 6 hours later.
The fix: We added idempotency keys and a "seen events" log:
# ✅ FIXED: Idempotent with deduplication
event_key = f"{subscription_id}_{event.type}_{event.timestamp}"
if event_key in subscription.processed_events: