Apache Airflow 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

