AWS Architecture Blog

Building resilient real-time streaming workers with Amazon DynamoDB leases

8.5内容质量
Building resilient real-time streaming workers with Amazon DynamoDB leases

TL;DR · AI 摘要

AWS通过DynamoDB租约机制实现WebSocket连接的自动故障转移,将数据恢复时间从分钟级缩短至秒级。

核心要点

  • 使用DynamoDB条件写入实现租约所有权竞争,确保连接不丢失
  • 孤儿连接自动回收机制减少故障恢复时间至2-3秒
  • 优雅关闭策略使部署停机时间降低90%以上

结构提纲

按章节快速跳转。

  1. 500并发会议转录场景暴露WebSocket连接管理难题

  2. WebSocket状态性导致故障时连接丢失且无法自动恢复

  3. DynamoDB条件写入实现租约竞争与所有权转移

  4. 三重保障机制:租约管理+孤儿回收+优雅关闭

  5. Fargate容器化方案适配多种计算层架构

  6. 恢复时间从分钟级降至秒级且无需人工干预

思维导图

用一张图看清主题之间的关系。

查看大纲文本(无障碍 / 无 JS 友好)
  • WebSocket租约管理
    • 核心挑战
      • 连接状态性
      • 故障恢复
      • 部署中断
    • 解决方案
      • DynamoDB租约
      • 条件写入
      • 孤儿回收
      • 优雅关闭
    • 效果指标
      • 恢复时间<2秒
      • 停机时间降低90%

金句 / Highlights

值得收藏与分享的关键句。

#AWS#DynamoDB#实时流处理#WebSocket
打开原文

使用 Amazon DynamoDB 租约构建高可用实时流处理工作器 | AWS 架构博客

使用 Amazon DynamoDB 租约构建高可用实时流处理工作器

设想一个实时语音转录服务正在处理 500 个并发会议。每个处理这些会议的工作器都需要一个专用的出站 WebSocket 连接到上游流媒体源。当单个工作器发生故障时,会导致 100+ 连接中断,每个连接需要操作人员手动重启服务后才能恢复,数据丢失时间长达 2 至 3 分钟。

构建能够维护数百个持久 WebSocket 连接的实时流处理工作器存在协调挑战:当工作器意外停止时,其连接会变得无人管理,数据流会中断。必须确保每个连接只能由一个工作器独占,但工作器会独立发生故障、重新部署和扩展。如果没有机制来跟踪所有权并自动将健康工作器可声明的连接进行转移,操作人员必须手动介入处理每次故障。

这种模式可减少故障时的人工干预,将连接恢复时间从分钟级缩短至秒级,并在部署期间无需外部协调服务即可最小化停机时间。

在本文中,您将学习如何在 Amazon Elastic Container Service (Amazon ECS) 和 AWS Fargate 上构建 WebSocket 舰队管理系统。Amazon DynamoDB 是该解决方案中用于管理分布式租约所有权、协调和故障转移的主要服务。对于计算层,本文使用 AWS Fargate 上的 Amazon ECS 运行工作器舰队。不过,您可以将此模式适配到任何计算层(例如 Amazon Elastic Kubernetes Service (Amazon EKS) 或带有自动扩展组的 Amazon Elastic Compute Cloud (Amazon EC2)),而无需更改核心租约逻辑。您将学习如何通过条件写入实现基于租约的所有权管理,通过孤儿协调实现自动故障转移,以及通过优雅关闭实现低停机时间部署。

挑战:管理长期存在的 WebSocket 连接

WebSocket 连接本质上与 HTTP 请求截然不同。HTTP 请求到达、被处理并返回响应,服务器在请求之间不保存状态。相比之下,WebSocket 连接是一个持久的双向通道。工作器必须保持 TCP 连接开放,处理上游源发送的消息,并响应上游源的保活 ping。

这种状态性引入了多个运营挑战:

