Skip to main content

steps.parallel

The steps.parallel function runs a collection of functions or steps.task descriptors concurrently, waits for all of them, and returns their results as a list in input order.

Usage​

steps.parallel(functions = None, tasks = None, max_concurrency = 4, fail_fast = False)

Supply exactly one of functions or tasks.

Arguments​

functions

(Optional) A list or tuple of callables that take no arguments. Pass function references such as api, not calls such as api(). Each task is named after its function and position, for example api[0].

tasks

(Optional) A list or tuple of steps.task values. Use tasks to pass arguments and to attach a retry policy or a timeout. Task names must be unique within the group.

max_concurrency

(Optional) The most tasks that run at the same time. Defaults to 4. It must be a positive integer. Nested groups each have their own limit.

fail_fast

(Optional) Defaults to False. By default every independent task finishes and all failures are reported together. With True, a task that fails after exhausting its retries cancels the running tasks and skips the queued ones. The group still waits for the cancellations to finish.

Returns​

A list with one entry per task, in the order you supplied them, regardless of the order in which tasks finish. An empty collection returns an empty list. Each result is frozen, so it is read-only.

Behavior​

  • Everything reachable becomes read-only. When the group starts, Atmos freezes the functions, their arguments, their closures, and the script's global values. Build lists and dictionaries before the call, or create them inside each function and return them. Attempting to change a frozen value fails with cannot insert into frozen hash table or a similar message.
  • Output is attributed. With more than one task, each line a task prints, logs, or streams from a subprocess is prefixed with the task name, and lines from different tasks never interleave mid-line. A group with a single task inherits the surrounding prefix.
  • Failures are aggregated. One failing task is reported with its name, message, and traceback. Several failures are reported together as N tasks failed.
  • Nesting works. A task can start its own group. When the outer group has several tasks, the inner lines carry both names, for example [outer] [inner[0]] message.
  • Process settings stay local. Directories and environment overrides belong to each call and never change Atmos's own working directory or environment.
  • dependencies.tools must be called before the group starts, not inside a task.

Errors​

  • Passing both functions and tasks, or neither, fails with provide exactly one of functions or tasks.
  • An entry in functions that is not callable, an entry in tasks that is not a steps.task, or a duplicate task name fails before any task starts.
  • A task that exceeds its timeout fails with task "<name>" timed out after <duration>.

Examples​

Run functions concurrently​

def api():
return {"service": "api", "region": env["REGION"]}

def worker():
return {"service": "worker", "region": env["REGION"]}

output = steps.parallel(
functions = [api, worker],
max_concurrency = 2,
)

Pass arguments with tasks​

def deploy(name):
return exec.run(["./deploy.sh", name], output = "capture").stdout

output = steps.parallel(
tasks = [
steps.task(name = name, function = deploy, args = [name], timeout = "5m")
for name in ["api", "worker", "scheduler"]
],
max_concurrency = 2,
)

Stop early on the first failure​

def check(stack):
atmos.validate("component", "vpc", flags = {"stack": stack})

steps.parallel(
tasks = [steps.task(name = s, function = check, args = [s]) for s in ["dev", "staging", "prod"]],
fail_fast = True,
)

Prefixed output​

def work(n):
print("working", n)
ui.info("step " + str(n))

steps.parallel(tasks = [steps.task(name = "t%d" % i, function = work, args = [i]) for i in range(2)])
[t0] working 0
[t1] working 1
[t0] ▶ step 0
[t1] ▶ step 1

The order of lines from different tasks varies from run to run.