【问题标题】:Apache Spark and Apache Airflow connection in Docker based solution基于 Docker 的解决方案中的 Apache Spark 和 Apache Airflow 连接
【发布时间】:2018-09-20 06:12:59
【问题描述】:

我有 Spark 和 Airflow 集群,我想将 Spark 作业从 Airflow 容器发送到 Spark 容器。但是我是 Airflow 的新手,我不知道我需要执行哪种配置。我复制了插件文件夹下的spark_submit_operator.py

from airflow import DAG

from airflow.contrib.operators.spark_submit_operator import SparkSubmitOperator
from datetime import datetime, timedelta

    args = {
        'owner': 'airflow',
        'start_date': datetime(2018, 7, 31)
    }
    dag = DAG('spark_example_new', default_args=args, schedule_interval="*/10 * * * *")

    operator = SparkSubmitOperator(
        task_id='spark_submit_job',
        conn_id='spark_default',
        java_class='Simple',
        application='/spark/abc.jar',
        total_executor_cores='1',
        executor_cores='1',
        executor_memory='2g',
        num_executors='1',
        name='airflow-spark-example',
        verbose=False,
        driver_memory='1g',
        application_args=["1000"],
        conf={'master':'spark://master:7077'},
        dag=dag,
    )

ma​​ster 是我们的 Spark Master 容器的主机名。当我运行 dag 时,它会产生以下错误:

[2018-09-20 05:57:46,637] {{models.py:1569}} INFO - Executing <Task(SparkSubmitOperator): spark_submit_job> on 2018-09-20T05:57:36.756154+00:00
[2018-09-20 05:57:46,637] {{base_task_runner.py:124}} INFO - Running: ['bash', '-c', 'airflow run spark_example_new spark_submit_job 2018-09-20T05:57:36.756154+00:00 --job_id 4 --raw -sd DAGS_FOLDER/firstJob.py --cfg_path /tmp/tmpn2hznb5_']
[2018-09-20 05:57:47,002] {{base_task_runner.py:107}} INFO - Job 4: Subtask spark_submit_job [2018-09-20 05:57:47,001] {{settings.py:174}} INFO - setting.configure_orm(): Using pool settings. pool_size=5, pool_recycle=1800
[2018-09-20 05:57:47,312] {{base_task_runner.py:107}} INFO - Job 4: Subtask spark_submit_job [2018-09-20 05:57:47,311] {{__init__.py:51}} INFO - Using executor CeleryExecutor
[2018-09-20 05:57:47,428] {{base_task_runner.py:107}} INFO - Job 4: Subtask spark_submit_job [2018-09-20 05:57:47,428] {{models.py:258}} INFO - Filling up the DagBag from /usr/local/airflow/dags/firstJob.py
[2018-09-20 05:57:47,447] {{base_task_runner.py:107}} INFO - Job 4: Subtask spark_submit_job [2018-09-20 05:57:47,447] {{cli.py:492}} INFO - Running <TaskInstance: spark_example_new.spark_submit_job 2018-09-20T05:57:36.756154+00:00 [running]> on host e6dd59dc595f
[2018-09-20 05:57:47,471] {{logging_mixin.py:95}} INFO - [2018-09-20 05:57:47,470] {{spark_submit_hook.py:283}} INFO - Spark-Submit cmd: ['spark-submit', '--master', 'yarn', '--conf', 'master=spark://master:7077', '--num-executors', '1', '--total-executor-cores', '1', '--executor-cores', '1', '--executor-memory', '2g', '--driver-memory', '1g', '--name', 'airflow-spark-example', '--class', 'Simple', '/spark/ugur.jar', '1000']

[2018-09-20 05:57:47,473] {{models.py:1736}} ERROR - [Errno 2] No such file or directory: 'spark-submit': 'spark-submit'
Traceback (most recent call last):
  File "/usr/local/lib/python3.6/site-packages/airflow/models.py", line 1633, in _run_raw_task
    result = task_copy.execute(context=context)
  File "/usr/local/lib/python3.6/site-packages/airflow/contrib/operators/spark_submit_operator.py", line 168, in execute
    self._hook.submit(self._application)
  File "/usr/local/lib/python3.6/site-packages/airflow/contrib/hooks/spark_submit_hook.py", line 330, in submit
    **kwargs)
  File "/usr/local/lib/python3.6/subprocess.py", line 709, in __init__
    restore_signals, start_new_session)
  File "/usr/local/lib/python3.6/subprocess.py", line 1344, in _execute_child
    raise child_exception_type(errno_num, err_msg, err_filename)
FileNotFoundError: [Errno 2] No such file or directory: 'spark-submit': 'spark-submit'

正在运行的命令:

Spark-Submit cmd: ['spark-submit', '--master', 'yarn', '--conf', 'master=spark://master:7077', '--num-executors', '1', '--total-executor-cores', '1', '--executor-cores', '1', '--executor-memory', '2g', '--driver-memory', '1g', '--name', 'airflow-spark-example', '--class', 'Simple', '/spark/ugur.jar', '1000']

但我没有使用纱线。

【问题讨论】:

    标签: docker apache-spark airflow


    【解决方案1】:

    如果您使用 SparkSubmitOperator,则默认情况下与 master 的连接将设置为“Yarn”,无论您在 python 代码中设置哪个 master,但是,您可以通过其构造函数指定 conn_id 并在您已经创建的条件下覆盖 master上述conn_id 在 Airflow Web 界面的“Admin->Connection”菜单中。希望对你有帮助。

    【讨论】:

      【解决方案2】:

      我猜你需要在你的连接中为这个 spark conn id (spark_default) 设置额外的选项。默认情况下是纱线,所以试试 {"master":"your-con"} 另外,您的气流用户在类路径中有 spark-submit 吗?从日志中看起来没有。

      【讨论】:

      • 在我的集群中,spark 在独立于气流容器的容器中运行。因此,气流容器不包含 spark-submit。如何将在另一个容器中找到的 spark-submit 添加到气流容器的类路径中?
      • 您已经撤消了对帖子的令人信服的改进编辑。请向@AndrzejSydor 解释您的原因。
      • 哪一个? @Yunnosch
      • @AndrzejSydor 很抱歉让您感到困惑。我的评论针对 OP,只提到了你。我使用了 ping 语法以便通知您。我希望 OP 向您解释为什么他们取消了您的体面编辑,当然希望他们自己改进答案。
      猜你喜欢
      • 2019-11-11
      • 2011-08-17
      • 2020-09-23
      • 2018-12-16
      • 1970-01-01
      • 2015-10-05
      相关资源
      最近更新 更多