工作器故障。当工作器进程意外停止或其容器终止时,工作器会断开 WebSocket 连接。上游源可能会短暂缓冲数据,但如果没有机制检测到故障并将连接重新分配给健康工作器,系统会丢失数据。

滚动部署。ECS 滚动部署会终止旧任务并启动新任务。每个终止的任务都会断开连接。如果没有协调机制,会出现连接无所有者的窗口期。

水平扩展。添加工作器相对简单。新任务启动后即可接管工作。移除工作器则更复杂。您需要从即将离开的工作器中排出连接,并在任务退出前验证其他工作器已接管。

双重声明。如果两个工作线程都认为自己拥有相同的连接,它们都会尝试连接到同一个上游源。这可能导致数据重复处理、协议错误,或上游服务拒绝连接。

由于工作线程是作为WebSocket客户端主动发起连接到上游源的,因此需要在应用层而非网络层实现协调机制。

解决方案概述

该架构使用六种AWS服务来协调WebSocket工作线程舰队:

图1:WebSocket舰队管理架构

  • Amazon API Gateway:通过REST API接收来自外部系统的START和STOP事件。START事件表示新的流式会话(如会议或直播)已开始,需要专用的WebSocket连接。STOP事件表示流式会话已结束,应释放连接。
  • AWS Lambda(事件路由器):用于将连接状态写入Amazon DynamoDB,并向Amazon Simple Queue Service(Amazon SQS)发送通知。
  • Amazon DynamoDB:用于存储连接状态和租约所有权。条件写入(仅在满足指定条件时成功的原子操作)可以在不依赖外部协调服务的情况下提供分布式锁定功能。
  • Amazon SQS:用于向工作线程分发任务通知,以便快速获取新连接。
  • AWS Fargate上的Amazon ECS:用于运行工作线程舰队。每个工作线程都会轮询Amazon SQS,管理WebSocket连接,并通过心跳信号续租。
  • Amazon CloudWatch:用于收集自定义指标(活跃连接数),驱动ECS自动扩展。

关键洞察是DynamoDB的条件写入可以作为分布式锁,无需单独的协调服务。每个连接都有一个租约:一个有时间限制的所有权声明。工作线程必须持续续租。如果工作线程意外停止,租约将过期,其他工作线程会接管。

为什么不能仅使用SQS或现有锁客户端?

SQS在此架构中作为快速通知通道发挥重要作用,但不能作为唯一的协调机制。SQS的设计目的是任务执行,将工作单元发送给一个消费者。WebSocket连接所有权不是一次性任务,而是需要在整个连接生命周期内持续维护和续租的状态。SQS没有跟踪当前连接所有者的机制,无法查询没有活跃所有者的连接,也无法表示管理连接所需的领域状态(desired_state、ws_url、last_seq)。DynamoDB通过持久化项、条件写入和二级索引提供了所有这些功能。

AWS发布的amazon-dynamodb-lock-client库在DynamoDB上实现了类似的分布式锁定原语。但该库是为Java环境设计的,没有将领域特定的连接状态整合到锁记录中。本方案采用异步Python实现以匹配工作线程架构,将锁所有权和连接元数据合并到单个DynamoDB项中以减少读取操作,并使用GSI启用舰队范围的对账查询,这是通用锁客户端无法提供的功能。

租约模式

租约是 DynamoDB 中的一行数据,用于跟踪连接的所有者以及所有权过期时间。该表使用以下架构:

属性 | 类型 | 描述 --- | --- | --- Pk | 字符串(分区键) | 连接 ID,例如 CONN#meeting-123 desired_state | 字符串 | STARTED 或 STOPPED ws_url | 上游 WebSocket 连接地址 lease_owner | 当前拥有该连接的工作者 ID lease_expires_at_ms | 数字 | 租约过期的 Unix 时间戳(毫秒) last_seq | 最后处理的序列号(用于恢复)

一个全局二级索引(GSI)允许通过 desired_state(分区键)和 lease_expires_at_ms(排序键)查询非主键属性。该索引可高效查询未管理的连接:即 desired_state 为 STARTED 且租约已过期的连接。

