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¶
-
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"onupdate_item. If two workers race, DynamoDB evaluates the condition atomically and exactly one succeeds; the loser getsConditionalCheckFailedExceptionand backs off. This is the same primitive the wiki already documents for object-store CAS, applied here to a mutable item (patterns/conditional-write). -
The lease is a time-bounded ownership claim renewed by heartbeat. A worker writes its
lease_owner(worker ID) and a futurelease_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 bylease_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. -
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, andlast_seq(for resumption). The post explicitly contrasts this with the Javaamazon-dynamodb-lock-client, which does not integrate domain-specific connection state; combining lock + metadata in one item reduces read operations. -
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). -
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 setslease_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. -
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).
-
Clocks must be synchronized because DynamoDB doesn't supply the
now. Lease expiry is evaluated against the:nowvalue 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). -
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; afinallyblock always releases the lease and removes local tracking. -
Graceful shutdown fans out cleanup in parallel. A shared
asyncio.Eventset by the SIGTERM handler is checked by everywhile not shutdown_event.is_set()loop, so one signal propagates to all connections. WebSocket closes and lease releases run concurrently viaasyncio.gather, keeping total shutdown time roughly constant regardless of connection count (the fast-path release in concepts/distributed-lease). -
Autoscaling is driven by a custom per-worker CloudWatch metric. Each worker publishes
ActiveConnections(namespaceWsFleet) 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¶
- Systems: DynamoDB (lease store + GSI + messages table), ECS on Fargate (worker fleet), SQS (work notifications), API Gateway (START/STOP ingress), Lambda (event router), CloudWatch (custom metric + autoscaling signal).
- Concepts: distributed lease (the central primitive), leader election (per-connection single-owner is a degenerate leader election), GSI / secondary index, at-least-once delivery (SQS), idempotent operations (lease-before-connect dedup), stateless-vs-stateful compute (WebSocket statefulness is the whole problem), clock skew and graceful shutdown (both covered in concepts/distributed-lease).
- Patterns: conditional write (compare-and-set) — the enforcement mechanism for every lease transition.
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()GetItemloop is fine for small fleets; at scale, replace with a single centralizedBatchGetItemcovering 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¶
- Original: https://aws.amazon.com/blogs/architecture/building-resilient-real-time-streaming-workers-with-amazon-dynamodb-leases/
- Raw markdown:
raw/aws/2026-09-10-building-resilient-real-time-streaming-workers-with-amazon-d-e8b2c01b.md - Reference implementation: aws-samples/sample-web-socket-fleet-management-system
Related¶
- concepts/distributed-lease — the central primitive; this article is a canonical instance
- patterns/conditional-write — the compare-and-set enforcement mechanism
- concepts/leader-election — per-connection single-owner as a fine-grained leader election
- concepts/secondary-index — the GSI that makes orphan reconciliation a query, not a scan
- systems/dynamodb · systems/amazon-ecs · systems/aws-fargate · systems/aws-sqs · systems/aws-cloudwatch