Skip to main content

Airflow connector to Ocean for Apache Spark

An Airflow plugin and provider to launch and monitor Spark applications on Ocean for Apache Spark.

Installation

pip install ocean-spark-airflow-provider

Usage

For general usage of Ocean for Apache Spark, refer to the official documentation.

Setting up the connection

In the connection menu, register a new connection of type Ocean for Apache Spark. The default connection name is ocean_spark_default. You will need to have:

  • The Ocean Spark cluster ID of the cluster you just created (of the format osc-e4089a00). You can find it in the Spot console in the list of clusters, or by using the Cluster List API.
  • A Spot token to interact with the Spot API.

connection setup dialog

The Ocean for Apache Spark connection type is not available for Airflow 1, instead create an HTTP connection and fill your cluster id as host, and your API token as password.

You will need to create a separate connection for each Ocean Spark cluster that you want to use with Airflow. In the OceanSparkOperator, you can select which Ocean Spark connection to use with the connection_name argument (defaults to ocean_spark_default). For example, you may choose to have one Ocean Spark cluster per environment (dev, staging, prod), and you can easily target an environment by picking the correct Airflow connection.

Using the Spark operator

from ocean_spark.operators import OceanSparkOperator

# DAG creation

spark_pi_task = OceanSparkOperator(
    job_id="spark-pi",
    task_id="compute-pi",
    dag=dag,
    config_overrides={
        "type": "Scala",
        "sparkVersion": "3.2.0",
        "image": "gcr.io/datamechanics/spark:platform-3.2-latest",
        "imagePullPolicy": "IfNotPresent",
        "mainClass": "org.apache.spark.examples.SparkPi",
        "mainApplicationFile": "local:///opt/spark/examples/jars/examples.jar",
        "arguments": ["10000"],
        "driver": {
            "cores": 1,
            "spot": false
        },
        "executor": {
            "cores": 4,
            "instances": 1,
            "spot": true,
            "instanceSelector": "r5"
        },
    },
)

Using the Spark Connect operator (available since airflow 2.6.2)

from airflow import DAG, utils
from ocean_spark.operators import (
    OceanSparkConnectOperator,
)

args = {
    "owner": "airflow",
    "depends_on_past": False,
    "start_date": utils.dates.days_ago(0, second=1),
}


dag = DAG(dag_id="spark-connect-task", default_args=args, schedule_interval=None)

spark_pi_task = OceanSparkConnectOperator(
    task_id="spark-connect",
    dag=dag,
)

Trigger the DAG with config, such as

{
  "sql": "select random()"
}

more examples are available for Airflow 2.

Test locally

You can test the plugin locally using the docker compose setup in this repository. Run make serve_airflow at the root of the repository to launch an instance of Airflow 2 with the provider already installed.

Metadata

Release files for ocean-spark-airflow-provider 1.1.4

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for ocean-spark-airflow-provider 1.1.4
File Size Uploaded
ocean_spark_airflow_provider-1.1.4.tar.gz 12.4 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for ocean-spark-airflow-provider 1.1.4
File Interpreter ABI Platform
ocean_spark_airflow_provider-1.1.4-py3-none-any.whl Python 3 none any Details

Total release size: 28.9 kB

Release files / ocean_spark_airflow_provider-1.1.4.tar.gz

Download URL ocean_spark_airflow_provider-1.1.4.tar.gz
Size 12.4 kB
Tags Source
SHA-256 checksum
How to use checksums
6db73eb59f77a4582e7a6933db444596420140596015b391ed16d6d147e4c863
BLAKE2b-256 checksum
How to use checksums
79d276989946167a053307dcd061fdff682aabec3182354be3df9d1df17a5ae8
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/5.1.1 CPython/3.12.6

Release files / ocean_spark_airflow_provider-1.1.4-py3-none-any.whl

Download URL ocean_spark_airflow_provider-1.1.4-py3-none-any.whl
Size 16.5 kB
Tags Python 3
SHA-256 checksum
How to use checksums
9dcd05da2cbd988c5c084e53da1fb9b201a71123167183d2c4c33f9bb48741d3
BLAKE2b-256 checksum
How to use checksums
079a2196110e185db76e3ece2a74875dfd37d2e9dbb96c919dd24edbb7978ce7
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/5.1.1 CPython/3.12.6

Release history Release notifications | RSS feed

This release

1.1.4 This release

2 release files

1.1.3

2 release files

1.1.2

2 release files

1.1.1

2 release files

1.1.0

2 release files

1.0.0

2 release files

0.1.6

2 release files

0.1.5

2 release files

0.1.4

2 release files

0.1.3

2 release files

0.1.2

2 release files

0.1.1

2 release files

0.1.0

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