У вас растёт объём данных, но ручной запуск скриптов и cron'ов уже не справляется? Мы сталкивались с ситуацией, когда N+1 запросов в ETL-пайплайне приводил к часовым задержкам, а мониторинг отсутствовал. Пайплайн загрузки 5 млн заказов на LocalExecutor выполнялся 4 часа, а при пиковых нагрузках время простоя достигало 30 минут. Мы решаем эту проблему с помощью Apache Airflow — зрелой платформы для оркестрации и автоматизации процессов обработки данных. Наши инженеры построили десятки DAG-пайплайнов для компаний из e-commerce, fintech и логистики. Каждый пайплайн проходит нагрузочное тестирование: типичный DAG обрабатывает до 10 млн записей за ночь, укладываясь в SLA при 99.9% времени выполнения. Стабильность подтверждена 50+ проектами. Получите консультацию по вашему проекту — оценим за 1 день.
Сравнение Airflow с Temporal и Camunda
Airflow оптимизирован для batch-обработки данных:
- ETL/ELT пайплайны (PostgreSQL → трансформация → Data Warehouse)
- Ежедневные отчёты и выгрузки
- ML-пайплайны (подготовка данных → обучение → деплой модели)
- Периодические агрегации и синхронизации
Если у вас событийно-управляемые бизнес-процессы с human tasks — стоит присмотреться к Temporal или Camunda. Airflow не умеет ждать пользователя часами. Но для data-инженеров это лучший выбор: он в 3-5 раз быстрее Temporal в batch-сценариях.
Почему мы выбираем KubernetesExecutor?
| Executor | Масштабирование | Изоляция | Управление ресурсами |
|---|---|---|---|
| LocalExecutor | Ограничено одной нодой | Нет | Ручное |
| CeleryExecutor | Горизонтальное через workers | Средняя | Требует Redis/RabbitMQ |
| KubernetesExecutor | Автоматическое | Каждая задача в Pod | Через requests/limits |
KubernetesExecutor даёт изоляцию на уровне задач: каждая запускается в отдельном Pod с собственными CPU и памятью. При пиковых нагрузках Kubernetes автоматически поднимает Pod'ы, а после — утилизирует. Мы используем этот подход в production и считаем его стандартом для современных data-пайплайнов.
Как KubernetesExecutor ускоряет обработку данных?
При пиковой нагрузке KubernetesExecutor масштабирует ресурсы горизонтально: вместо одной ноды — кластер из 10+ Pod'ов. Время выполнения ETL-пайплайна сокращается в 3-5 раз по сравнению с LocalExecutor. Например, пайплайн загрузки 5 млн заказов на LocalExecutor выполняется 4 часа, на KubernetesExecutor — 50 минут. Разница очевидна.
Как мы строим ETL-пайплайн: развёрнутый кейс
Рассмотрим пример для интернет-магазина: ежедневная загрузка заказов из PostgreSQL в DWH (PostgreSQL) с трансформацией и агрегацией. Ниже — DAG, который мы развернули заказчику за 2 дня.
Установка через Helm
helm repo add apache-airflow https://airflow.apache.org
helm upgrade --install airflow apache-airflow/airflow \
--namespace airflow \
--create-namespace \
--set executor=KubernetesExecutor \
--set postgresql.enabled=true \
--set redis.enabled=true \
--values airflow-values.yaml
# airflow-values.yaml
aiflow:
image:
repository: apache/airflow
tag: 2.8.0
config:
AIRFLOW__CORE__DAGS_FOLDER: /opt/airflow/dags
AIRFLOW__CORE__MAX_ACTIVE_RUNS_PER_DAG: "3"
AIRFLOW__SCHEDULER__MIN_FILE_PROCESS_INTERVAL: "30"
dags:
gitSync:
enabled: true
repo: https://github.com/company/airflow-dags.git
branch: main
subPath: dags/
DAG — пример ETL-пайплайна
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow.providers.postgres.hooks.postgres import PostgresHook
from datetime import datetime, timedelta
import pandas as pd
default_args = {
'owner': 'data-team',
'depends_on_past': False,
'start_date': datetime(2024, 1, 1),
'retries': 2,
'retry_delay': timedelta(minutes=5),
'email_on_failure': True,
'email': ['[email protected]'],
}
with DAG(
'daily_orders_etl',
default_args=default_args,
schedule_interval='0 2 * * *',
catchup=False,
tags=['etl', 'orders'],
description='Загрузка и трансформация заказов в DWH',
) as dag:
def extract_orders(**context):
hook = PostgresHook(postgres_conn_id='production_db')
ds = context['ds']
df = hook.get_pandas_df(f"""
SELECT o.id, o.customer_id, o.total, o.status,
o.created_at, c.email, c.country
FROM orders o
JOIN customers c ON c.id = o.customer_id
WHERE o.created_at::date = '{ds}'
AND o.status IN ('paid', 'shipped', 'delivered')
""")
context['ti'].xcom_push(key='orders_count', value=len(df))
df.to_parquet(f'/tmp/orders_{ds}.parquet')
return len(df)
def transform_orders(**context):
ds = context['ds']
df = pd.read_parquet(f'/tmp/orders_{ds}.parquet')
df['order_date'] = pd.to_datetime(df['created_at']).dt.date
df['revenue_usd'] = df['total'] / 100
df['is_international'] = df['country'] != 'RU'
df['customer_tier'] = df['revenue_usd'].apply(
lambda x: 'vip' if x >= 500 else 'regular'
)
df.to_parquet(f'/tmp/orders_transformed_{ds}.parquet')
def load_to_dwh(**context):
ds = context['ds']
df = pd.read_parquet(f'/tmp/orders_transformed_{ds}.parquet')
hook = PostgresHook(postgres_conn_id='datawarehouse')
engine = hook.get_sqlalchemy_engine()
df.to_sql('fact_orders', engine, schema='dwh',
if_exists='append', index=False,
method='multi', chunksize=1000)
aggregate_metrics = PostgresOperator(
task_id='aggregate_metrics',
postgres_conn_id='datawarehouse',
sql="""
INSERT INTO dwh.daily_metrics (date, total_revenue, orders_count, avg_order)
SELECT
'{{ ds }}'::date,
SUM(revenue_usd),
COUNT(*),
AVG(revenue_usd)
FROM dwh.fact_orders
WHERE order_date = '{{ ds }}'
ON CONFLICT (date) DO UPDATE SET
total_revenue = EXCLUDED.total_revenue,
orders_count = EXCLUDED.orders_count,
avg_order = EXCLUDED.avg_order;
""",
)
extract = PythonOperator(task_id='extract_orders', python_callable=extract_orders)
transform = PythonOperator(task_id='transform_orders', python_callable=transform_orders)
load = PythonOperator(task_id='load_to_dwh', python_callable=load_to_dwh)
extract >> transform >> load >> aggregate_metrics
Параллельное выполнение и Sensors
from airflow.utils.task_group import TaskGroup
with TaskGroup('process_regions') as process_regions:
for region in ['EU', 'US', 'APAC']:
PythonOperator(
task_id=f'process_{region.lower()}',
python_callable=process_region_data,
op_kwargs={'region': region}
)
extract >> process_regions >> aggregate_all
Для ожидания внешних событий используем Sensors: FileSensor для файлов, HttpSensor для API. Это стандартный паттерн для production.
KubernetesExecutor в действии
executor_config = {
'KubernetesExecutor': {
'request_memory': '2Gi',
'request_cpu': '500m',
'limit_memory': '4Gi',
'image': 'custom-airflow:2.8.0-pandas',
}
}
heavy_transform = PythonOperator(
task_id='heavy_transform',
python_callable=transform_large_dataset,
executor_config=executor_config
)
Как обрабатываются сбои DAG?
При сбое задачи Airflow автоматически выполняет повторные попытки (retries) с экспоненциальной задержкой. Мы настраиваем алерты в Telegram/Slack на каждый сбой и длительное выполнение. Дополнительно интегрируем метрики в Prometheus/Grafana: дашборды показывают время выполнения, количество успешных/упавших задач, загрузку ресурсов. Это позволяет быстро реагировать на инциденты.
Минимальные требования к инфраструктуре:
- Kubernetes кластер версии 1.24+
- PostgreSQL 13+ для metadata database
- Redis (опционально, для CeleryExecutor)
- Объем хранилища: от 100 ГБ для логов и артефактов
Процесс работы
- Анализ: изучаем ваши источники данных, объёмы, SLA, текущие проблемы, настройку Airflow и инфраструктуру.
- Проектирование: описываем DAG'и, выбираем Executor, продумываем обработку ошибок и алерты.
- Реализация: пишем код DAG'ов, конфиги Helm, CI/CD, тесты.
- Деплой: разворачиваем инфраструктуру (Kubernetes или Docker), настраиваем GitSync и мониторинг.
- Тестирование: прогоняем backfill на исторических данных, проверяем корректность.
- Документация и обучение: передаём runbook, учим вашу команду.
- Поддержка: 2 недели после запуска для стабилизации.
Что входит в работу
- Разработанные DAG'и (Python-код, готовый к production).
- Конфигурация Airflow (values.yaml, переменные, connections).
- CI/CD pipeline для автоматического деплоя DAG'ов через Git.
- Docker-образ с зависимостями (библиотеки, драйверы).
- Документация: описание DAG'ов, инструкция по запуску, troubleshooting.
- Обучение команды (2-3 сессии).
- Мониторинг: алерты в Telegram/Slack, дашборды Grafana.
Сроки реализации
| Этап | Срок |
|---|---|
| Airflow deployment (Helm/Docker) + первый DAG | 3–5 дней |
| ETL-пайплайн с 5–8 задачами, трансформациями и DWH-загрузкой | 1–2 недели |
| Сложный пайплайн с параллелизмом, sensors и backfill | 2–4 недели |
Наш опыт и гарантии
У нас 10+ лет опыта в разработке data-инфраструктуры. Мы реализовали 50+ проектов на Airflow для e-commerce, fintech и логистики. Используем best practices: GitSync, KubernetesExecutor, алертинг, версионирование DAG'ов. Гарантируем стабильную работу пайплайнов и документацию уровня enterprise.
Мы готовы разработать аналогичный пайплайн для ваших данных. Свяжитесь с нами для оценки — ответим за 1 день. Или закажите консультацию — обсудим детали без обязательств.







