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.

Release files for flytekitplugins-awsemrserverless 1.16.28

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.28
File Size Uploaded
flytekitplugins_awsemrserverless-1.16.28.tar.gz 38.8 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for flytekitplugins-awsemrserverless 1.16.28
File Interpreter ABI Platform
flytekitplugins_awsemrserverless-1.16.28-py3-none-any.whl Python 3 none any Details

Total release size: 65.2 kB

Release files / flytekitplugins_awsemrserverless-1.16.28.tar.gz

Download URL flytekitplugins_awsemrserverless-1.16.28.tar.gz
Size 38.8 kB
Tags Source
SHA-256 checksum
How to use checksums
62500d517879c9d32a00e7a94d7fba3b8c766c0c684d98be7eb420d8f7f4d276
BLAKE2b-256 checksum
How to use checksums
1e2ffd538bf1ad7f0b9a95ea50c14de1d5d0f61a22e5dc5756c3563f55274589
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.28-py3-none-any.whl

Download URL flytekitplugins_awsemrserverless-1.16.28-py3-none-any.whl
Size 26.3 kB
Tags Python 3
SHA-256 checksum
How to use checksums
fbd26ee69e43b84a3d1302902189f39ff070b9a5ac5c34a21d945d24e8ad85fd
BLAKE2b-256 checksum
How to use checksums
3d4ec7bc6aef2c7bbba0131265124997a28cbe8bdb4431e0f5cad43ea50dacd8
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.28 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