实时流处理服务经常需要维护数百条持久 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 恢复订阅,并周期性续租。续租失败意味着所有权已经丢失,进程必须停止发布数据并关闭连接,不能继续以“旧主人”的身份运行。
租约通常形成以下状态循环:
- 工作进程通过条件写入竞争空闲或已过期的分片。
- 获胜者建立 WebSocket,并在后台定期续租。
- 每处理一批消息,工作进程更新可恢复的检查点。
- 任务崩溃后不再续租,租约自然过期。
- 其他任务获取租约,使用检查点重新连接和补读。
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 权限应限制到目标表和必要操作,例如
GetItem、UpdateItem、Query。 - 应监控租约获取失败率、续租延迟、过期租约数量、WebSocket 重连次数、检查点滞后和条件写入冲突。
上线前,至少验证三类场景:强制停止一个 Fargate 任务、让工作进程暂停超过租约周期、在持续流量下执行滚动发布。验收标准不只是连接最终恢复,还要检查消息缺口、重复率、接管时间,以及旧所有者是否在丢失租约后立即停止产生副作用。