改造前后对比
改造前
- 40 台按需 c5.4xlarge 常驻,夜间几乎空转
- 单个视频串行转码,2 小时素材要跑 50 分钟
- 队列积压时无法快速扩容,靠人工加机器
- 失败任务需要人工重跑,没有自动重试
- 转码进度不可见,客服无法回答「什么时候能好」
改造后
- AWS Batch 管理 Spot 池,无任务时容量降到 0
- 视频切片后并行转码,2 小时素材 24 分钟完成
- 队列深度驱动自动扩容,峰值可拉到 800 并发
- Spot 中断自动重试,分片级重试而非整个任务
- Step Functions 记录每个分片状态,进度实时可查
转码流水线
-
01
上传与探测
素材上传 S3 触发 Lambda,探测视频规格(分辨率、码率、时长、音轨)并写入任务表。
交付物:任务记录、分片计划
-
02
智能分片
按关键帧位置切片,每片 3-5 分钟。短视频不切片,避免分片开销大于收益。
交付物:分片清单、分片文件
-
03
并行转码
分片投递到 SQS,AWS Batch 在 Spot 池上并行处理。每片输出多档码率。
交付物:各档位分片
-
04
合并与打包
所有分片完成后合并,生成 HLS/DASH 清单,写入 CDN 源桶。
交付物:HLS 切片、播放清单
-
05
校验与发布
自动校验时长、码率、音视频同步,通过后标记可播放并通知业务系统。
交付物:质检报告、发布通知
Spot 中断的分片级重试
"""分片转码任务的中断处理。
关键设计:
1. 分片是幂等的,重跑不会产生副作用
2. Spot 中断通知触发主动上报,而不是等超时
3. 重试计数在 DynamoDB 里,避免无限重试
4. 连续失败的分片降级到按需实例,保证最终完成
"""
from __future__ import annotations
import json
import os
from datetime import datetime, timezone
import boto3
ddb = boto3.resource("dynamodb")
sqs = boto3.client("sqs")
TASK_TABLE = ddb.Table(os.environ["TASK_TABLE"])
SPOT_QUEUE = os.environ["SPOT_QUEUE_URL"]
ONDEMAND_QUEUE = os.environ["ONDEMAND_QUEUE_URL"]
MAX_SPOT_RETRY = 3
def requeue_shard(job_id: str, shard_id: str, reason: str) -> dict:
"""重新投递分片。超过 Spot 重试上限则转到按需队列。"""
resp = TASK_TABLE.update_item(
Key={"job_id": job_id, "shard_id": shard_id},
UpdateExpression=(
"SET #st = :pending, last_failure = :reason, updated_at = :now "
"ADD retry_count :one"
),
ExpressionAttributeNames={"#st": "status"},
ExpressionAttributeValues={
":pending": "pending",
":reason": reason,
":now": datetime.now(timezone.utc).isoformat(),
":one": 1,
},
ReturnValues="ALL_NEW",
)
item = resp["Attributes"]
retries = int(item.get("retry_count", 0))
# 重试次数用完就上按需,宁可贵一点也要把任务做完
queue = ONDEMAND_QUEUE if retries > MAX_SPOT_RETRY else SPOT_QUEUE
sqs.send_message(
QueueUrl=queue,
MessageBody=json.dumps({
"job_id": job_id,
"shard_id": shard_id,
"input_key": item["input_key"],
"output_prefix": item["output_prefix"],
"profile": item["profile"],
"attempt": retries + 1,
}),
MessageGroupId=job_id,
)
return {
"job_id": job_id,
"shard_id": shard_id,
"retry_count": retries,
"queue": "on-demand" if retries > MAX_SPOT_RETRY else "spot",
}
def lambda_handler(event, context): # noqa: ARG001
"""处理两类事件:Spot 中断警告,以及 Batch 任务失败。"""
source = event.get("source",")
detail = event.get("detail", {})
results = []
if source == "aws.ec2" and event.get("detail-type") == "EC2 Spot Instance Interruption Warning":
instance_id = detail["instance-id"]
# 找出该实例上正在跑的所有分片,立即重新投递
running = TASK_TABLE.query(
IndexName="instance-index",
KeyConditionExpression=boto3.dynamodb.conditions.Key("instance_id").eq(instance_id),
FilterExpression=boto3.dynamodb.conditions.Attr("status").eq("running"),
)["Items"]
for shard in running:
results.append(requeue_shard(
shard["job_id"], shard["shard_id"], "spot-interruption"
))
elif source == "aws.batch" and detail.get("status") == "FAILED":
results.append(requeue_shard(
detail["parameters"]["job_id"],
detail["parameters"]["shard_id"],
detail.get("statusReason", "batch-failed"),
))
return {"requeued": len(results), "shards": results}
Spot 中断有 2 分钟通知窗口,主动重新投递比等任务超时快得多,用户几乎感知不到。
为什么不直接用 MediaConvert
评估过。MediaConvert 更省心,但客户有自定义的水印、片头片尾拼接和特定编码参数需求,用自建 FFmpeg 流水线更灵活。标准转码需求我们仍然推荐直接用 MediaConvert,省下的运维精力通常比省下的机器钱更值。
技术栈
- EC2 Spot
- AWS Batch
- S3
- SQS
- Lambda
- Step Functions
- MediaConvert
- CloudFront
- DynamoDB