MWAA Airflow运行高频长时脚本出现Negsignal.SIGKILL报错如何扩容
根因说明
你遇到的报错:
[2021-09-25 20:33:16,472] {{local_task_job.py:102}} INFO - Task exited with return code Negsignal.SIGKILL
本质是任务进程被系统OOM(内存不足)杀手强制终止,你观测到的worker集群整体负载未达上限是平均指标,实际单个任务运行所在的执行单元内存配额不足,直接触发了强制kill。
另外当前单任务单次运行需要8分钟,远高于1分钟的调度间隔,若不做优化,每分钟会新增1个待处理任务,任务积压会随时间持续加重,进一步拉长所有任务的运行耗时。
优化方案
1. DAG任务逻辑优化(优先级最高)
- 拆分单一大任务为多个并行子任务:使用Airflow动态任务映射功能,将200MB JSON的处理逻辑按文件、数据分片维度拆为多个小任务并行执行,尽可能将单次DAG运行的总耗时压缩到1分钟以内,从根源解决任务积压问题
- 优化内存占用逻辑:不要全量加载200MB JSON到内存后再处理,改用
ijson+流式S3读取接口逐行迭代解析数据,处理过程中及时释放无用变量的内存引用,避免内存溢出 - 移除不必要的中间结果缓存,不要将全量过滤结果暂存在内存中,直接流式写入S3桶B
2. MWAA环境参数调整
- 调整Celery worker单任务内存配额:在MWAA的自定义Airflow配置项中新增
celery.worker_max_memory_per_child,按你任务的实际内存占用设置合理值,比如设置为2048000(对应2GB单任务内存上限),避免任务被容器的内存限制提前kill - 调整单worker并发数:mw1.large规格的worker默认并发数为4,你可以将
celery.worker_concurrency调整为8~10,在不新增worker节点的前提下提升单节点的任务处理能力 - 优化调度器配置:调大
scheduler.parsing_processes和scheduler.max_dagruns_to_create_per_loop参数,提升调度器分发任务的速度,避免任务在调度队列中长时间排队
3. 调度规则优化
- 给DAG设置合理的最大活跃运行数:配置
max_active_runs=12(可按你集群的并行处理能力调整),避免无限制生成待运行任务抢占资源,导致所有任务运行速度都下降 - 关闭DAG的历史补跑功能:设置
catchup=False,避免调度器生成大量历史待运行任务挤占现有资源 - 给单个任务设置超时时间:配置
execution_timeout=timedelta(minutes=10),异常卡住的任务会自动终止释放资源,不会长期占用worker算力
内容的提问来源于stack exchange,提问作者Kei
相关产品推荐
相关产品推荐

