Flytekit AWS EMR Serverless Plugin
A Flyte connector for AWS EMR Serverless that submits Spark and Hive jobs to an EMR Serverless application and tracks them through to completion.
Features
- Pythonic Spark mode: write a Flyte
@taskwhose body is regular PySpark; the plugin packages the user code, uploads it to S3, and runs it on EMR Serverless. No long-lived cluster to manage. - Script Spark mode: point at an existing
main.py(or JAR) already in S3 and submit it directly. - Hive mode: submit a Hive query (inline or from S3) against an EMR Serverless application configured for Hive.
- Async connector lifecycle (
create/get/delete) so the connector pod stays light and many jobs can be tracked concurrently. - Honours Flyte task retries, timeouts, and cancellation, and surfaces EMR Serverless logs through the Flyte UI when log URIs are available.
Installation
pip install flytekitplugins-awsemrserverless
The connector is registered automatically with flytekit via the plugin entry point. Deploy it on a flyteconnector pod that has this package installed and an IAM identity allowed to call EMR Serverless StartJobRun / GetJobRun / CancelJobRun and to read/write the script-staging S3 prefix.
Usage
Pythonic Spark task
from flytekit import task, workflow
from flytekitplugins.awsemrserverless import EMRServerless, EMRServerlessSparkJobDriver
@task(
task_config=EMRServerless(
application_id="00fhabc12345",
execution_role_arn="arn:aws:iam::123456789012:role/EMRServerlessRole",
region="us-east-1",
job_driver=EMRServerlessSparkJobDriver(
spark_submit_parameters="--conf spark.executor.cores=2 --conf spark.executor.memory=4g",
),
),
)
def spark_count() -> int:
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()
return spark.range(1_000_000).count()
@workflow
def wf() -> int:
return spark_count()
The plugin serializes the task body, uploads it to S3, and EMR Serverless runs it inside the worker image you have associated with the application.
Script Spark task
@task(
task_config=EMRServerless(
application_id="00fhabc12345",
execution_role_arn="arn:aws:iam::123456789012:role/EMRServerlessRole",
region="us-east-1",
job_driver=EMRServerlessSparkJobDriver(
entry_point="s3://my-bucket/scripts/main.py",
entry_point_arguments=["--date", "2025-01-01"],
spark_submit_parameters="--conf spark.executor.memory=4g",
),
),
)
def submit_script():
...
Hive task
from flytekitplugins.awsemrserverless import EMRServerless, EMRServerlessHiveJobDriver
@task(
task_config=EMRServerless(
application_id="00fhabc12345",
execution_role_arn="arn:aws:iam::123456789012:role/EMRServerlessRole",
region="us-east-1",
job_driver=EMRServerlessHiveJobDriver(
query="SELECT COUNT(*) FROM my_table",
),
),
)
def hive_query():
...
Worker image
For Pythonic Spark tasks the worker image must contain flytekit and this plugin so the executor can rehydrate the task object on the EMR Serverless side. Script Spark and Hive jobs do not require flytekit on the worker.
A reference Dockerfile is shipped alongside this plugin; it builds on the public EMR Serverless Spark base image and installs the matching flytekit and flytekitplugins-awsemrserverless versions:
docker build \
--build-arg VERSION=<flytekit-version> \
-t <registry>/emr-serverless-flytekit:<tag> \
plugins/flytekit-aws-emr-serverless
Override the base with --build-arg EMR_BASE_IMAGE=... to track a different EMR release.
IAM
The connector pod's IAM principal needs:
emr-serverless:StartJobRun,emr-serverless:GetJobRun,emr-serverless:CancelJobRunon the target applications3:GetObject/s3:PutObjecton the script-staging prefix (Pythonic mode)iam:PassRolefor the EMR Serverless execution role
The execution role attached to the EMR Serverless application is the role the workers run as and needs whatever data-access permissions your jobs require.
Discussion
Tracking issue: flyteorg/flyte#7286.
Metadata
Release files for flytekitplugins-awsemrserverless 1.16.27
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| flytekitplugins_awsemrserverless-1.16.27.tar.gz | 38.8 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| flytekitplugins_awsemrserverless-1.16.27-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 65.2 kB
Release files / flytekitplugins_awsemrserverless-1.16.27.tar.gz
| Download URL | flytekitplugins_awsemrserverless-1.16.27.tar.gz |
|---|---|
| Size | 38.8 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
dd8e3ff172d87bb86c5d013c39482114756641b509d0d91e1b7e28a2afdba385
|
|
BLAKE2b-256 checksum How to use checksums |
f5f384cf2a9c396c7a2e5402401807f251b3bad8a8795873a731a244af4af926
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.12.13
|
Release files / flytekitplugins_awsemrserverless-1.16.27-py3-none-any.whl
| Download URL | flytekitplugins_awsemrserverless-1.16.27-py3-none-any.whl |
|---|---|
| Size | 26.3 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
71bc10dc5f7f89392d87f5325d6b6a994f183f38ad426400c936acd04c0ffe88
|
|
BLAKE2b-256 checksum How to use checksums |
714848d68e6090594f07c88a123a85a5f005291561795f44f968ee798d3f762b
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.12.13
|