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

Apache Airflow SparkSQLOperator持续打印空日志无法终止求助

Airflow 1.9.0 SparkSqlOperator 任务成功后持续输出空日志无法终止

我最近在使用Airflow v1.9.0搭建调度任务时遇到了一个头疼的问题:编写的DAG用SparkSqlOperator执行SQL查询,任务能成功返回正确结果,但执行airflow scheduler调度后,任务完成后会持续打印空日志,直到手动用Ctrl+C终止调度器。

我的DAG代码如下:

import airflow
from airflow import DAG
from airflow.contrib.operators.spark_sql_operator import SparkSqlOperator
from datetime import timedelta
from datetime import datetime as dt

default_args = {
    'owner': 'zxy',
    'depends_on_past': False,
    'email_on_failure': True,
    'email_on_retry': True,
    'retries': 1,
    'retry_delay': timedelta(minutes=5)
}

dag = DAG(
    'my_first_dag',
    default_args=default_args,
    #start_date=dt.strptime('2018-05-16', '%Y-%m-%d'),
    start_date=airflow.utils.dates.days_ago(2),
    description='My First Airflow DAG',
    schedule_interval=timedelta(minutes=5))

sql = r'''select count(u) from some_table where time=20180513 and platform='iOS' '''

t1 = SparkSqlOperator(task_id='Count_Ads_U', conn_id='spark_default',sql=sql, dag=dag)

问题日志表现:

任务执行完成后,日志会先输出Spark正常关闭的信息,然后开始持续打印空日志:

[2018-05-16 06:33:07,505] {base_task_runner.py:98} INFO - Subtask: [2018-05-16 06:33:07,505] {spark_sql_hook.py:142} INFO - b'18/05/16 06:33:07 INFO spark. SparkContext: Successfully stopped SparkContext\n'
[2018-05-16 06:33:07,506] {base_task_runner.py:98} INFO - Subtask: [2018-05-16 06:33:07,506] {spark_sql_hook.py:142} INFO - b'18/05/16 06:33:07 INFO util. ShutdownHookManager: Shutdown hook called\n'
[2018-05-16 06:33:07,506] {base_task_runner.py:98} INFO - Subtask: [2018-05-16 06:33:07,506] {spark_sql_hook.py:142} INFO - b'18/05/16 06:33:07 INFO util. ShutdownHookManager: Deleting directory /tmp/spark-fbb4089c-338b-4b0e-a394-975f45b307a8\n'
[2018-05-16 06:33:07,509] {base_task_runner.py:98} INFO - Subtask: [2018-05-16 06:33:07,509] {spark_sql_hook.py:142} INFO - b'18/05/16 06:33:07 INFO util. ShutdownHookManager: Deleting directory /apps/data/spark/temp/spark-f6b6695f-24e4-4db0-ae2b-29b6836ab9c3\n'
[2018-05-16 06:33:07,902] {base_task_runner.py:98} INFO - Subtask: [2018-05-16 06:33:07,902] {spark_sql_hook.py:142} INFO - b''
[2018-05-16 06:33:07,903] {base_task_runner.py:98} INFO - Subtask: [2018-05-16 06:33:07,902] {spark_sql_hook.py:142} INFO - b''
...(后续重复输出空日志,直至手动终止)

问题原因

这是Airflow 1.9.0版本中SparkSqlHook的一个已知bug:在读取Spark作业的输出流时,没有正确判断流的结束条件,导致循环持续读取空内容并打印空日志。

解决方案

针对这个问题,有几种可行的解决办法:

1. 升级Airflow到1.10.0及以上版本

官方在Airflow 1.10.0版本中已经修复了这个日志读取的问题,直接升级到稳定的新版本是最省心的方案。执行以下命令完成升级:

pip install --upgrade apache-airflow

升级完成后重启Airflow scheduler,问题即可解决。

2. 手动修改SparkSqlHook源码(临时修复)

如果暂时无法升级版本,可以手动修改Airflow的spark_sql_hook.py文件来修复:

  • 找到Airflow安装目录下的airflow/contrib/hooks/spark_sql_hook.py文件
  • 定位到日志读取的循环代码(大概在140行左右的位置),添加对空内容的判断,当读取到空字节串时退出循环:
# 原代码(类似):
while True:
    line = proc.stdout.readline()
    self.log.info(line)

# 修改为:
while True:
    line = proc.stdout.readline()
    if not line:  # 添加空内容判断,退出循环
        break
    self.log.info(line)
  • 保存文件后重启Airflow scheduler服务,空日志循环的问题就会消失。

3. 改用SparkSubmitOperator替代SparkSqlOperator

如果不想修改源码或升级,也可以换用SparkSubmitOperator来实现相同的SQL查询任务,避免触发SparkSqlOperator的日志bug。示例代码如下:

from airflow.contrib.operators.spark_submit_operator import SparkSubmitOperator

# 定义SparkSubmit任务
t1 = SparkSubmitOperator(
    task_id='Count_Ads_U',
    conn_id='spark_default',
    application='/path/to/your/count_ads_u.py',  # 你的Spark SQL脚本路径
    dag=dag
)

对应的count_ads_u.py脚本内容:

from pyspark.sql import SparkSession

# 初始化SparkSession
spark = SparkSession.builder.appName("CountAdsUser").getOrCreate()
# 执行SQL查询
spark.sql("select count(u) from some_table where time=20180513 and platform='iOS'").show()
# 关闭SparkSession
spark.stop()

内容的提问来源于stack exchange,提问作者Derrick Zhang

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 04:12:55