关于时钟精度的说明

租约过期机制依赖于工作者进程使用本地系统时钟生成的 Unix 时间戳(毫秒)。DynamoDB 会根据调用工作者提供的 now 值评估租约过期条件,而不是使用 DynamoDB 服务器端时钟。这意味着所有工作者必须保持合理的时钟同步,才能保证租约模式的正确行为。

运行在相同 AWS 区域的 AWS Fargate 任务会通过 Amazon Time Sync Service 实现时钟同步,任务间的时钟偏差可控制在几毫秒内。这远低于默认 20 秒租约时长和 5 秒心跳间隔提供的安全余量。如果在 AWS Fargate 以外的计算基础设施上部署此模式,请确认已配置 NTP 同步并监控时钟漂移。对于无法保证时钟精度的环境,应根据预期的最大时钟偏差增加租约时长,以防止租约误过期。

租约生命周期包含四个状态。图 2 展示了租约状态机。

图 2:租约生命周期

获取(Acquire)

工作者通过将自身工作者 ID(lease_owner)和未来过期时间戳(lease_expires_at_ms)写入 DynamoDB 租约记录来声明连接所有权。条件表达式确保只有一个工作者能成功:它检查租约尚不存在(attribute_not_exists)或现有租约已过期(lease_expires_at_ms < :now)。当两个工作者同时尝试获取同一连接时,DynamoDB 会原子性地评估此条件,仅允许一个工作者成功。另一个工作者将收到 ConditionalCheckFailedException 异常并优雅地回退。

以下代码示例来自工作者应用(worker.py),该应用在启动时会初始化 Amazon DynamoDB 表客户端、工作者 ID 和配置。完整实现可在 GitHub 仓库中找到。

code
async def try_acquire_lease(pk: str) -> Optional[dict]:
    """尝试获取连接的租约。"""
    try:
        resp = table.update_item(
            Key={"pk": pk},
            UpdateExpression=(
                "SET lease_owner = :w, "
                "lease_expires_at_ms = :exp, "
                "updated_at_ms = :now"
            ),
            ConditionExpression=(
                "attribute_not_exists(lease_expires_at_ms) "
                "OR lease_expires_at_ms < :now"
            ),
            ExpressionAttributeValues={
                ":w": WORKER_ID,
                ":exp": now_ms() + LEASE_SECONDS * 1000,
                ":now": now_ms(),
            },
            ReturnValues="ALL_NEW",
        )
        return resp["Attributes"]
    except ClientError as e:
        if e.response["Error"]["Code"] == "ConditionalCheckFailedException":
            return None  # 另一个工作线程已经拥有此连接
        raise

ConditionExpression 是关键部分:当租约不存在时(attribute_not_exists)或已过期时(lease_expires_at_ms < :now),该表达式会成功执行。

续订

拥有租约的工作线程每隔几秒(心跳)会续订租约。条件表达式会验证工作线程是否仍然拥有租约:

code
async def renew_lease(pk: str) -> bool:
    """续订已拥有的连接租约。"""
    try:
        table.update_item(
            Key={"pk": pk},
            UpdateExpression=(
                "SET lease_expires_at_ms = :exp, "
                "updated_at_ms = :now"
            ),
            ConditionExpression="lease_owner = :w",
            ExpressionAttributeValues={
                ":w": WORKER_ID,
                ":exp": now_ms() + LEASE_SECONDS * 1000,
                ":now": now_ms(),
            },
        )
        return True
    except ClientError:
        return False  # 已失去所有权

如果续订返回 False,工作线程会知道已失去所有权(可能有其他工作线程获取了过期的租约),并会优雅退出。

释放

在优雅关闭期间,工作线程会显式释放其租约,以便其他工作线程可以立即获取,而无需等待租约过期:

