Django+Celery+SQS场景下ECS容器合理缩容方案问询
解决Celery Worker长任务期间ECS缩容问题的方案
针对ECS缩容会杀掉正在执行长任务的Celery Worker问题,这里提供几个实用的落地方案,可单独或组合使用:
方案一:暴露Worker活跃任务指标,基于指标触发缩容
核心思路是让ECS只在Worker没有活跃任务时才执行缩容,需要给每个Worker上报自定义的活跃任务数指标到CloudWatch(ECS原生支持基于CloudWatch指标做扩缩容)。
实现步骤:
- 在Celery Worker中添加信号监听,追踪任务的开始与结束,维护当前Worker的活跃任务数。
- 定期将活跃任务数上报到CloudWatch。
- 调整ECS的缩容策略,将触发条件从
NumberOfMessagesSent改为自定义的Celery/Workers/ActiveTasks指标(建议设置为该指标持续5分钟为0时触发缩容)。
代码示例(Worker端):
from celery import signals import boto3 import os import time from threading import Thread cloudwatch = boto3.client('cloudwatch') # 用ECS容器ID作为Worker唯一标识,方便追踪 WORKER_ID = os.environ.get('ECS_CONTAINER_ID', f'worker-{os.getpid()}') ACTIVE_TASKS = 0 # 监听任务开始信号 @signals.task_prerun.connect def increment_active_tasks(sender=None, **kwargs): global ACTIVE_TASKS ACTIVE_TASKS += 1 # 监听任务结束信号(成功/失败都触发) @signals.task_postrun.connect def decrement_active_tasks(sender=None, **kwargs): global ACTIVE_TASKS ACTIVE_TASKS = max(0, ACTIVE_TASKS - 1) # 后台线程定期上报指标(每30秒一次) def report_metrics(): while True: cloudwatch.put_metric_data( Namespace='Celery/Workers', MetricData=[ { 'MetricName': 'ActiveTasks', 'Dimensions': [{'Name': 'WorkerID', 'Value': WORKER_ID}], 'Value': ACTIVE_TASKS, 'Unit': 'Count' } ] ) time.sleep(30) # 启动指标上报线程 Thread(target=report_metrics, daemon=True).start()
方案二:配置Worker优雅退出,等待长任务完成
ECS在缩容时会先给容器发送SIGTERM信号,默认等待30秒后强制杀掉容器。我们可以让Celery Worker捕获该信号,停止接收新任务,等待所有正在执行的任务完成后再退出,同时调整ECS的停止超时时间匹配你的最长任务时长。
实现步骤:
- 调整Celery Worker的启动参数,启用优雅关闭并设置超时时间。
- 在Django的Celery配置中设置Worker shutdown等待超时。
- 修改ECS任务定义,将容器的停止超时时间(
stopTimeout)设置为你的最长任务时长(比如4小时=14400秒)。
配置示例:
- Worker启动命令:
celery -A your_django_app worker --loglevel=info --concurrency=4 --soft-timeout=14400 --time-limit=18000
--soft-timeout:任务软超时,超过时间后发送SIGTERM给任务进程--time-limit:任务硬超时,超过时间直接强制终止Django settings.py中的Celery配置:
CELERY_WORKER_SHUTDOWN_WAIT_TIMEOUT = 14400 # 最大等待4小时,确保长任务完成 CELERY_WORKER_DISABLE_RATE_LIMITS = True
- ECS任务定义中的容器配置(JSON片段):
"containerDefinitions": [ { "name": "celery-worker", "image": "your-image:tag", "stopTimeout": 14400, ... } ]
方案三:利用ECS生命周期钩子,延迟终止有活跃任务的Worker
通过ECS生命周期钩子,在容器即将终止前触发Lambda函数,检查该Worker是否有正在执行的任务,有则延迟终止,无则允许缩容。
实现步骤:
- 给ECS集群添加生命周期钩子,绑定到容器的TERMINATING状态。
- 给每个Worker容器打上
WorkerID标签,用于关联到Celery Worker实例。 - 编写Lambda函数,查询Celery的活跃任务列表,判断是否允许终止容器:
- 如果Worker有活跃任务,调用
RecordLifecycleActionHeartbeat延迟终止(每次心跳可延长30分钟) - 如果无活跃任务,调用
CompleteLifecycleAction允许终止
- 如果Worker有活跃任务,调用
Lambda核心代码示例:
import boto3 from celery import Celery ecs_client = boto3.client('ecs') # 初始化Celery客户端,与Worker使用相同的配置 celery_app = Celery('your_django_app', broker='sqs://') def lambda_handler(event, context): # 解析生命周期钩子事件参数 lifecycle_token = event['detail']['LifecycleActionToken'] cluster_name = event['detail']['ClusterName'] instance_arn = event['detail']['ContainerInstanceARN'] instance_id = instance_arn.split('/')[-1] # 获取容器的WorkerID标签 instance_details = ecs_client.describe_container_instances( cluster=cluster_name, containerInstances=[instance_id] )['containerInstances'][0] worker_id = next((tag['value'] for tag in instance_details['tags'] if tag['key'] == 'WorkerID'), None) if not worker_id: # 无法获取WorkerID,直接允许终止 ecs_client.complete_lifecycle_action( Cluster=cluster_name, LifecycleActionToken=lifecycle_token, LifecycleHookName='celery-worker-shutdown-hook', LifecycleActionResult='CONTINUE' ) return # 查询该Worker的活跃任务 inspect_result = celery_app.control.inspect([worker_id]).active() has_active_tasks = bool(inspect_result and inspect_result.get(worker_id)) if has_active_tasks: # 有活跃任务,发送心跳延迟终止 ecs_client.record_lifecycle_action_heartbeat( Cluster=cluster_name, LifecycleActionToken=lifecycle_token, LifecycleHookName='celery-worker-shutdown-hook' ) else: # 无活跃任务,允许终止 ecs_client.complete_lifecycle_action( Cluster=cluster_name, LifecycleActionToken=lifecycle_token, LifecycleHookName='celery-worker-shutdown-hook', LifecycleActionResult='CONTINUE' )
推荐组合方案
优先使用方案一+方案二的组合:
- 用方案一的自定义指标确保ECS只在Worker无任务时触发缩容,从源头避免不必要的终止。
- 用方案二的优雅退出作为兜底,即使ECS触发了缩容,Worker也会等待现有任务完成后再退出。
内容的提问来源于stack exchange,提问作者Adrian
相关产品推荐
相关产品推荐

