AWS CreateClusterAndSubmitSteps

AWS CreateClusterAndSubmitSteps

Certified

Create an EMR cluster and run steps

Launches a new EMR cluster, optionally submits initial steps, and returns the jobFlowId. Defaults: releaseLabel emr-5.20.0, serviceRole EMR_DefaultRole, jobFlowRole EMR_EC2_DefaultRole, visibleToAllUsers true, keepJobFlowAliveWhenNoSteps false. Can wait until the cluster reaches WAITING or TERMINATED, polling every completionCheckInterval until waitUntilCompletion.

yaml
type: io.kestra.plugin.aws.emr.CreateClusterAndSubmitSteps

Create an EMR Cluster, submit a Spark job, wait until the job is terminated

yaml
id: aws_emr_create_cluster
namespace: company.team

tasks:
  - id: create_cluster
    type: io.kestra.plugin.aws.emr.CreateClusterAndSubmitSteps
    accessKeyId: "{{ secret('AWS_ACCESS_KEY_ID') }}"
    secretKeyId: "{{ secret('AWS_SECRET_KEY_ID') }}"
    region: eu-west-3
    clusterName: "Spark job cluster"
    logUri: "s3://my-bucket/test-emr-logs"
    keepJobFlowAliveWhenNoSteps: true
    applications:
        - Spark
    masterInstanceType: m5.xlarge
    slaveInstanceType: m5.xlarge
    instanceCount: 3
    ec2KeyName: my-ec2-ssh-key-pair-name
    steps:
        - name: Spark_job_test
          jar: "command-runner.jar"
          actionOnFailure: CONTINUE
          commands:
            - spark-submit s3://mybucket/health_violations.py --data_source s3://mybucket/food_establishment_data.csv --output_uri s3://mybucket/test-emr-output
    wait: true
Properties

Cluster name

Name assigned to the EMR cluster.

Instance count

Total number of instances in the cluster.

Master instance type

EC2 instance type for the primary node.

Core/Task instance type

EC2 instance type for core/task nodes.

Access Key Id in order to connect to AWS

If no credentials are defined, we will use the default credentials provider chain to fetch credentials.

SubTypestring

Applications

EMR applications to install, e.g., Hive, Spark, Ganglia.

Enable compatibility mode

Use it to connect to S3 bucket with S3 compatible services that don't support the new transport client.

DefaultPT10S

Check interval

Polling frequency while waiting; default 10s.

EC2 key pair

Existing EC2 key pair for SSH access to the master as user hadoop.

EC2 subnet ID

Applies to clusters that use the uniform instance group configuration. To launch the cluster in Amazon Virtual Private Cloud (Amazon VPC), set this parameter to the identifier of the Amazon VPC subnet where you want the cluster to launch. If you do not specify this value and your account supports EC2-Classic, the cluster launches in EC2-Classic.

The endpoint with which the SDK should communicate

This property allows you to use a different S3 compatible storage backend.

Force path style access

Must only be used when compatibilityMode is enabled.

DefaultEMR_EC2_DefaultRole

Job flow role

EC2 instance profile for the cluster nodes; default EMR_EC2_DefaultRole and must pre-exist.

Defaultfalse

Keep cluster alive

If true, cluster stays in WAITING after steps; default false terminates when done.

Log URI

S3 URI for cluster logs; leave empty to disable logging.

Reference (ref) of the pluginDefaults to apply to this task.

AWS region with which the SDK should communicate

Defaultemr-5.20.0

Release label

EMR release version such as emr-6.15.0; default emr-5.20.0.

Secret Key Id in order to connect to AWS

If no credentials are defined, we will use the default credentials provider chain to fetch credentials.

DefaultEMR_DefaultRole

Service role

IAM role assumed by EMR to access AWS resources; default EMR_DefaultRole.

AWS session token, retrieved from an AWS token service, used for authenticating that this user has received temporary permissions to access a given resource

If no credentials are defined, we will use the default credentials provider chain to fetch credentials.

Steps to run

Optional steps submitted with RunJobFlow; executed in order.

Definitions
actionOnFailure*Requiredstring
Possible Values
TERMINATE_CLUSTERCANCEL_AND_WAITCONTINUETERMINATE_JOB_FLOW

Action on failure

Behavior when the step fails: TERMINATE_CLUSTER, CANCEL_AND_WAIT, CONTINUE, or TERMINATE_JOB_FLOW.

jar*Requiredstring

JAR path

JAR executed for the step, e.g., command-runner.jar.

name*Requiredstring

Step name

Label for the step, e.g., Run Spark job.

commandsarray
SubTypestring

Arguments

List of arguments; each string is split on spaces before being passed to the step.

mainClassstring

Main class

Entry class name; omit if the JAR manifest defines Main-Class.

The AWS STS endpoint with which the SDKClient should communicate

AWS STS Role

The Amazon Resource Name (ARN) of the role to assume. If set the task will use the StsAssumeRoleCredentialsProvider. If no credentials are defined, we will use the default credentials provider chain to fetch credentials.

AWS STS External Id

A unique identifier that might be required when you assume a role in another account. This property is only used when an stsRoleArn is defined.

DefaultPT15M

AWS STS Session duration

The duration of the role session (default: 15 minutes, i.e., PT15M). This property is only used when an stsRoleArn is defined.

AWS STS Session name

This property is only used when an stsRoleArn is defined.

Defaulttrue

Visible to all users

When true (default), any IAM principal in the account with permissions can manage the cluster.

Defaultfalse

Wait for completion

When true, poll cluster status until TERMINATED, TERMINATED_WITH_ERRORS, or WAITING; default false.

DefaultPT1H

Completion timeout

Maximum time to wait for WAITING or TERMINATED; default 1h.

Job flow ID

Identifier of the created cluster.