code
async def release_lease(pk: str):
    """释放连接的租约。"""
    try:
        table.update_item(
            Key={"pk": pk},
            UpdateExpression=(
                "SET lease_owner = :empty, "
                "lease_expires_at_ms = :zero"
            ),
            ConditionExpression="lease_owner = :w",
            ExpressionAttributeValues={
                ":w": WORKER_ID,
                ":empty": "",
                ":zero": 0,
            },
        )
    except ClientError:
        pass  # 已释放或被其他工作线程获取

过期

当工作线程因容器崩溃、网络分区或进程故障等意外情况停止时,将无法续订租约。与优雅关闭不同,工作线程没有机会显式释放所有权。租约会保留在 DynamoDB 中,但 lease_expires_at_ms 时间戳会因未续订而过期。

此过期租约表示一个没有活跃所有者的连接:desired_state 仍为 STARTED(连接应处于活动状态),但没有健康的工作线程在管理它。该连接现在成为一个孤儿。

协调循环通过查询 GSI 中 desired_state = STARTED 且 lease_expires_at_ms < now 的记录来检测此条件。任何健康的工作者发现此类记录后,都可以使用与初始获取时相同的条件写入操作尝试获取该记录。由于 lease_expires_at_ms < :now 是获取操作的有效条件之一,过期的租约将被视为与未声明的租约完全相同。

Expired 状态是瞬时状态:它存在于租约停止续订的时刻与协调循环运行并由新工作者成功获取的时刻之间。连接在 Expired 状态下的最大持续时间由协调间隔(默认:60 秒)所限制。

技术实现

以下部分将逐步介绍系统各个组件,从事件如何进入处理管道开始,到舰队如何扩展结束。

事件摄入

当外部系统需要启动或停止流式连接时,它会通过 API Gateway 向 Lambda 事件路由器发送事件。Lambda 函数将连接状态写入 DynamoDB 并向 SQS 队列发送通知:

code
def handler(event, context):
    payload = json.loads(event.get("body", "{}"))
    event_type = payload["event_type"].upper()
    connection_id = payload["connection_id"]
    pk = f"CONN#{connection_id}"

    if event_type == "START":
        table.put_item(Item={
            "pk": pk,
            "desired_state": "STARTED",
            "ws_url": payload["ws_url"],
            "last_seq": 0,
            "lease_owner": "",
            "lease_expires_at_ms": 0,
            "updated_at_ms": now_ms(),
        })
        sqs.send_message(
            QueueUrl=QUEUE_URL,
            MessageBody=json.dumps({"pk": pk})
        )

    elif event_type == "STOP":
        table.update_item(
            Key={"pk": pk},
            UpdateExpression="SET desired_state = :s, updated_at_ms = :t",
            ExpressionAttributeValues={
                ":s": "STOPPED", ":t": now_ms()
            },
        )

    return {"statusCode": 200, "body": "OK"}

DynamoDB 是连接状态的权威数据源。Amazon SQS 作为快速通知通道。当 START 事件到达时,SQS 消息会立即通知可用工作者可以声明新连接,因此工作者无需等待下一次协调周期(默认:60 秒)即可发现并获取新连接。如果没有 SQS,新连接仅会在协调循环下次运行时查询 GSI 以查找未管理连接时才会被处理。

工作者轮询

每个 ECS Fargate 工作者都会运行一个持续的 SQS 轮询循环以获取新连接通知。在启动新的 WebSocket 连接之前,该循环会执行以下四个步骤:

  1. 容量检查

在接受任何新任务之前,工作者会检查是否已达到最大连接限制(MAX_CONNECTIONS)。如果工作者已满载,它将暂停 5 秒并跳过当前轮询周期。这可以防止单个工作者被压垮,而舰队中的其他工作者却未充分利用。

  1. 去重

如果工作者已经管理了 SQS 消息中引用的连接(在本地连接字典中跟踪),它将删除该消息并继续处理。这可以处理同一连接生成多个 SQS 通知的情况,例如在重试或重新投递期间。

  1. WebSocket启动前的租约获取

