You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于SQS Queue触发ECS Fargate任务的架构实现及Boto3代码示例

基于ECS Fargate + SQS的按需计算任务架构实现(Python Boto3)

整体架构说明

核心流程:

  • 应用后端将计算任务参数写入SQS标准队列
  • ECS Fargate服务运行的任务持续轮询SQS队列,获取任务并执行计算
  • 借助ECS自动扩缩容能力,根据SQS队列的待处理消息数自动调整Fargate任务数量
  • 计算完成后可将结果写入S3/数据库或另一个SQS队列供后端读取

前置AWS资源准备

开始编码前,需先创建以下AWS资源:

  • SQS队列:存储待处理计算任务,建议开启长轮询以减少空轮询
  • ECS集群:托管Fargate任务的集群
  • ECS任务定义:定义Fargate任务的镜像、资源配置、环境变量(需包含SQS队列URL等参数)
  • IAM角色:
    • 后端角色:拥有SQS SendMessage权限
    • Fargate任务角色:拥有SQS ReceiveMessage/DeleteMessage权限,以及计算所需的其他权限(如S3读写)

代码实现

1. 应用后端:向SQS写入任务(Python Boto3)

这段代码负责将计算任务参数序列化后发送到SQS队列:

import boto3
import json

# 初始化SQS客户端
sqs = boto3.client('sqs', region_name='us-east-1')
QUEUE_URL = '你的SQS队列URL'

def send_compute_task(task_data):
    """发送计算任务到SQS队列"""
    try:
        # 序列化任务数据为JSON字符串
        message_body = json.dumps(task_data)
        response = sqs.send_message(
            QueueUrl=QUEUE_URL,
            MessageBody=message_body
            # 可选:FIFO队列需添加MessageGroupId
            # MessageGroupId='compute-group'
        )
        print(f"任务已发送,消息ID: {response['MessageId']}")
        return response
    except Exception as e:
        print(f"发送任务失败: {str(e)}")
        raise

# 示例:发送一个计算任务
if __name__ == "__main__":
    sample_task = {
        "task_id": "task_001",
        "compute_type": "image_processing",
        "input_path": "s3://my-bucket/inputs/image.jpg",
        "output_path": "s3://my-bucket/outputs/image_processed.jpg"
    }
    send_compute_task(sample_task)

2. Fargate任务:轮询SQS并执行计算(Python Boto3)

这段代码是Fargate任务的核心运行逻辑,持续轮询SQS获取任务并执行计算:

import boto3
import json
import time
from my_compute_module import run_computation  # 替换为你的计算逻辑模块

# 初始化SQS客户端
sqs = boto3.client('sqs', region_name='us-east-1')
QUEUE_URL = '你的SQS队列URL'
WAIT_TIME_SECONDS = 20  # 长轮询等待时间,减少空请求

def poll_sqs_and_process():
    """持续轮询SQS并处理任务"""
    while True:
        try:
            # 长轮询获取消息
            response = sqs.receive_message(
                QueueUrl=QUEUE_URL,
                MaxNumberOfMessages=1,
                WaitTimeSeconds=WAIT_TIME_SECONDS,
                MessageAttributeNames=['All']
            )

            if 'Messages' in response:
                for message in response['Messages']:
                    try:
                        # 解析任务数据
                        task_data = json.loads(message['Body'])
                        print(f"开始处理任务: {task_data['task_id']}")

                        # 执行计算逻辑
                        result = run_computation(task_data)
                        print(f"任务处理完成: {task_data['task_id']}, 结果: {result}")

                        # 处理完成后删除消息,避免重复消费
                        sqs.delete_message(
                            QueueUrl=QUEUE_URL,
                            ReceiptHandle=message['ReceiptHandle']
                        )
                    except Exception as e:
                        print(f"处理任务失败: {str(e)}")
                        # 可选:将失败任务发送到死信队列或重试
                        continue
            else:
                print("当前队列无消息,等待下一轮轮询")
        except Exception as e:
            print(f"轮询SQS失败: {str(e)}")
            time.sleep(5)  # 出错后等待5秒再重试

if __name__ == "__main__":
    poll_sqs_and_process()

注意:my_compute_module.py是自定义的计算逻辑模块,比如图像处理、数据运算等,需根据实际需求实现。

3. 配置ECS自动扩缩容

要实现按需扩缩容,需通过ECS服务的自动扩缩容策略,基于SQS队列的ApproximateNumberOfMessagesVisible指标调整任务数量:

  1. 进入ECS控制台,找到目标Fargate服务,打开"自动扩缩容"配置
  2. 创建新的扩缩容策略:
    • 选择CloudWatch指标,指标名称为ApproximateNumberOfMessagesVisible,命名空间为AWS/SQS,维度为目标队列名称
    • 设置目标值:比如每个Fargate任务处理10条消息,目标值设为10
    • 配置扩缩容范围:最小任务数1,最大任务数10(可根据实际需求调整)
  3. 保存配置后,ECS会自动根据队列长度调整Fargate任务数量

关键注意事项

  • SQS长轮询:必须开启长轮询(WaitTimeSeconds>0),减少空轮询次数,降低成本
  • 消息幂等性:SQS可能存在重复消息,计算任务需实现幂等性,确保重复执行不会产生错误结果
  • 错误处理:处理失败的任务建议配置死信队列,便于后续排查和重试
  • 资源配置:根据计算任务的资源需求,合理设置Fargate任务的CPU和内存配置,避免资源不足或浪费

内容的提问来源于stack exchange,提问作者karkir subu

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.06 13:15:51