Skip to main content

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 @task whose 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:CancelJobRun on the target application
  • s3:GetObject / s3:PutObject on the script-staging prefix (Pythonic mode)
  • iam:PassRole for 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)

Source distribution for flytekitplugins-awsemrserverless 1.16.27
File Size Uploaded
flytekitplugins_awsemrserverless-1.16.27.tar.gz 38.8 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for flytekitplugins-awsemrserverless 1.16.27
File Interpreter ABI Platform
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

Release history Release notifications | RSS feed

This release

1.16.27 This release

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page