Flyte — the friendly introFlyte — дружелюбное введение

A Kubernetes-native orchestrator for data & ML pipelines, built around strong typing, versioning and caching. Plain words, code-first. Kubernetes-native оркестратор для data- и ML-пайплайнов, построенный вокруг строгой типизации, версионирования и кэширования. Простыми словами, код в центре.
🚀 Flyte 📅 AirflowAirflow 🏗️ Terraform

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 — каждый шаг в своём контейнере, с ресурсами, ретраями и кэшированием «из коробки».

AnalogyАналогия Think of Flyte as a kitchen pass with strict recipe cards. Each cook (task) gets a card that says exactly what ingredients (typed inputs) come in and what dish (typed outputs) goes out. If the same dish was already cooked from the same ingredients, the pass just hands back the cached plate instead of cooking again. Представь Flyte как раздачу на кухне со строгими рецептурными карточками. Каждый повар (таска) получает карточку: какие именно ингредиенты (типизированные входы) приходят и какое блюдо (типизированные выходы) выходит. Если такое блюдо уже готовили из тех же ингредиентов — раздача просто отдаёт сохранённую тарелку, а не готовит заново.
One lineОдной строкой Flyte = typed, versioned, cached Python pipelines that run as containers on Kubernetes. Reproducibility and ML-scale are its whole point. Flyte = типизированные, версионируемые, кэшируемые Python-пайплайны, бегущие как контейнеры на Kubernetes. Воспроизводимость и ML-масштаб — его суть.

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 — первоклассные граждане.
The selling pointГлавный плюс Reproducibility by construction: same code version + same inputs → same outputs, and Flyte can prove it and reuse them. Воспроизводимость по построению: та же версия кода + те же входы → те же выходы, и Flyte это гарантирует и переиспользует.

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
What Flyte does with thisЧто Flyte с этим делает Builds the DAG 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) и отслеживает состояние.

Python (@task/@workflow) └─▶ flytekit packages + registers a versioned workflow └─▶ FlytePropeller (K8s operator) plans the DAG └─▶ each task ─▶ a Pod (its own container + resources) └─▶ typed inputs/outputs via S3/GCS └─▶ cache lookup: hit ⇒ skip, miss ⇒ run + store
Mental modelМодель в голове Flyte = a typed contract layer + a Kubernetes execution engine. The contract gives reproducibility; K8s gives scale and isolation. Flyte = слой типизированных контрактов + движок исполнения на Kubernetes. Контракт даёт воспроизводимость; K8s даёт масштаб и изоляцию.

6. Flyte vs Airflow vs Dataswarm6. Flyte vs Airflow vs Dataswarm

AirflowFlyte
Primary worldGeneral scheduling / ETLData + ML, reproducibility-first
TypingUntyped (XComs, dicts)Strongly typed task interfaces
CachingNone built-inNative, keyed by inputs + version
IsolationPer-operator (varies)Per-task container on K8s
VersioningDAG file = current stateEvery workflow/task version pinned
Dynamic graphsHarder (dynamic task mapping newer)First-class (@dynamic, map_task)
AirflowFlyte
Основной мирОбщее планирование / ETLData + ML, воспроизводимость во главе
ТипизацияБез типов (XCom, словари)Строго типизированные интерфейсы тасок
КэшированиеНет «из коробки»Нативное, по ключу входы + версия
ИзоляцияПо оператору (варьируется)Контейнер на таску в K8s
ВерсионированиеDAG-файл = текущее состояниеВерсия каждого воркфлоу/таски зафиксирована
Динамические графыСложнее (dynamic mapping появился позже)Первокласно (@dynamic, map_task)
Meta contextКонтекст Meta Meta's internal Dataswarm/Chronos plays a similar orchestration role; Flyte is the popular open-source equivalent (Lyft, Spotify) you'd reach for outside Meta. Same job — schedule a typed DAG — different ecosystem. Внутренние Dataswarm/Chronos в Meta играют похожую роль оркестрации; Flyte — популярный опенсорсный аналог (Lyft, Spotify), который берут вне Meta. Та же задача — запустить типизированный DAG — другая экосистема.

7. Quick self-check7. Быстрая самопроверка

Answer in your head, then tap to flip.Ответь про себя, потом нажми, чтобы перевернуть.

What's the one-line definition of Flyte?Определение Flyte одной строкой?
tapнажми
A Kubernetes-native orchestrator for typed, versioned, cached data/ML pipelines written in Python.Kubernetes-native оркестратор для типизированных, версионируемых, кэшируемых data/ML-пайплайнов на Python.
Task vs Workflow?Task vs Workflow?
tapнажми
A task (@task) is one containerized step with typed I/O. A workflow (@workflow) wires tasks into a DAG via their data dependencies.Task (@task) — один контейнеризованный шаг с типизированным I/O. Workflow (@workflow) связывает таски в DAG через их зависимости по данным.
Why does Flyte cache, and on what key?Зачем Flyte кэширует и по какому ключу?
tapнажми
To skip recomputation. Outputs are keyed by (inputs, cache_version) — same inputs + version ⇒ cached result returned instantly.Чтобы не пересчитывать. Выходы по ключу (входы, cache_version) — те же входы + версия ⇒ сразу отдаётся сохранённый результат.
Biggest advantage over classic Airflow?Главное преимущество над классическим Airflow?
tapнажми
Strong typing + native caching + per-task containers = reproducibility by construction. Airflow is untyped with no built-in caching.Строгая типизация + нативный кэш + контейнер на таску = воспроизводимость по построению. Airflow без типов и без встроенного кэша.
When do you use @dynamic / map_task?Когда нужны @dynamic / map_task?
tapнажми
When the DAG shape depends on data: @dynamic builds the graph at runtime; map_task fans one task out over a list in parallel.Когда форма DAG зависит от данных: @dynamic строит граф в рантайме; map_task распараллеливает одну таску по списку.
What runs Flyte tasks under the hood?Что исполняет таски Flyte под капотом?
tapнажми
Kubernetes. FlytePropeller (a K8s operator) turns each task into a pod and passes typed I/O via blob storage (S3/GCS).Kubernetes. FlytePropeller (оператор K8s) превращает каждую таску в pod и прокидывает типизированный I/O через blob-хранилище (S3/GCS).