用 DynamoDB 分布式租约构建可自动接管的 Fargate WebSocket 工作集群

2026-09-11 22 预计阅读时间: 1 分钟
来源: aws.amazon.com AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:10 分钟

实时流处理服务经常需要维护数百条持久 WebSocket 连接。把这些连接分散到 Amazon ECS 和 AWS Fargate 任务后,单个任务崩溃、扩缩容或滚动发布都可能让部分数据流暂时无人消费。解决问题的关键不是保存 WebSocket 连接本身,而是把每条逻辑数据流的所有权、租约期限和消费进度放进一个可协调、可恢复的控制面。

DynamoDB 的条件写入适合承担这个控制面:多个工作进程可以竞争同一条租约,但只有一个进程能够成功获得所有权;租约过期后,其他进程自动接管并重新建立连接。

管理逻辑数据流,而不是容器

WebSocket 是进程内资源,不能从一个 Fargate 任务直接迁移到另一个任务。系统真正需要管理的是可重建的逻辑单元,例如交易对、租户、设备组或上游订阅分片。

可以为每个逻辑分片保存一条 DynamoDB 记录:

字段 用途
stream_id 分区键,唯一标识逻辑数据流
owner_id 当前持有租约的 ECS 任务或工作进程
lease_expires_at 租约到期时间,使用 Unix 毫秒
fence_token 每次重新获取租约时递增的隔离令牌
checkpoint 已确认处理的序列号或游标
updated_at 诊断和告警所需的更新时间

工作进程取得租约后建立 WebSocket、从 checkpoint 恢复订阅,并周期性续租。续租失败意味着所有权已经丢失,进程必须停止发布数据并关闭连接,不能继续以“旧主人”的身份运行。

租约通常形成以下状态循环:

  1. 工作进程通过条件写入竞争空闲或已过期的分片。
  2. 获胜者建立 WebSocket,并在后台定期续租。
  3. 每处理一批消息,工作进程更新可恢复的检查点。
  4. 任务崩溃后不再续租,租约自然过期。
  5. 其他任务获取租约,使用检查点重新连接和补读。

DynamoDB TTL 可以用于清理历史记录,但不能用于判断租约是否可接管。TTL 删除是异步的,接管条件必须直接比较 lease_expires_at

用条件写入保证单一所有者

下面是一个可以直接改造的 Python 示例。运行前需要安装 boto3、配置 AWS 凭证,并准备一张以 stream_id 为字符串分区键的 DynamoDB 表。

python -m pip install boto3

aws dynamodb create-table \
  --table-name websocket-leases \
  --attribute-definitions AttributeName=stream_id,AttributeType=S \
  --key-schema AttributeName=stream_id,KeyType=HASH \
  --billing-mode PAY_PER_REQUEST

将下面代码保存为 lease_demo.py。修改区域、表名和 stream_id 后即可测试租约竞争与续租:

import os
import socket
import time

import boto3
from botocore.exceptions import ClientError

TABLE_NAME = os.getenv("LEASE_TABLE", "websocket-leases")
REGION = os.getenv("AWS_REGION", "us-east-1")
OWNER_ID = os.getenv("OWNER_ID", socket.gethostname())
LEASE_SECONDS = 30

table = boto3.resource("dynamodb", region_name=REGION).Table(TABLE_NAME)


def now_ms() -> int:
    return int(time.time() * 1000)


def acquire(stream_id: str) -> int | None:
    now = now_ms()
    try:
        result = table.update_item(
            Key={"stream_id": stream_id},
            UpdateExpression=(
                "SET owner_id = :owner, lease_expires_at = :expires, "
                "fence_token = if_not_exists(fence_token, :zero) + :one, "
                "updated_at = :now"
            ),
            ConditionExpression=(
                "attribute_not_exists(owner_id) OR "
                "lease_expires_at < :now OR owner_id = :owner"
            ),
            ExpressionAttributeValues={
                ":owner": OWNER_ID,
                ":expires": now + LEASE_SECONDS * 1000,
                ":now": now,
                ":zero": 0,
                ":one": 1,
            },
            ReturnValues="ALL_NEW",
        )
        return int(result["Attributes"]["fence_token"])
    except ClientError as exc:
        if exc.response["Error"]["Code"] == "ConditionalCheckFailedException":
            return None
        raise


