1. What is Flyte?1. Что такое Flyte?
Flyte is an open-source workflow orchestrator for data and machine-learning pipelines. You write your steps as plain Python functions, decorate them, and Flyte runs them as a versioned DAG (directed acyclic graph) on top of Kubernetes — each step in its own container, with resources, retries and caching handled for you.
Flyte — это опенсорсный оркестратор воркфлоу для data- и ML-пайплайнов. Ты пишешь шаги как обычные Python-функции, навешиваешь декораторы, и Flyte запускает их как версионируемый DAG (направленный ацикличный граф) поверх Kubernetes — каждый шаг в своём контейнере, с ресурсами, ретраями и кэшированием «из коробки».
2. Why does it exist?2. Зачем он нужен?
Classic orchestrators (like early Airflow) struggle with ML/data work in a few ways. Flyte was built (originally at Lyft) to fix them:
Классические оркестраторы (как ранний Airflow) плохо тянут ML/data-задачи по нескольким причинам. Flyte (изначально в Lyft) сделан, чтобы их закрыть:
- No type safety. Passing the wrong shape between steps blows up at runtime, hours in. Flyte checks types before running.
- No real caching. Re-running recomputes everything. Flyte caches outputs by inputs + version and skips unchanged steps.
- Weak versioning. Hard to reproduce "the run from last Tuesday". Flyte versions every workflow and task.
- Resource pain. GPUs, memory, parallel fan-out are awkward. Flyte makes per-task resources and dynamic fan-out first-class.
- Нет проверки типов. Передал не ту структуру между шагами — падает в рантайме, спустя часы. Flyte проверяет типы до запуска.
- Нет настоящего кэша. Повторный запуск всё пересчитывает. Flyte кэширует выходы по входам + версии и пропускает неизменившиеся шаги.
- Слабое версионирование. Трудно воспроизвести «тот запуск с прошлого вторника». Flyte версионирует каждый воркфлоу и таску.
- Боль с ресурсами. GPU, память, параллельный fan-out неудобны. Во Flyte ресурсы на таску и динамический fan-out — первоклассные граждане.
3. Core concepts3. Ключевые понятия
① Task① Task (таска)
The smallest unit of work: one Python function decorated with @task. It runs in its own container, declares its inputs/outputs by type, and can request CPU/GPU/memory.
Минимальная единица работы: одна Python-функция с декоратором @task. Бежит в своём контейнере, объявляет входы/выходы по типам, может запросить CPU/GPU/память.
② Workflow② Workflow (воркфлоу)
A function decorated with @workflow that wires tasks together. Flyte reads the data dependencies between tasks to build the DAG and decide what can run in parallel.
Функция с декоратором @workflow, которая связывает таски. Flyte читает зависимости по данным между тасками, строит DAG и решает, что можно запустить параллельно.
③ Caching & versioning③ Кэширование и версионирование
With cache=True + a cache_version, Flyte stores outputs keyed by (inputs, version). Re-run with the same inputs → it returns the cached result instantly instead of recomputing.
С cache=True + cache_version Flyte сохраняет выходы по ключу (входы, версия). Повторный запуск с теми же входами → мгновенно возвращает сохранённый результат вместо пересчёта.
④ Dynamic workflows & map tasks④ Динамические воркфлоу и map-таски
When the shape of the DAG depends on data (e.g. "process N files, N unknown until runtime"), @dynamic builds the graph at runtime, and map_task fans the same task out over a list in parallel.
Когда форма DAG зависит от данных (например «обработать N файлов, N неизвестно до запуска»), @dynamic строит граф в рантайме, а map_task распараллеливает одну таску по списку.
⑤ Launch plans & schedules⑤ Launch plans и расписания
A launch plan binds a workflow to specific inputs and an optional schedule (cron) — how you productionize and trigger a workflow.
A launch plan привязывает воркфлоу к конкретным входам и опциональному расписанию (cron) — так воркфлоу выводят в прод и запускают.
4. A minimal example4. Минимальный пример
A two-step pipeline: clean the data (cached), then average it. Notice the types on every edge.
Пайплайн из двух шагов: почистить данные (с кэшем), затем усреднить. Обрати внимание на типы на каждом ребре.
from flytekit import task, workflow # cached by (inputs, cache_version): same input → no recompute @task(cache=True, cache_version="1.0", requests=Resources(cpu="1", mem="1Gi")) def clean_data(raw: list[int]) -> list[int]: return [x for x in raw if x > 0] @task def average(nums: list[int]) -> float: return sum(nums) / len(nums) @workflow def pipeline(raw: list[int]) -> float: cleaned = clean_data(raw=raw) # Flyte sees clean_data → average dependency return average(nums=cleaned) # type mismatch here = caught before running
clean_data → average, runs each task in a container, type-checks the edges, and caches clean_data so the next run with the same raw skips straight to average.
Строит DAG clean_data → average, запускает каждую таску в контейнере, проверяет типы на рёбрах и кэширует clean_data, так что следующий запуск с тем же raw сразу переходит к average.
5. How it works (high level)5. Как он устроен (в целом)
You register a workflow with the Flyte backend (FlytePropeller, a Kubernetes operator). On execution, Flyte turns each task into a pod, passes typed inputs/outputs through blob storage (S3/GCS), and tracks state.
Ты регистрируешь воркфлоу в бэкенде Flyte (FlytePropeller — оператор Kubernetes). При запуске Flyte превращает каждую таску в pod, прокидывает типизированные входы/выходы через blob-хранилище (S3/GCS) и отслеживает состояние.
6. Flyte vs Airflow vs Dataswarm6. Flyte vs Airflow vs Dataswarm
| Airflow | Flyte | |
|---|---|---|
| Primary world | General scheduling / ETL | Data + ML, reproducibility-first |
| Typing | Untyped (XComs, dicts) | Strongly typed task interfaces |
| Caching | None built-in | Native, keyed by inputs + version |
| Isolation | Per-operator (varies) | Per-task container on K8s |
| Versioning | DAG file = current state | Every workflow/task version pinned |
| Dynamic graphs | Harder (dynamic task mapping newer) | First-class (@dynamic, map_task) |
| Airflow | Flyte | |
|---|---|---|
| Основной мир | Общее планирование / ETL | Data + ML, воспроизводимость во главе |
| Типизация | Без типов (XCom, словари) | Строго типизированные интерфейсы тасок |
| Кэширование | Нет «из коробки» | Нативное, по ключу входы + версия |
| Изоляция | По оператору (варьируется) | Контейнер на таску в K8s |
| Версионирование | DAG-файл = текущее состояние | Версия каждого воркфлоу/таски зафиксирована |
| Динамические графы | Сложнее (dynamic mapping появился позже) | Первокласно (@dynamic, map_task) |
7. Quick self-check7. Быстрая самопроверка
Answer in your head, then tap to flip.Ответь про себя, потом нажми, чтобы перевернуть.