Skip to content

Flows

Maturity labels

  • Now: Stable and supported in current releases.
  • Preview: Usable today, but behavior and APIs may evolve.
  • Planned: Not yet implemented.

Note

Status: Now

Linear flows

>>> import omnipy as om
>>> @om.TaskTemplate()
... def plus_one(number: int) -> int:
...     return number + 1
>>> @om.LinearFlowTemplate(plus_one, plus_one)
... def plus_two(number: int) -> int:
...     ...
>>> plus_two.run(10)

DAG flows

>>> import omnipy as om
>>> @om.TaskTemplate()
... def make_inputs() -> dict[str, int]:
...     return {'number': 21}
>>> @om.TaskTemplate()
... def double(number: int) -> int:
...     return number * 2
>>> @om.DagFlowTemplate(make_inputs, double)
... def my_dag() -> int:
...     ...
>>> my_dag.run()

FuncFlow

>>> import omnipy as om
>>> @om.TaskTemplate()
... def plus_one(number: int) -> int:
...     return number + 1
>>> @om.FuncFlowTemplate()
... def repeat_plus_one(number: int, n: int) -> int:
...     for _ in range(n):
...         number = plus_one(number=number)
...     return number
>>> repeat_plus_one.run(0, 3)

Generator flow bodies may exist only to define the callable signature while child jobs perform the actual execution.

@LinearFlowTemplate(task_one, generator_task_two)
def my_sync_generator_flow(number: int) -> Generator[int, None, None]:
    yield from Void()  # For generator signature only; never run.
@DagFlowTemplate(task_one, generator_task_two)
async def my_async_generator_flow(number: int) -> AsyncGenerator[int, None]:
    async for _ in Void():  # For generator signature only; never run.
        yield _