【发布时间】:2020-03-29 08:50:36
【问题描述】:
我是 Spark 和 Airflow 的新手,我正在尝试创建一个在 pyspark 中运行 spark 提交作业的 DAG。
在我的 Ubuntu 系统中,我创建了一个名为“hadoopusr”的用户,通过它我手动运行我的 spark 提交。所有环境变量都设置在该用户下的/.bashrc中。
当我从终端手动运行 spark-submit 时,作业运行成功。
我创建了一个示例 DAG 文件,如下所示:
"""
Apache Airflow DAG Script: Ingestion MVP basic version
"""
#imports
from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from datetime import datetime, timedelta
import os
import sys
#Paths
os.environ['SPARK_HOME'] = '/usr/local/spark/spark-2.4.4-bin-hadoop2.7'
sys.path.append(os.path.join(os.environ['SPARK_HOME'], 'bin'))
#Parameters
filename="CDC_DAY1"
file_extension="csv"
batch_id="20191203"
## Define the DAG object
default_args = {
'owner': 'hadoopusr',
'depends_on_past': False,
'start_date': datetime(2019, 12, 3),
'retries': 5,
'retry_delay': timedelta(minutes=1),
}
dag = DAG('MVP_BASIC_2', default_args=default_args, schedule_interval=timedelta(1))
#List of tasks
LandingToRaw = BashOperator(
task_id='Landing_to_Raw',
bash_command="spark-submit --master yarn-client /home/soham/Documents/mvp_landing_raw_1.py "+filename+" "+file_extension+" "+batch_id,
dag=dag,
run_as_user='hadoopusr')
RawToStaging = BashOperator(
task_id='Raw_to_Staging',
bash_command="spark-submit --master yarn-client /home/soham/Documents/mvp_raw_staging_1.py "+filename+" "+file_extension+" "+batch_id+" FULL emp_id",
dag=dag,
run_as_user='hadoopusr')
#Dependencies
LandingToRaw >> RawToStaging
当我在同一用户 (hadoopusr) 下使用以下命令测试 DAG 的第一个任务时,它会引发如下异常:
Running command: spark-submit --master yarn-client /home/soham/Documents/mvp_landing_raw_1.py CDC_DAY1 csv 20191203
[2019-12-04 19:27:52,527] {bash_operator.py:124} INFO - Output:
[2019-12-04 19:27:54,937] {bash_operator.py:128} INFO - Exception in thread "main" org.apache.spark.SparkException: When running with master 'yarn-client' either HADOOP_CONF_DIR or YARN_CONF_DIR must be set in the environment.
[2019-12-04 19:27:54,938] {bash_operator.py:128} INFO - at org.apache.spark.deploy.SparkSubmitArguments.error(SparkSubmitArguments.scala:657)
[2019-12-04 19:27:54,938] {bash_operator.py:128} INFO - at org.apache.spark.deploy.SparkSubmitArguments.validateSubmitArguments(SparkSubmitArguments.scala:290)
[2019-12-04 19:27:54,938] {bash_operator.py:128} INFO - at org.apache.spark.deploy.SparkSubmitArguments.validateArguments(SparkSubmitArguments.scala:251)
[2019-12-04 19:27:54,938] {bash_operator.py:128} INFO - at org.apache.spark.deploy.SparkSubmitArguments.<init>(SparkSubmitArguments.scala:120)
[2019-12-04 19:27:54,939] {bash_operator.py:128} INFO - at org.apache.spark.deploy.SparkSubmit$$anon$2$$anon$1.<init>(SparkSubmit.scala:907)
[2019-12-04 19:27:54,939] {bash_operator.py:128} INFO - at org.apache.spark.deploy.SparkSubmit$$anon$2.parseArguments(SparkSubmit.scala:907)
[2019-12-04 19:27:54,939] {bash_operator.py:128} INFO - at org.apache.spark.deploy.SparkSubmit.doSubmit(SparkSubmit.scala:81)
[2019-12-04 19:27:54,939] {bash_operator.py:128} INFO - at org.apache.spark.deploy.SparkSubmit$$anon$2.doSubmit(SparkSubmit.scala:920)
[2019-12-04 19:27:54,940] {bash_operator.py:128} INFO - at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:929)
[2019-12-04 19:27:54,940] {bash_operator.py:128} INFO - at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)
[2019-12-04 19:27:54,973] {bash_operator.py:132} INFO - Command exited with return code 1
[2019-12-04 19:27:54,987] {taskinstance.py:1058} ERROR - Bash command failed
Traceback (most recent call last):
File "/usr/local/lib/python3.6/dist-packages/airflow/models/taskinstance.py", line 930, in _run_raw_task
result = task_copy.execute(context=context)
File "/usr/local/lib/python3.6/dist-packages/airflow/operators/bash_operator.py", line 136, in execute
raise AirflowException("Bash command failed")
airflow.exceptions.AirflowException: Bash command failed
[2019-12-04 19:27:54,990] {taskinstance.py:1081} INFO - Marking task as UP_FOR_RETRY
Traceback (most recent call last):
File "/usr/local/bin/airflow", line 37, in <module>
args.func(args)
File "/usr/local/lib/python3.6/dist-packages/airflow/utils/cli.py", line 74, in wrapper
return f(*args, **kwargs)
File "/usr/local/lib/python3.6/dist-packages/airflow/bin/cli.py", line 688, in test
ti.run(ignore_task_deps=True, ignore_ti_state=True, test_mode=True)
File "/usr/local/lib/python3.6/dist-packages/airflow/utils/db.py", line 74, in wrapper
return func(*args, **kwargs)
File "/usr/local/lib/python3.6/dist-packages/airflow/models/taskinstance.py", line 1020, in run
session=session)
File "/usr/local/lib/python3.6/dist-packages/airflow/utils/db.py", line 70, in wrapper
return func(*args, **kwargs)
File "/usr/local/lib/python3.6/dist-packages/airflow/models/taskinstance.py", line 930, in _run_raw_task
result = task_copy.execute(context=context)
File "/usr/local/lib/python3.6/dist-packages/airflow/operators/bash_operator.py", line 136, in execute
raise AirflowException("Bash command failed")
airflow.exceptions.AirflowException: Bash command failed
我是不是在配置或其他什么地方出错了?
【问题讨论】:
标签: apache-spark pyspark airflow directed-acyclic-graphs spark-submit