SQS消息仅作为所有权的提示,而非保证。在启动WebSocket连接前,工作者必须通过try_acquire_lease成功获取DynamoDB租约。如果其他工作者已声明该连接,try_acquire_lease将返回None,当前工作者将跳过该操作。这确保了任意时间每个连接仅由一个工作者拥有。

  1. 任务创建

若租约已获取且desired_state为STARTED,工作者将创建异步任务来管理WebSocket连接。无论租约是否成功获取,SQS消息都会被删除,防止重复处理相同通知。

以下代码展示了完整的轮询循环实现:

code
async def poll_sqs():
    while not shutdown_event.is_set():
        if len(connections) >= MAX_CONNECTIONS:
            await asyncio.sleep(5)
            continue

        resp = await asyncio.to_thread(
            sqs.receive_message,
            QueueUrl=QUEUE_URL,
            MaxNumberOfMessages=1,
            WaitTimeSeconds=10,
            VisibilityTimeout=30,
        )

        for msg in resp.get("Messages", []):
            body = json.loads(msg["Body"])
            pk = body["pk"]

            if pk in connections:
                sqs.delete_message(
                    QueueUrl=QUEUE_URL,
                    ReceiptHandle=msg["ReceiptHandle"]
                )
                continue

            conn_data = await try_acquire_lease(pk)
            if conn_data and conn_data.get("desired_state") == "STARTED":
                asyncio.create_task(
                    manage_websocket(
                        pk, conn_data["ws_url"],
                        conn_data.get("last_seq", 0)
                    )
                )
            sqs.delete_message(
                QueueUrl=QUEUE_URL,
                ReceiptHandle=msg["ReceiptHandle"]
            )

连接管理

工作者获取租约后,会向上游源建立WebSocket连接,并为该连接的整个生命周期运行三个并发异步任务。这三个任务协同工作以保持连接存活、处理入站数据,并检测连接应停止的时机。

  1. 心跳循环

心跳循环每隔HEARTBEAT_EVERY秒调用renew_lease。如果续约失败(意味着其他工作者已接管所有权或租约记录已变更),循环将立即退出。这是工作者检测在连接进行中失去所有权的机制。

  1. 接收循环

接收循环处理来自上游WebSocket源的每条消息。每条消息都会写入单独的DynamoDB消息表,包含连接ID、时间戳、消息数据和工作者ID。循环将持续运行直到WebSocket连接关闭或发生错误。

  1. 期望状态检查器

每10秒,期望状态检查器会从DynamoDB读取连接记录。如果desired_state被设置为STOPPED(意味着外部系统通过API发送了STOP事件),循环将退出,表明即使WebSocket本身仍保持打开状态,该连接也应被关闭。

三个任务的交互方式

所有三个任务均通过 asyncio.gather 并发执行。当三个任务中的任何一个返回结果或抛出异常时,asyncio.gather 会完成执行并进入 finally 块。这意味着只要触发单个事件(如租约丢失、WebSocket 关闭或 STOP 事件),无论其他两个任务的状态如何,都可以干净地结束连接。

清理

finally 块始终会执行,无论连接如何结束。它会释放 DynamoDB 租约,使其他工作节点可以立即获取连接,并从工作节点的本地跟踪字典中移除该连接。

以下代码展示了完整的连接管理实现:

code
async def manage_websocket(pk: str, ws_url: str, last_seq: int):
    connections[pk] = {"pk": pk, "ws_url": ws_url, "ws": None}

    try:
        async with websockets.connect(ws_url) as ws:
            connections[pk]["ws"] = ws

            async def heartbeat_loop():
                while not shutdown_event.is_set():
                    await asyncio.sleep(HEARTBEAT_EVERY)
                    if not await renew_lease(pk):
                        print(f"[{pk}] Lost lease, closing")
                        return

            async def receive_loop():
                async for msg in ws:
                    data = json.loads(msg)
                    messages_table.put_item(Item={
                        "pk": pk,
                        "sk": str(now_ms()),
                        "message_data": data.get("data", str(data)),
                        "timestamp_ms": now_ms(),
                        "worker_id": WORKER_ID,
                    })

            async def check_desired_state():
                while not shutdown_event.is_set():
                    await asyncio.sleep(10)
                    resp = table.get_item(Key={"pk": pk})
                    if resp.get("Item", {}).get("desired_state") == "STOPPED":
                        return

            await asyncio.gather(
                heartbeat_loop(),
                receive_loop(),
                check_desired_state()
            )

    except Exception as e:
        print(f"[{pk}] WebSocket error: {e}")
    finally:
        await release_lease(pk)
        connections.pop(pk, None)

