【问题标题】:Airflow/Luigi for AWS EMR automatic cluster creation and pyspark deployment用于 AWS EMR 自动集群创建和 pyspark 部署的 Airflow/Luigi
【发布时间】:2019-04-16 11:56:23
【问题描述】:

我是气流自动化的新手,我现在不知道是否可以使用 apache 气流(或 luigi 等)来做到这一点,或者我应该只制作一个长的 bash 文件来做到这一点。

我想为此构建 dag

  1. 在 AWS EMR 上创建/克隆集群
  2. 安装 python 要求
  3. 安装pyspark相关库
  4. 从 github 获取最新代码
  5. 提交 Spark 作业
  6. 完成时终止集群

对于各个步骤,我可以制作如下所示的 .sh 文件(不确定这样做是否好),但不知道如何在气流中进行

1) 使用cluster.sh 创建一个集群

 aws emr create-cluster \
    --name "1-node dummy cluster" \
    --instance-type m3.xlarge \
    --release-label emr-4.1.0 \
    --instance-count 1 \
    --use-default-roles \
    --applications Name=Spark \
    --auto-terminate

2 & 3 & 4) 克隆 git 并安装要求 codesetup.sh

git clone some-repo.git
pip install -r requirements.txt
mv xyz.jar /usr/lib/spark/xyz.jar

5) 运行 spark 作业 sparkjob.sh

aws emr add-steps --cluster-id <Your EMR cluster id> --steps Type=spark,Name=TestJob,Args=[--deploy-mode,cluster,--master,yarn,--conf,spark.yarn.submit.waitAppCompletion=true,pythonjob.py,s3a://your-source-bucket/data/data.csv,s3a://your-destination-bucket/test-output/],ActionOnFailure=CONTINUE

6) 不确定,可能是这个

  terminate-clusters
--cluster-ids <value> [<value>...]

最后,这一切都可以作为一个 .sh 文件执行。我需要知道气流/luigi 的好方法。

我发现了什么:

我发现这篇文章很接近,但它已经过时(2016 年)并且错过了剧本的连接和代码

https://www.agari.com/email-security-blog/automated-model-building-emr-spark-airflow/

【问题讨论】:

    标签: amazon-web-services apache-spark pyspark airflow luigi


    【解决方案1】:

    我发现,可以有两种选择

    1) 我们可以在 emr create-clusteraddstep 的帮助下制作一个 bash 脚本,然后使用气流 Bashoperator 来调度它

    或者,这两个有一个包装器,称为sparksteps

    他们文档中的一个例子

    sparksteps examples/episodes.py \
      --s3-bucket $AWS_S3_BUCKET \
      --aws-region us-east-1 \
      --release-label emr-4.7.0 \
      --uploads examples/lib examples/episodes.avro \
      --submit-args="--deploy-mode client --jars /home/hadoop/lib/spark-avro_2.10-2.0.2-custom.jar" \
      --app-args="--input /home/hadoop/episodes.avro" \
      --tags Application="Spark Steps" \
      --debug
    

    您可以使用您选择的默认选项创建.sh script。准备好此脚本后,您可以从气流 bashoperator 调用它,如下所示

    create_command = "sparkstep_custom.sh "    
    
    t1 = BashOperator(
            task_id= 'create_file',
            bash_command=create_command,
            dag=dag
       )
    

    2) 您可以使用气流自己的操作符来执行此操作。

    EmrCreateJobFlowOperator(用于启动集群)EmrAddStepsOperator(用于提交spark作业) EmrStepSensor(跟踪步骤何时完成) EmrTerminateJobFlowOperator(在步骤完成时终止集群)

    创建集群和提交步骤的基本示例

    my_step=[
    
        {
            'Name': 'setup - copy files',
            'ActionOnFailure': 'CANCEL_AND_WAIT',
            'HadoopJarStep': {
                'Jar': 'command-runner.jar',
                'Args': ['aws', 's3', 'cp', S3_URI + 'test.py', '/home/hadoop/']
            }
        },
    {
            'Name': 'setup - copy files 3',
            'ActionOnFailure': 'CANCEL_AND_WAIT',
            'HadoopJarStep': {
                'Jar': 'command-runner.jar',
                'Args': ['aws', 's3', 'cp', S3_URI + 'myfiledependecy.py', '/home/hadoop/']
            }
        },
     {
            'Name': 'Run Spark',
            'ActionOnFailure': 'CANCEL_AND_WAIT',
            'HadoopJarStep': {
                'Jar': 'command-runner.jar',
                'Args': ['spark-submit','--jars', "jar1.jar,jar2.jar", '--py-files','/home/hadoop/myfiledependecy.py','/home/hadoop/test.py']
            }
        }
        ]
    
    
    cluster_creator = EmrCreateJobFlowOperator(
        task_id='create_job_flow2',
        job_flow_overrides=JOB_FLOW_OVERRIDES,
        aws_conn_id='aws_default',
        emr_conn_id='emr_default',
        dag=dag
    )
    
    step_adder_pre_step = EmrAddStepsOperator(
        task_id='pre_step',
        job_flow_id="{{ task_instance.xcom_pull('create_job_flow2', key='return_value') }}",
        aws_conn_id='aws_default',
        steps=my_steps,
        dag=dag
    )
    step_checker = EmrStepSensor(
        task_id='watch_step',
        job_flow_id="{{ task_instance.xcom_pull('create_job_flow2', key='return_value') }}",
        step_id="{{ task_instance.xcom_pull('pre_step', key='return_value')[0] }}",
        aws_conn_id='aws_default',
        dag=dag
    )
    
    cluster_remover = EmrTerminateJobFlowOperator(
        task_id='remove_cluster',
        job_flow_id="{{ task_instance.xcom_pull('create_job_flow2', key='return_value') }}",
        aws_conn_id='aws_default',
        dag=dag
    )
    

    另外,将代码上传到 s3(我很想从 github 获取最新代码_可以使用 s3boto3Pythonoperator 来完成)

    简单示例

    S3_BUCKET = 'you_bucket_name'
    S3_URI = 's3://{bucket}/'.format(bucket=S3_BUCKET)
    def upload_file_to_S3(filename, key, bucket_name):
        s3.Bucket(bucket_name).upload_file(filename, key)
    
    upload_to_S3_task = PythonOperator(
        task_id='upload_to_S3',
        python_callable=upload_file_to_S3,
        op_kwargs={
            'filename': configdata['project_path']+'test.py',
            'key': 'test.py',
            'bucket_name': 'dep-buck',
        },
        dag=dag)
    

    【讨论】:

      【解决方案2】:

      Airflow 有这方面的运营商。 airflow doc

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2015-03-21
        • 1970-01-01
        • 2020-01-23
        • 2020-01-24
        • 1970-01-01
        • 2020-12-28
        • 2016-02-19
        相关资源
        最近更新 更多