Розробка Workflow-двигуна на Apache Airflow під ключ
У вас зростає обсяг даних, але ручний запуск скриптів і cron'ів вже не справляється? Ми стикалися з ситуацією, коли N+1 запитів в ETL-пайплайні призводив до годинних затримок, а моніторинг був відсутній. Пайплайн завантаження 5 млн замовлень на LocalExecutor виконувався 4 години, а при пікових навантаженнях час простою сягав 30 хвилин. Ми вирішуємо цю проблему за допомогою Apache Airflow — зрілої платформи для оркестрації та автоматизації процесів обробки даних. Наші інженери побудували десятки DAG-пайплайнів для компаній з e-commerce, fintech та логістики. Кожен пайплайн проходить навантажувальне тестування: типовий DAG обробляє до 10 млн записів за ніч, вкладаючись в SLA при 99.9% часу виконання. Стабільність підтверджена 50+ проектами. Отримайте консультацію по вашому проекту — оцінимо за 1 день.
Чому варто обрати KubernetesExecutor?
Airflow оптимізований для batch-обробки даних:
- ETL/ELT пайплайни (PostgreSQL → трансформація → Data Warehouse)
- Щоденні звіти та вивантаження
- ML-пайплайни (підготовка даних → навчання → деплой моделі)
- Періодичні агрегації та синхронізації
Якщо у вас подійно-керовані бізнес-процеси з human tasks — варто придивитися до Temporal або Camunda. Airflow не вміє чекати користувача годинами. Але для data-інженерів це найкращий вибір: Airflow в 3-5 разів швидший за Temporal в batch-сценаріях, а також краще підходить для ETL-пайплайнів.
Переваги KubernetesExecutor
| Executor | Масштабування | Ізоляція | Управління ресурсами |
|---|---|---|---|
| LocalExecutor | Обмежено однією нодою | Ні | Ручне |
| CeleryExecutor | Горизонтальне через workers | Середня | Вимагає Redis/RabbitMQ |
| KubernetesExecutor | Автоматичне | Кожне завдання в Pod | Через requests/limits |
KubernetesExecutor дає ізоляцію на рівні завдань: кожне запускається в окремому Pod з власними CPU та пам'яттю. При пікових навантаженнях Kubernetes автоматично піднімає Pod'и, а після — утилізує. Ми використовуємо цей підхід у production і вважаємо його стандартом для сучасних data-пайплайнів.
Як Airflow пришвидшує обробку даних?
При піковому навантаженні 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:
persistence:
enabled: true
gitSync:
enabled: true
repo: https://github.com/yourcompany/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),
}
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 в дії на Apache Airflow
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 та логістики. Середня uptime пайплайнів — 99.9%. Використовуємо best practices: GitSync, KubernetesExecutor, сповіщення, версіонування DAG'ів. Гарантуємо стабільну роботу пайплайнів та документацію рівня enterprise.
Ми готові розробити аналогічний пайплайн для ваших даних. Зв'яжіться з нами для оцінки — відповімо за 1 день. Або замовте консультацію — обговоримо деталі без зобов'язань.