生产环境注意事项:代码示例中使用 print() 仅用于清晰度。在生产环境中,应将其替换为结构化日志(Python logging 模块或 Amazon CloudWatch Logs),并为租约获取失败和重新连接事件生成 CloudWatch 指标,以支持操作告警。

扩展性注意事项:此处展示的每个连接的 check_desired_state() 循环适用于小型集群。在大规模部署中,应将单个 GetItem 调用替换为单个集中式循环,使用 BatchGetItem 一次性检查所有活动连接的状态,将 DynamoDB 读取操作从每 10 秒 N 次调用减少到 1 次批量调用。

孤儿连接协调

协调循环是系统的安全网。它在每个工作节点上定期运行,独立于 SQS 轮询循环。其唯一目的是查找应处于活动状态但当前没有所有者的连接,并重新获取这些连接。

该循环会查询 GSI 索引中所有 desired_state = STARTED 且 lease_expires_at_ms 小于当前时间的记录。这些是外部系统请求为活动状态的连接,但其租约从未被认领或已过期未续期,表明之前的拥有者已不再运行。

对于每个发现的孤儿连接,工作者会调用 try_acquire_lease 方法。由于 try_acquire_lease 使用了 DynamoDB 的条件写入,多个工作者可以安全地并行执行对账操作,而无需担心重复声明的问题。每个连接恰好只有一个工作者能成功获取租约,其余工作者会收到 ConditionalCheckFailedException 异常并继续执行后续操作。

对账间隔时间(默认:60 秒)决定了意外工作者终止后的最大恢复时间。如果工作者在未执行优雅关闭处理程序的情况下崩溃,其租约将在 LEASE_SECONDS(默认:20 秒)后自然过期。随后对账循环会在下一个 60 秒周期内重新拾取这些连接,最坏情况下的恢复时间约为 80 秒(20 秒租约过期时间加上最多 60 秒的对账间隔)。

以下代码展示了完整实现:

code
async def reconcile_orphaned_connections():
    while not shutdown_event.is_set():
        await asyncio.sleep(RECONCILE_EVERY)

        if len(connections) >= MAX_CONNECTIONS:
            continue

        resp = table.query(
            IndexName=GSI_NAME,
            KeyConditionExpression=(
                "desired_state = :state "
                "AND lease_expires_at_ms < :now"
            ),
            ExpressionAttributeValues={
                ":state": "STARTED",
                ":now": now_ms()
            },
            Limit=RECONCILE_PAGE_SIZE,
        )

        for item in resp.get("Items", []):
            pk = item["pk"]
            if pk not in connections and len(connections) < MAX_CONNECTIONS:
                conn_data = await try_acquire_lease(pk)
                if conn_data:
                    asyncio.create_task(
                        manage_websocket(
                            pk, conn_data["ws_url"],
                            conn_data.get("last_seq", 0)
                        )
                    )

图3:通过孤儿对账实现的自动故障转移

优雅关闭

当 ECS 在滚动部署或缩容事件期间发送 SIGTERM 信号时,工作者在容器被强制终止前只有有限的时间进行清理。与其突然断开连接并等待租约自然过期,工作者会通过三个步骤执行协调关闭。

步骤1:信号传播

