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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 07:45:24