
AWS CreateClusterAndSubmitSteps
CertifiedCreate an EMR cluster and run steps
AWS CreateClusterAndSubmitSteps
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.
type: io.kestra.plugin.aws.emr.CreateClusterAndSubmitStepsExamples
Create an EMR Cluster, submit a Spark job, wait until the job is terminated
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
clusterName *Requiredstring
Cluster name
Name assigned to the EMR cluster.
instanceCount *Requiredintegerstring
Instance count
Total number of instances in the cluster.
masterInstanceType *Requiredstring
Master instance type
EC2 instance type for the primary node.
slaveInstanceType *Requiredstring
Core/Task instance type
EC2 instance type for core/task nodes.
accessKeyId string
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.
applications array
Applications
EMR applications to install, e.g., Hive, Spark, Ganglia.
compatibilityMode booleanstring
Enable compatibility mode
Use it to connect to S3 bucket with S3 compatible services that don't support the new transport client.
completionCheckInterval string
PT10SCheck interval
Polling frequency while waiting; default 10s.
ec2KeyName string
EC2 key pair
Existing EC2 key pair for SSH access to the master as user hadoop.
ec2SubnetId string
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.
endpointOverride string
The endpoint with which the SDK should communicate
This property allows you to use a different S3 compatible storage backend.
forcePathStyle booleanstring
Force path style access
Must only be used when compatibilityMode is enabled.
jobFlowRole string
EMR_EC2_DefaultRoleJob flow role
EC2 instance profile for the cluster nodes; default EMR_EC2_DefaultRole and must pre-exist.
keepJobFlowAliveWhenNoSteps booleanstring
falseKeep cluster alive
If true, cluster stays in WAITING after steps; default false terminates when done.
logUri string
Log URI
S3 URI for cluster logs; leave empty to disable logging.
pluginDefaultsRef Non-dynamicstring
Reference (ref) of the pluginDefaults to apply to this task.
region string
AWS region with which the SDK should communicate
releaseLabel string
emr-5.20.0Release label
EMR release version such as emr-6.15.0; default emr-5.20.0.
secretKeyId string
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.
serviceRole string
EMR_DefaultRoleService role
IAM role assumed by EMR to access AWS resources; default EMR_DefaultRole.
sessionToken string
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 Non-dynamicarray
Steps to run
Optional steps submitted with RunJobFlow; executed in order.
io.kestra.plugin.aws.emr.models.StepConfig
TERMINATE_CLUSTERCANCEL_AND_WAITCONTINUETERMINATE_JOB_FLOWAction on failure
Behavior when the step fails: TERMINATE_CLUSTER, CANCEL_AND_WAIT, CONTINUE, or TERMINATE_JOB_FLOW.
JAR path
JAR executed for the step, e.g., command-runner.jar.
Step name
Label for the step, e.g., Run Spark job.
Arguments
List of arguments; each string is split on spaces before being passed to the step.
Main class
Entry class name; omit if the JAR manifest defines Main-Class.
stsEndpointOverride string
The AWS STS endpoint with which the SDKClient should communicate
stsRoleArn string
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.
stsRoleExternalId string
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.
stsRoleSessionDuration string
PT15MAWS 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.
stsRoleSessionName string
AWS STS Session name
This property is only used when an stsRoleArn is defined.
visibleToAllUsers booleanstring
trueVisible to all users
When true (default), any IAM principal in the account with permissions can manage the cluster.
wait booleanstring
falseWait for completion
When true, poll cluster status until TERMINATED, TERMINATED_WITH_ERRORS, or WAITING; default false.
waitUntilCompletion string
PT1HCompletion timeout
Maximum time to wait for WAITING or TERMINATED; default 1h.
Outputs
jobFlowId string
Job flow ID
Identifier of the created cluster.