Tasks and workflows
The biggest structural change in Flyte 2 is that everything is a task. The @task, @workflow, and @dynamic decorators all collapse into a single @env.task on a flyte.TaskEnvironment, and a “workflow” is just a task that calls other tasks. This page covers the basic structure; see
Task configuration for the environment settings and
Migration for the big picture.
Hello world: tasks and workflows
A @task plus @workflow becomes two @env.tasks, where the entrypoint task calls the others. Sequential calls are naturally ordered — no >> operator required.
import flyte
env = flyte.TaskEnvironment(name="hello_world")
@env.task
def say_hello(name: str) -> str:
return f"Hello, {name}!"
@env.task
def to_upper(greeting: str) -> str:
return greeting.upper()
# The "workflow" is now just a task that calls other tasks.
@env.task
def main(name: str) -> str:
greeting = say_hello(name)
return to_upper(greeting)
Chaining and ordering
In Flyte 1 you sometimes used >> to force ordering between tasks with no data dependency. In Flyte 2, sequential (synchronous) calls run in the order they’re written, and awaiting async tasks in sequence does the same. The >> operator is gone.
from flytekit import task, workflow
@task
def clear_staging_table() -> None:
# Side effect only: truncate the staging table.
print("cleared staging table")
@task
def load_into_staging() -> None:
# Side effect only: load fresh rows into staging.
print("loaded staging table")
@task
def publish_to_prod() -> None:
# Side effect only: swap staging into the production table.
print("published to prod")
@workflow
def main() -> None:
clear = clear_staging_table()
load = load_into_staging()
publish = publish_to_prod()
# These tasks pass no data between them, so use the >> operator to force
# ordering: clear must finish before load, which must finish before publish.
clear >> load >> publish
import flyte
env = flyte.TaskEnvironment(name="staging_publish")
@env.task
def clear_staging_table() -> None:
print("cleared staging table")
@env.task
def load_into_staging() -> None:
print("loaded staging table")
@env.task
def publish_to_prod() -> None:
print("published to prod")
# Sequential (synchronous) calls run in the order they're written, even when no
# data flows between them. The Flyte 1 `>>` ordering operator is gone.
@env.task
def main() -> None:
clear_staging_table()
load_into_staging()
publish_to_prod()
Subworkflows
A @workflow invoked by another @workflow (for example, a reusable preprocessing pipeline) becomes a task that calls other tasks — nest them as deeply as you like.
from flytekit import task, workflow
@task
def impute(value: float) -> float:
# Replace missing/negative sentinel values with 0.
return value if value >= 0 else 0.0
@task
def scale(value: float) -> float:
return value / 100.0
@workflow
def preprocess(value: float) -> float:
imputed = impute(value=value)
return scale(value=imputed)
@workflow
def main(raw_value: float) -> float:
return preprocess(value=raw_value)
import flyte
env = flyte.TaskEnvironment(name="subworkflow")
@env.task
def impute(value: float) -> float:
# Replace missing/negative sentinel values with 0.
return value if value >= 0 else 0.0
@env.task
def scale(value: float) -> float:
return value / 100.0
# A preprocessing "subworkflow" is just a task that calls other tasks.
@env.task
def preprocess(value: float) -> float:
imputed = impute(value)
return scale(imputed)
@env.task
def main(raw_value: float) -> float:
return preprocess(raw_value)
TaskEnvironment configuration
The TaskEnvironment holds the configuration that Flyte 1 spread across the @task decorator. The task decorator can still override a few settings per-task.
import flyte
env = flyte.TaskEnvironment(
name="my_env", # Required: unique name
image=flyte.Image.from_debian_base(), # Or a string, or "auto"
resources=flyte.Resources(
cpu="2",
memory="4Gi",
gpu="A100:1",
disk="10Gi",
),
env_vars={"LOG_LEVEL": "INFO"},
secrets=[flyte.Secret(key="api-key", as_env_var="API_KEY")],
cache="auto", # "auto", "override", "disable", or a Cache object
reusable=flyte.ReusePolicy(replicas=5, idle_ttl=60),
interruptible=True,
)
# The task decorator can override some settings:
@env.task(
short_name="my_task", # Display name
cache="disable", # Override cache
retries=3, # Retry count
timeout=3600, # Seconds or a timedelta
report=True, # Generate an HTML report
)
def my_task(x: int) -> int:
return xParameter mapping: @task → TaskEnvironment + @env.task
Flyte 1 @task parameter |
Flyte 2 location | Notes |
|---|---|---|
container_image |
TaskEnvironment(image=...) |
Env-level only |
requests |
TaskEnvironment(resources=...) |
Env-level only |
limits |
TaskEnvironment(resources=...) |
Combined with requests (single value) |
environment |
TaskEnvironment(env_vars=...) |
Env-level only |
secret_requests |
TaskEnvironment(secrets=...) |
Env-level only |
cache |
Both | Can override at task level |
cache_version |
flyte.Cache(version_override=...) |
In a Cache object |
retries |
@env.task(retries=...) |
Task-level only |
timeout |
@env.task(timeout=...) |
Task-level only |
interruptible |
Both | Can override at task level |
pod_template |
Both | Can override at task level |
deprecated |
N/A | Not in Flyte 2 |
docs |
@env.task(docs=...) |
Task-level only |
For image, resource, secret, and caching detail, see Task configuration.
Next
- Task configuration — image, resources, caching, secrets, and scheduling
- Control flow — conditionals, dynamic behavior, and error handling
- Parallelism and fan-out — running tasks in parallel