Skip to content

AWS 2026-09-10 Tier 1

Read original ↗

Building resilient real-time streaming workers with Amazon DynamoDB leases

Summary

An AWS Architecture Blog post that walks through a production pattern for coordinating a fleet of stateful WebSocket workers — each holding hundreds of long-lived outbound connections to upstream streaming sources — without a dedicated coordination service (no ZooKeeper, etcd, or Consul). The core idea is a time-bounded lease stored as a single DynamoDB item per connection; workers claim, renew, and release leases using DynamoDB conditional writes (patterns/conditional-write), which act as an atomic compare-and-set. A GSI on (desired_state, lease_expires_at_ms) lets any healthy worker query for orphaned connections and reclaim them. Workers run on ECS on Fargate, are notified of new work through SQS, and autoscale on a custom CloudWatch ActiveConnections metric. The lease mechanism handles the full failure spectrum: unexpected crashes (lease expires, reconciliation reclaims), rolling deployments and scale-in (graceful shutdown releases leases immediately), and double-claiming (conditional write guarantees exactly one owner).

The motivating problem: a real-time transcription service processing 500 concurrent meetings where a single worker failure dropped 100+ WebSocket connections and caused 2–3 minutes of data loss per connection until operators manually restarted services.

Key takeaways

  1. DynamoDB conditional writes are a distributed lock without a lock service. Lease acquisition uses ConditionExpression: "attribute_not_exists(lease_expires_at_ms) OR lease_expires_at_ms < :now" on update_item. If two workers race, DynamoDB evaluates the condition atomically and exactly one succeeds; the loser gets ConditionalCheckFailedException and backs off. This is the same primitive the wiki already documents for object-store CAS, applied here to a mutable item (patterns/conditional-write).

  2. The lease is a time-bounded ownership claim renewed by heartbeat. A worker writes its lease_owner (worker ID) and a future lease_expires_at_ms; it must keep renewing (default every 5 s) before the lease (default 20 s) expires. Renewal is itself a conditional write guarded by lease_owner = :w — if it returns false, the worker has lost ownership mid-flight and exits cleanly. The 4:1 lease-to-heartbeat ratio gives ~4 renewal attempts before expiry.

  3. Colocating domain state with the lock record saves reads. The single DynamoDB item carries desired_state (STARTED/STOPPED), ws_url, lease_owner, lease_expires_at_ms, and last_seq (for resumption). The post explicitly contrasts this with the Java amazon-dynamodb-lock-client, which does not integrate domain-specific connection state; combining lock + metadata in one item reduces read operations.

  4. A GSI turns fleet-wide reconciliation into an efficient query. The GSI keyed on (desired_state [PK], lease_expires_at_ms [SK]) lets a reconciliation loop query directly for orphans: desired_state = STARTED AND lease_expires_at_ms < now. Without the secondary index this would require a scan (concepts/secondary-index).

  5. The failure spectrum is covered by two complementary mechanisms. Unexpected crash: the crashed worker can't renew, the lease expires after LEASE_SECONDS, and the orphan-reconciliation loop (default every 60 s) reclaims it — worst-case recovery ≈ 80 s (20 s expiry + up to 60 s reconciliation). Graceful shutdown (SIGTERM on rolling deploy / scale-in): the worker sets lease_expires_at_ms = 0 (already-expired) so other workers pick it up on the next reconciliation cycle without waiting for natural expiry — the fast path.

  6. SQS is a fast-notification hint, not the source of ownership. A START event enqueues an SQS message so workers pick up new connections in seconds instead of waiting for the next 60 s reconciliation cycle. But the SQS message is only a hint — a worker must still acquire the lease before opening the WebSocket, so redeliveries and duplicate notifications can't cause double-claiming (concepts/at-least-once-delivery, concepts/idempotent-operations).

  7. Clocks must be synchronized because DynamoDB doesn't supply the now. Lease expiry is evaluated against the :now value the calling worker supplies from its local clock, not a DynamoDB server clock. Fargate tasks in one Region get Amazon Time Sync (skew within a few ms), well inside the 20 s lease margin. Off Fargate, verify NTP and increase lease duration by max expected clock skew to avoid false expirations (see the clock-skew discussion in concepts/distributed-lease).

  8. Connection management runs three concurrent async tasks joined by asyncio.gather. Per connection: a heartbeat loop (renews lease / detects lost ownership), a receive loop (writes each upstream message to a separate DynamoDB messages table), and a desired-state checker (polls for STOPPED every 10 s). Any one returning ends the connection; a finally block always releases the lease and removes local tracking.

  9. Graceful shutdown fans out cleanup in parallel. A shared asyncio.Event set by the SIGTERM handler is checked by every while not shutdown_event.is_set() loop, so one signal propagates to all connections. WebSocket closes and lease releases run concurrently via asyncio.gather, keeping total shutdown time roughly constant regardless of connection count (the fast-path release in concepts/distributed-lease).

  10. Autoscaling is driven by a custom per-worker CloudWatch metric. Each worker publishes ActiveConnections (namespace WsFleet) every 30 s; an Application Auto Scaling target-tracking policy scales on the fleet average (e.g. 700 connections/task). Scale-in is safe because of the lease pattern: SIGTERM → release leases → other workers reclaim via reconciliation.

Operational numbers

  • Motivating scenario: 500 concurrent meetings, 100+ connections dropped per worker failure, 2–3 min data loss per connection before manual restart.
  • Lease duration: 20 s (survives brief network hiccups; fast failover).
  • Heartbeat interval: 5 s (4× safety margin vs lease).
  • Reconciliation interval: 60 s (recovery speed vs DynamoDB read cost).
  • Worst-case unexpected-crash recovery: ~80 s (20 s lease expiry + up to 60 s reconciliation).
  • Max connections per task: 700 (memory/CPU profiled; each WebSocket ≈ 2–5 MB).
  • Scale-out cooldown: 2 min; scale-in cooldown: 15 min.
  • Desired-state check: every 10 s (per connection); metric publish: every 30 s.
  • DynamoDB cost (heartbeat writes dominate; 1 WCU per renewal): 100 conns ≈ 20 WCU/s ≈ $65/mo on-demand / $10/mo provisioned; 500 conns ≈ $325 / $47; 2,000 conns ≈ $1,300 / $190. GSI reconciliation queries default to eventually-consistent reads (half the RCU of strong reads).

Extracted vocabulary

Caveats

  • The code samples use print(); the post says to replace with structured logging and emit CloudWatch metrics for lease-acquisition failures and reconnection events.
  • The per-connection check_desired_state() GetItem loop is fine for small fleets; at scale, replace with a single centralized BatchGetItem covering all active connections (N reads/10 s → 1 batched read).
  • Applies specifically to systems holding long-lived connections at scale (transcription, IoT ingestion, financial feeds, live event streaming). It is not a general task-queue pattern — SQS alone cannot track continuous ownership.
  • Suggested enhancements: AWS X-Ray distributed tracing across workers, and reconnection with upstream replay / offset-based resumption (via last_seq) to close the data gap between failure and recovery.

Source

Last updated · 766 distilled / 2,225 read