当接收到 SIGTERM 信号时,signal_handler 函数会设置一个共享的 shutdown_event。该事件会被所有活动连接中的运行循环检查。心跳循环、期望状态检查器和对账循环在事件被设置后会立即退出 while not shutdown_event.is_set() 循环。无需额外的每连接关闭逻辑,共享事件会自动将关闭信号传播到所有并发任务。

步骤2:并行清理

与其按顺序关闭连接并释放租约(这会随着活动连接数量增加而耗时更长),工作者会使用 asyncio.gather 并行关闭所有 WebSocket 连接并释放所有租约。对于管理数百个连接的工作者来说,这使得总关闭时间大致保持恒定,与连接数量无关。

步骤3:立即释放租约

优雅关闭期间,每个释放的连接都会被设置 lease_expires_at_ms = 0。值为 0 表示租约对运行协调查询的任何工作节点来说都已过期。舰队中的其他工作节点会在下一次协调周期中立即获取这些释放的连接,而不是等待原始租约持续时间(默认:20 秒)自然到期。

与意外终止的对比

优雅关闭是快速路径。当工作节点通过 SIGTERM 清晰退出时,连接在一次协调周期内即可被重新获取。当工作节点在未运行关闭处理程序的情况下意外崩溃时,租约会在 LEASE_SECONDS(默认:20 秒)后自然过期,随后由协调循环获取。两种路径最终结果相同:另一个工作节点会获取连接,但优雅关闭显著更快。

以下代码展示了完整的优雅关闭实现:

code
shutdown_event = asyncio.Event()

def signal_handler(signum, frame):
    shutdown_event.set()

async def graceful_shutdown():
    await shutdown_event.wait()
    tasks = []
    for pk, conn in list(connections.items()):
        if conn.get("ws"):
            tasks.append(conn["ws"].close())
        tasks.append(release_lease(pk))
    await asyncio.gather(*tasks, return_exceptions=True)

设置 shutdown_event 会使心跳循环和状态检查器退出其 while not shutdown_event.is_set() 循环。优雅关闭函数随后并行关闭活动的 WebSocket 连接并释放租约。释放的租约具有 lease_expires_at_ms = 0,这意味着其他工作节点的协调循环会在下一次周期中立即获取这些租约,而不是等待原始租约过期。

扩展舰队规模

每个工作节点都会发布一个自定义的 CloudWatch 指标,包含其活动连接数:

code
async def publish_metrics():
    while not shutdown_event.is_set():
        await asyncio.sleep(30)
        cw.put_metric_data(
            Namespace="WsFleet",
            MetricData=[{
                "MetricName": "ActiveConnections",
                "Value": len(connections),
                "Unit": "Count",
                "Dimensions": [
                    {"Name": "ServiceName", "Value": SERVICE_NAME}
                ],
            }],
        )

AWS 应用程序自动扩展的目标跟踪策略会根据所有工作节点的平均 ActiveConnections 指标扩展舰队规模。当平均值超过目标值(例如每个任务 700 连接)时,ECS 会启动额外任务。新任务启动 SQS 轮询和协调循环,获取新连接并重新平衡舰队。

图 4:基于活动连接数的自动扩展

由于租约模式,缩容是安全的。当 ECS 终止任务时,工作节点会收到 SIGTERM,释放其租约,其他工作节点通过协调获取释放的连接。

配置

值

理由

租约持续时间

20 秒

足够长以应对短暂的网络波动,又足够短以实现快速故障转移

心跳间隔

5 秒

在过期前提前更新(4 倍安全余量)

协调间隔

60 秒

在恢复速度和 DynamoDB 读取成本之间取得平衡

每任务最大连接数

700

基于每连接的内存和 CPU 分析

扩展冷却时间

2 分钟

防止流量高峰期间频繁扩展

缩容冷却时间

15 分钟

允许连接稳定后再移除容量

