Skip to main content

pacer

About

pacer is a lightweight Python package for implementing distributed data processing workflows. Instead of defining a DAG which models the data flow from sources to a final result pacer uses a pull model which is very similar to nesting function calls. Running such a workflow starts on the result node and recursively delegates work to the inputs.

Originally we developed pacer for running analysis pipelines in emzed, a framework for analyzing LCMS data.

How does pacer work ?

Under the hood pacer has two core components:

  • one for managing distributed computations of chained computations:

    Processing steps in pacer are just Python functions with some additional annotations. pacer tries to compute as many processing steps as possible in parallel, either because such a function has to be applied to different data sets, or it has more than one input and those are computed concurrently.

  • a distributed cache which is retained on the file system

    In case of partial modifications of the inputs a pacer workflow does not determine needed update computations but uses a distributed cache for mapping the input values of single processing steps to their final result. So a repeated run of the workflow with unchanged inputs will run the full workflow with all processing steps returning already known results immediately. Running the workflow with unknown or modified inputs will only execute the needed computations and update the cache.

These two components are independent and can be used seperately.

Examples

We provide some simple examples which show how easy it is to use pacer. You find these examples which we extended to print more logging information in the examples/ folder in the git repository.

In a real world LCMS workflow we would not use as simple functions as used below but longer running computation steps such as running a LCMS peak picker and a subsequent peak aligner.

How to declare a pipeline

In this case our input sources are a list of Python strings ["a", "bc", "def"] and a tuple of numbers (1, 2). The very simple example workflow computes the length of each string and multiplies it with every number from the tuple. This very simple example could be implemented in pure Python as follows:

import itertools

def main():

    def length(what):
        return len(what)

    def multiply(a, b):
        return a * b

    words = ["a", "bc", "def"]
    multipliers = (1, 2)

    result = [multiply(length(w), v) for (w, v) in itertools.product(words, multipliers)]
    assert result == [1, 2, 2, 4, 3, 6]

if __name__ == "__main__":
    main()

In order to transform this computations to a smart parallel processing pipeline we use the apply and join function decorators from pacer and declare the dependencies among the single steps using function calls.

from pacer import apply, output, join, Engine

def main():

    @apply
    def length(what):
        return len(what)

    @output
    @join
    def multiply(a, b):
        return a * b

    words = ["a", "bc", "def"]
    multipliers = (1, 2)

    # now we DECLARE the workflow (no execution at that time):
    workflow = multiply(length(words), multipliers)

if __name__ == "__main__":
    main()

Running this workflow on three CPU cores is easy now. In this case the computation steps are run in parallel:

Engine.set_number_of_processes(3)
workflow.start_computations()
result = workflow.get_all_in_order()

assert result == [1, 2, 2, 4, 3, 6]

pacers approach to compute needed updates in case of modified input data

As already stated above, pacer does not determine needed update computations in case of modified input data but uses a distributed cache instead. So running a workflow a second time will fetch the already known results of computations not affected by changes, and start computations with unknown input arguments.

We use decorators again. Leveraging the example above only needs few adjustments:

from pacer import apply, join, output, Engine, CacheBuilder

cache = CacheBuiler("/tmp/cache_000")

@apply
@cache
def length(what):
    return len(what)

@output
@join
@cache
def multiply(a, b):
    return a * b

# inputs to workflow
words = ["a", "bc", "def"]
multipliers = (1, 2)

workflow = multiply(length(words), multipliers)

# run workflow
Engine.set_number_of_processes(3)
workflow.start_computations()
result = workflow.get_all_in_order()

assert result == [1, 2, 2, 4, 3, 6]

If you run these examples from a command line you see logging results showing the parallel execution of single steps and cache hits avoiding recomputations.

Metadata

Release files for pacer 0.30.6

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

Source distribution (sdist)

Source distribution for pacer 0.30.6
File Size Uploaded
pacer-0.30.6.tar.gz 18.7 kB Details

Release files / pacer-0.30.6.tar.gz

Download URL pacer-0.30.6.tar.gz
Size 18.7 kB
Tags Source
SHA-256 checksum
How to use checksums
8ed4d8ca3e281a3fc0e02b31dd4bc91fb4137cdfbea48fb278c5128a26112faf
BLAKE2b-256 checksum
How to use checksums
00c8a256d5f2dc048f01bdd39812ba6d032d8a1121e1a2860d585a91979994ee
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No

Release history Release notifications | RSS feed

This release

0.30.6 This release

1 release file

0.30.5

1 release file

0.30.4

1 release file

0.30.3

1 release file

0.30.2

1 release file

0.30.1

1 release file

0.30.0

1 release file

0.29.1

1 release file

0.29.0

1 release file

0.28.2

1 release file

0.28.1

1 release file

0.28.0

1 release file

0.27.0

1 release file

0.23.0

1 release file

0.22.1

1 release file

0.22.0

1 release file

0.21.7

1 release file

0.21.6

1 release file

0.21.5

1 release file

0.21.4

1 release file

0.21.3

1 release file

0.21.2

1 release file

0.21.1

1 release file

0.21.0

1 release file

0.20.1

1 release file

0.20.0

1 release file

0.19.4

1 release file

0.19.3

1 release file

0.19.2

1 release file

0.19.1

1 release file

0.18.0

1 release file

0.17.5

1 release file

0.17.4

1 release file

0.17.3

1 release file

0.17.2

1 release file

0.17.1

1 release file

0.17.0

1 release file

0.16.0

1 release file

0.15.8

1 release file

0.15.7

1 release file

0.15.6

1 release file

0.15.5

1 release file

0.15.4

1 release file

0.15.3

1 release file

0.15.2

1 release file

0.15.1

1 release file

0.15.0

1 release file

0.13.0

1 release file

0.12.3

1 release file

0.12.2

1 release file

0.12.1

1 release file

0.12.0

1 release file

0.11.1

1 release file

0.11.0

1 release file

0.10.0

1 release file

0.9.0

1 release file

0.8.0

1 release file

0.7.1

1 release file

0.7.0

1 release file

0.6.0

1 release file

0.5.10

1 release file

0.5.9

1 release file

0.5.8

1 release file

0.5.7

1 release file

0.5.5

1 release file

0.5.4

1 release file

0.5.3

1 release file

0.5.2

1 release file

0.5.1

1 release file

0.5.0

1 release file

0.4.7

1 release file

0.4.6

1 release file

0.4.5

1 release file

0.4.4

1 release file

0.4.3

1 release file

0.4.2

1 release file

0.4.1

1 release file

0.4.0

1 release file

0.3.1

1 release file

0.3.0

1 release file

0.2.2

1 release file

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