Airflow+Docker Swarm部署下SIGTERM错误排查求助
问题排查:Airflow自定义Hive Operator任务SIGTERM异常
问题背景
- 部署环境:5台服务器通过Docker Swarm部署Airflow,稳定运行2个月后出现异常
- 异常现象:使用自定义Hive Operator的DAG频繁抛出SIGTERM错误,重试任务成功率随机;但Hive任务实际未失败——Airflow标记任务失败110分钟后,Hadoop日志显示Hive查询已完成,流程为:任务启动→510分钟→Airflow抛出SIGTERM→1~10分钟→Hive任务成功
- 已尝试操作:重启Airflow服务器,问题未解决
Airflow任务日志
[2023-01-09 08:06:07,583] {local_task_job.py:208} WARNING - State of this instance has been externally set to up_for_retry. Terminating instance. [2023-01-09 08:06:07,588] {process_utils.py:100} INFO - Sending Signals.SIGTERM to GPID 135213 [2023-01-09 08:06:07,588] {taskinstance.py:1236} ERROR - Received SIGTERM. Terminating subprocesses. [2023-01-09 08:13:42,510] {taskinstance.py:1463} ERROR - Task failed with exception Traceback (most recent call last): File "/opt/airflow/dags/common/operator/hive_q_operator.py", line 81, in execute cur.execute(statement) # hive query custom operator File "/home/airflow/.local/lib/python3.8/site-packages/pyhive/hive.py", line 454, in execute response = self._connection.client.ExecuteStatement(req) File "/home/airflow/.local/lib/python3.8/site-packages/TCLIService/TCLIService.py", line 280, in ExecuteStatement return self.recv_ExecuteStatement() File "/home/airflow/.local/lib/python3.8/site-packages/TCLIService/TCLIService.py", line 292, in recv_ExecuteStatement (fname, mtype, rseqid) = iprot.readMessageBegin() File "/home/airflow/.local/lib/python3.8/site-packages/thrift/protocol/TBinaryProtocol.py", line 134, in readMessageBegin sz = self.readI32() File "/home/airflow/.local/lib/python3.8/site-packages/thrift/protocol/TBinaryProtocol.py", line 217, in readI32 buff = self.trans.readAll(4) File "/home/airflow/.local/lib/python3.8/site-packages/thrift/transport/TTransport.py", line 62, in readAll chunk = self.read(sz - have) File "/home/airflow/.local/lib/python3.8/site-packages/thrift_sasl/__init__.py", line 173, in read self._read_frame() File "/home/airflow/.local/lib/python3.8/site-packages/thrift_sasl/__init__.py", line 177, in _read_frame header = self._trans_read_all(4) File "/home/airflow/.local/lib/python3.8/site-packages/thrift_sasl/__init__.py", line 210, in _trans_read_all return read_all(sz) File "/home/airflow/.local/lib/python3.8/site-packages/thrift/transport/TTransport.py", line 62, in readAll chunk = self.read(sz - have) File "/home/airflow/.local/lib/python3.8/site-packages/thrift/transport/TSocket.py", line 150, in read buff = self.handle.recv(sz) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 1238, in signal_handler raise AirflowException("Task received SIGTERM signal") airflow.exceptions.AirflowException: Task received SIGTERM signal
Flower日志
Traceback (most recent call last): File "/home/airflow/.local/lib/python3.8/site-packages/celery/app/trace.py", line 412, in trace_task R = retval = fun(*args, **kwargs) File "/home/airflow/.local/lib/python3.8/site-packages/celery/app/trace.py", line 704, in __protected_call__ return self.run(*args, **kwargs) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/executors/celery_executor.py", line 88, in execute_command _execute_in_fork(command_to_exec) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/executors/celery_executor.py", line 99, in _execute_in_fork raise AirflowException('Celery command failed on host: ' + get_hostname()) airflow.exceptions.AirflowException: Celery command failed on host: 8be4caa25d17
排查指导
1. 检查Celery Worker资源限制
- 查看Docker Swarm中Celery Worker容器的CPU、内存配额,确认是否因资源耗尽触发OOM Killer或Swarm调度终止进程
- 执行
docker stats实时监控容器资源使用,同时查看宿主机dmesg日志,排查是否有OOM事件记录
2. 排查Airflow与Celery的心跳超时配置
- 检查Airflow配置中
celery.worker_send_task_events、celery.task_track_started是否开启,确保Worker状态能被正确追踪 - 核对
broker_heartbeat、worker_prefetch_multiplier等参数,避免Worker长时间阻塞无响应被Broker标记为死亡,触发任务重试并终止原进程 - 检查Broker(RabbitMQ/Redis)的连接状态与消息堆积情况,确认是否存在连接中断
3. 调整自定义Hive Operator的阻塞逻辑
- 当前自定义Operator中
cur.execute(statement)为同步阻塞调用,Hive长查询会导致Worker进程长时间无响应,可能触发Celery的task_time_limit或task_soft_time_limit - 检查Airflow配置中
core.dag_run_timeout、celery.task_time_limit是否设置过短,若存在则延长超时时间 - 修改自定义Operator,改用异步提交Hive任务+轮询状态的方式,避免Worker进程长时间阻塞
4. 优化Thrift连接稳定性
- 日志显示错误发生在Thrift读取响应阶段,可能是Airflow与HiveServer2的网络连接不稳定导致超时或中断
- 查看HiveServer2日志,确认是否存在连接超时、断开的记录
- 在自定义Operator中添加pyhive的Thrift连接超时配置,增加重试机制
5. 排查Docker Swarm集群状态
- 检查Swarm集群节点的资源波动、网络分区情况,确认是否有Worker容器被迁移或网络中断
- 查看Swarm服务日志,排查是否存在容器重启、调度异常事件
内容的提问来源于stack exchange,提问作者mslee
相关产品推荐
相关产品推荐

