Skip to main content

Distributed Runner

This library allows for the easy construction and management of ephemeral Dask clusters on AWS via a simple context manager. This allows you to perform distributed and parallelized computations with ease.

license activity Code style: black

Dask Git GitLab Linux Python

Installation

pip install distrunner

Usage

In your scheduler (Airfow etc.) use something like this:

from distrunner import DistRunner
import whatismyip

def main():
    with DistRunner(workers=10) as cldr:
        results = cldr.client.map(print_ip, range(10))
        outcome = cldr.client.gather(results)

        print(outcome)

def print_ip(x):
    return f"My IP address is {whatismyip.whatismyip()}"

The "local" flag will determine whether a remote cluster is created. For example, the following will all run locally instead of spinning up infrastructure:

from distrunner import DistRunner
import whatismyip

def main():
    with DistRunner(workers=10, local=True) as cldr:
        results = cldr.client.map(print_ip, range(10))
        outcome = cldr.client.gather(results)

        print(outcome)

def print_ip(x):
    return f"My IP address is {whatismyip.whatismyip()}"

You will need to set the AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY environment variables to use the Fargate clusters (or run from an environment with an authorised IAM role).

Running in Managed Workflows for Apache Airflow (MWAA)

Running inside an AWS MWAA environment requires a little more setup than running locally. This is because of the way that requirements are handled in Airflow. You also need to ensure that any called functions are nested underneath the task in question. The above examples, refactored for Airflow, could read:

REQUIREMENTS = [
    "distrunner>=1.3.0",
    "whatismyip,
    "coiled",
    "dask[complete]",
]

@dag(
    default_args=DEFAULT_ARGS,
    schedule_interval="@daily",
    catchup=False,
    dagrun_timeout=timedelta(hours=16),
    start_date=datetime(2023, 4, 16),
    tags=["api"],
)
def main_task():
    @task.virtualenv(
        task_id="main_task",
        requirements=REQUIREMENTS,
        system_site_packages=True,
    )
    def entry_point(requirements):
        
        def print_ip(x):
            return f"My IP address is {whatismyip.whatismyip()}"

        from distrunner.distrunner import DistRunner

        with DistRunner(
            workers=1,
            requirements=requirements,
            application_name="test_application",
        ) as cldr:
            results = cldr.client.map(snapshot_routes_body, range(1))
            cldr.client.gather(results)

    entry_point(REQUIREMENTS)

main_task()

Features

  • Context manager handling of Dask Fargate clusters with scale-to-zero on complete
  • Easy ability to switch between local and distributed/remote development

What it Does

This library allows you to easily run functions across a Dask cluster.

Credits

Copyright © Crossref 2023

Metadata

Release files for distrunner 1.4.1

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

Source distribution (sdist)

Source distribution for distrunner 1.4.1
File Size Uploaded
distrunner-1.4.1.tar.gz 7.6 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for distrunner 1.4.1
File Interpreter ABI Platform
distrunner-1.4.1-py3-none-any.whl Python 3 none any Details

Total release size: 14.6 kB

Release files / distrunner-1.4.1.tar.gz

Download URL distrunner-1.4.1.tar.gz
Size 7.6 kB
Tags Source
SHA-256 checksum
How to use checksums
330b5566a6fb870f3c5f2d3722f0a48bba6a4aa9741d46bab5d72c5efbf529bc
BLAKE2b-256 checksum
How to use checksums
a49e6c2ee4fc2f2a346bef29bb9b84e22ebf0e6e6faa7143ff5097ba30fdd386
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/4.0.2 CPython/3.10.12

Release files / distrunner-1.4.1-py3-none-any.whl

Download URL distrunner-1.4.1-py3-none-any.whl
Size 7.0 kB
Tags Python 3
SHA-256 checksum
How to use checksums
144a6d6cd8a22e5fd0f491d7e7376af5e2b055e3da899dffd4951d8bb6dbc234
BLAKE2b-256 checksum
How to use checksums
386774cdc3e78895a77ecfc2dfb168e198ae9cbf99dd67bd7279d3c55f8f330a
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/4.0.2 CPython/3.10.12

Release history Release notifications | RSS feed

This release

1.4.1 This release

2 release files

1.4.0

2 release files

1.3.0

2 release files

1.2.0

2 release files

1.1.0

2 release files

1.0.0

2 release files

0.0.46

2 release files

0.0.45

2 release files

0.0.44

2 release files

0.0.43

2 release files

0.0.42

2 release files

0.0.41

2 release files

0.0.40

2 release files

0.0.39

2 release files

0.0.38

2 release files

0.0.37

2 release files

0.0.36

2 release files

0.0.35

2 release files

0.0.34

2 release files

0.0.33

2 release files

0.0.32

2 release files

0.0.31

2 release files

0.0.30

2 release files

0.0.29

2 release files

0.0.28

2 release files

0.0.27

2 release files

0.0.26

2 release files

0.0.25

2 release files

0.0.24

2 release files

0.0.23

2 release files

0.0.22

2 release files

0.0.21

2 release files

0.0.20

2 release files

0.0.19

2 release files

0.0.18

2 release files

0.0.17

2 release files

0.0.16

2 release files

0.0.9

2 release files

0.0.8

2 release files

0.0.7

2 release files

0.0.6

2 release files

0.0.5

2 release files

0.0.4

2 release files

0.0.3

2 release files

0.0.2

2 release files

0.0.1

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