调优指南。这些值仅作为起点,根据实际需求进行调整:

  • 租约持续时间:建议从20秒开始。需要更快故障转移时可减少该值,网络波动导致租约误过期时可增加该值。
  • 心跳间隔:需保持低于租约持续时间。租约与心跳间隔保持4:1比例(租约:心跳)时,租约过期前会有4次续租尝试。
  • 对账间隔:建议从60秒开始。需要更快恢复意外终止连接时可减少该值,需要降低DynamoDB读取成本时可增加该值。
  • 每个任务最大连接数:建议从100开始,并在CloudWatch Container Insights中监控内存和CPU使用情况逐步增加。每个WebSocket连接通常消耗2-5MB内存,具体取决于消息吞吐量。

DynamoDB成本考量

该架构的主要成本驱动因素是心跳写入操作。每个活跃连接每颗心跳间隔会产生一次update_item调用,消耗1个WCUs。在默认5秒心跳间隔下:

| 活跃连接数 | 每秒WCUs | 月度成本(按需) | 月度成本(预置) | |------------|-----------|------------------|------------------| | 100 | 20 | ~$65 | ~$10 | | 500 | ~$325 | ~$47 | | | 2,000 | 400 | ~$1,300 | ~$190 |

对于需要长期维持高连接数的生产环境部署,建议使用预置容量配合自动扩展而非按需计费。心跳写入具有可预测性和稳定性,这使其非常适配预置吞吐量。在预置容量上配置自动扩展,可跟踪舰队规模变化时的连接数变动。

为降低成本,可考虑以下调整:

  • 增加心跳间隔。将心跳间隔从5秒延长至10秒可使WCUs消耗减半。同时将租约持续时间也延长一倍以保持4:1的租约与心跳比例。这会按比例增加故障转移窗口。
  • 增加对账间隔。将对账间隔从60秒延长至120秒可使对账查询的RCUs消耗减半。这会降低意外终止后的恢复速度。
  • 使用BatchGetItem进行期望状态检查。将check_desired_state循环中的每个连接get_item调用替换为覆盖所有活跃连接的单个BatchGetItem调用。这可将每个周期的RCU消耗从N次读取降低到1次批量读取。对账期间的GSI查询默认使用最终一致性读取,相比强一致性读取可使RCU成本减半。在DynamoDB控制台监控GSI读取消耗,并调整对账分页大小和间隔以控制在成本目标范围内。

结论

规模化管理长生命周期WebSocket连接需要显式的所有权追踪、自动故障转移以及跨工作者舰队的协调。本文展示了如何利用DynamoDB条件写入作为分布式租约机制,解决这些挑战。

关键要点:

  • 可使用DynamoDB条件写入实现原子分布式协调,无需外部锁服务。update_item的ConditionExpression可确保每次只有一个工作者拥有每个连接。
  • 心跳与对账模式可处理完整的故障场景。租约过期可检测意外工作者终止,优雅关闭可处理滚动部署。新工作者获取租约,离线工作者释放租约,使扩展过程安全可靠。
  • 此模式适用于需要大规模管理长期WebSocket连接的系统:实时语音转录、物联网数据采集、金融数据流处理或直播事件流媒体。

Getting started

完整的实现方案(包括工作程序应用、Lambda事件路由器,以及DynamoDB表、SQS队列和ECS集群的Terraform模板)已发布在GitHub仓库中。按照仓库README中的说明部署基础设施,并通过少量测试连接验证租约生命周期。

如需进一步增强功能,可添加AWS X-Ray分布式追踪以实现跨工作程序的端到端可视化,并实现重连逻辑(上游回放或基于偏移量的恢复机制),以处理工作程序故障与恢复期间的数据缺口。

Further reading

  • Amazon DynamoDB 条件写入 – 了解条件表达式在实现分布式锁中的应用。
  • Amazon DynamoDB 全局二级索引 – 掌握如何高效查询过期租约。
  • AWS Fargate 上的 Amazon ECS – 学习如何在不管理服务器的情况下运行容器化工作程序。
  • 应用程序自动扩展 – 配置基于自定义指标的目标跟踪策略。
  • Amazon CloudWatch 自定义指标 – 发布用于扩展决策的应用程序层级指标。

'"` /think