def renew(stream_id: str, token: int) -> bool:
    now = now_ms()
    try:
        table.update_item(
            Key={"stream_id": stream_id},
            UpdateExpression="SET lease_expires_at = :expires, updated_at = :now",
            ConditionExpression=(
                "owner_id = :owner AND fence_token = :token "
                "AND lease_expires_at >= :now"
            ),
            ExpressionAttributeValues={
                ":owner": OWNER_ID,
                ":token": token,
                ":expires": now + LEASE_SECONDS * 1000,
                ":now": now,
            },
        )
        return True
    except ClientError as exc:
        if exc.response["Error"]["Code"] == "ConditionalCheckFailedException":
            return False
        raise


if __name__ == "__main__":
    stream_id = os.getenv("STREAM_ID", "orders-us-east")
    token = acquire(stream_id)
    if token is None:
        raise SystemExit(f"{stream_id} is owned by another worker")

    print(f"acquired {stream_id}: owner={OWNER_ID}, fence_token={token}")
    while True:
        time.sleep(10)
        if not renew(stream_id, token):
            raise SystemExit("lease lost; stop publishing and close the socket")
        print("lease renewed")

可以打开两个终端,使用不同的 OWNER_ID 运行同一个分片,观察条件写入只允许其中一个工作进程持有租约:

OWNER_ID=worker-a STREAM_ID=orders-us-east python lease_demo.py
OWNER_ID=worker-b STREAM_ID=orders-us-east python lease_demo.py

fence_token 用来处理短暂的双活窗口。例如旧任务发生长时间停顿,租约已经被新任务接管,但旧任务恢复后仍可能尝试发送缓冲区中的消息。输出事件应携带隔离令牌;下游存储或消费者需要拒绝小于当前令牌的写入。只有租约而没有下游隔离,无法彻底阻止过期所有者产生副作用。

把故障接管和滚动发布设计成同一条路径

在 ECS 服务中,新任务启动后先参与租约竞争,再建立自己拥有的数据流连接。旧任务收到 SIGTERM 时应进入排空状态:停止领取新租约、提交检查点、主动释放租约或停止续租,然后关闭 WebSocket。

下面的 ECS 服务部署参数可作为起点。它允许部署期间先启动新任务,同时保持现有容量:

aws ecs update-service \
  --cluster realtime-streaming \
  --service websocket-workers \
  --force-new-deployment \
  --deployment-configuration minimumHealthyPercent=100,maximumPercent=200

任务定义中还可以配置停止等待时间,让应用有机会完成排空:

{
  "name": "worker",
  "image": "123456789012.dkr.ecr.us-east-1.amazonaws.com/ws-worker:latest",
  "essential": true,
  "stopTimeout": 60,
  "environment": [
    {"name": "LEASE_TABLE", "value": "websocket-leases"},
    {"name": "LEASE_SECONDS", "value": "30"}
  ]
}

租约周期需要结合实际故障恢复目标设置。例如租约为 30 秒、每 10 秒续租,可以容忍一次短暂超时,并将任务崩溃后的理论接管延迟控制在约 30 秒附近。租约过短会放大 DynamoDB 瞬时错误和调度停顿,过长则增加无人消费的时间。生产实现应对限流和临时网络错误做带抖动的指数退避,但不能在本地假定续租成功。

分片发现也需要控制成本。少量分片可以扫描表;规模扩大后,可以维护待分配分片索引、按桶查询候选项,或由调度器分发竞争目标。无论如何,候选发现只负责提示“去尝试”,最终所有权仍必须由条件写入裁决。

数据完整性的边界

分布式租约解决的是所有权和自动接管,不会自动提供端到端恰好一次处理。可靠性还取决于上游协议和下游写入方式:

  • 上游必须支持按序列号恢复、历史补读或至少能够检测缺口,否则 WebSocket 断开期间的数据无法凭空找回。
  • 检查点应在下游写入成功后推进,避免确认尚未持久化的数据。
  • 接管可能重复读取最后一批消息,下游操作需要按事件 ID 幂等。
  • 机器时钟会参与租约到期判断,应保持时钟同步,并为时钟偏差和网络延迟预留余量。
  • DynamoDB 权限应限制到目标表和必要操作,例如 GetItemUpdateItemQuery
  • 应监控租约获取失败率、续租延迟、过期租约数量、WebSocket 重连次数、检查点滞后和条件写入冲突。

上线前,至少验证三类场景:强制停止一个 Fargate 任务、让工作进程暂停超过租约周期、在持续流量下执行滚动发布。验收标准不只是连接最终恢复,还要检查消息缺口、重复率、接管时间,以及旧所有者是否在丢失租约后立即停止产生副作用。


相关推荐