Serverless ночной ETL без Airflow и Kubernetes
Запускайте ночные ETL-пайплайны — извлечение из API или баз, трансформация на Node.js или Python, загрузка в хранилище данных — как шаги serverless-пайплайна с cron-триггерами, ретраями на уровне шагов и историей выполнений. Без Airflow DAG, за которым нужно следить, и без кластера Kubernetes под воркеры.
Обновлено: 2026-06-28
Кратко
Суть ответа
Serverless ночной ETL без Airflow и Kubernetes. Пайплайны Inquir соединяют узлы serverless-функций рёбрами. Узел extract забирает данные из источника, transform применяет бизнес-логику, load пишет в приёмник. У каждого узла свой бюджет ретраев; упавший load ретраится без повторного запуска extract.
Когда подходит
- У вас 1–10 ночных задач синхронизации данных, которые не оправдывают кластер Airflow
- Шаги ETL занимают 5–60 минут и нуждаются в ретраях на уровне шагов без перезапуска с нуля
На что обратить внимание
- Airflow требует базу метаданных PostgreSQL, процесс планировщика, воркеры и, опционально, Celery или Kubernetes для выполнения задач. Для ночного ETL из трёх шагов, который отрабатывает за 20 минут, стоимость инфраструктуры для большинства молодых команд превышает пользу.
- Даже managed Airflow (MWAA, Astronomer, Cloud Composer) несёт заметную фиксированную стоимость и модель написания DAG, отличную от остального serverless-бэкенда.
Нагрузка и где ломается
Варианты оркестрации ETL и их компромиссы
- Airflow: сильная оркестрация DAG, но нужно поддерживать кластер Kubernetes или managed-сервис
- Prefect / Dagster: DX лучше, чем у Airflow, но всё равно нужен отдельный деплой
- crontab: просто, но без истории прогонов, ретраев и зависимостей между шагами
- Lambda + EventBridge: возможно, но цепочка шагов и наблюдаемость требуют серьёзной обвязки
Большинство ночных ETL-задач в молодых компаниях недостаточно сложны, чтобы оправдать Airflow. Им нужны надёжное расписание, ретраи на уровне шагов, видимая история прогонов и общие секреты с остальным бэкендом. Пайплайны Inquir закрывают это без отдельного сервиса оркестрации.
Когда простых рецептов недостаточно
Почему Airflow избыточен для небольших ETL
Airflow требует базу метаданных PostgreSQL, процесс планировщика, воркеры и, опционально, Celery или Kubernetes для выполнения задач. Для ночного ETL из трёх шагов, который отрабатывает за 20 минут, стоимость инфраструктуры для большинства молодых команд превышает пользу.
Даже managed Airflow (MWAA, Astronomer, Cloud Composer) несёт заметную фиксированную стоимость и модель написания DAG, отличную от остального serverless-бэкенда.
Как помогает Inquir
Шаги пайплайна как лёгкий оркестратор ETL
Пайплайны Inquir соединяют узлы serverless-функций рёбрами. Узел extract забирает данные из источника, transform применяет бизнес-логику, load пишет в приёмник. У каждого узла свой бюджет ретраев; упавший load ретраится без повторного запуска extract.
Ночные ETL-пайплайны делят секреты воркспейса (URL баз, API-ключи) и наблюдаемость с HTTP-маршрутами и фоновыми задачами. Одна платформа для всех бэкенд-нагрузок — без отдельного кластера Airflow.
Что вы получаете
Из чего состоит ночной ETL-пайплайн
Многошаговый граф пайплайна
Extract → Transform → Load как узлы-функции, соединённые рёбрами. Упавший load ретраится независимо, без перезапуска extract.
Параллельное извлечение (fan-out)
Запустите N параллельных шагов extract (по одному на источник) и соберите их в один шаг transform. Общее время ETL сокращается.
Python для трансформаций
Python 3.12 с pandas, numpy и sqlalchemy для шагов трансформации. Node.js для шагов извлечения из API. Оба в одном пайплайне.
Алерт по SLO длительности
Настройте алерт, когда ночной ETL длится дольше ожидаемого, — до того, как устаревшие данные попадут в дашборды и отчёты.
Что дальше
Структура ночного ETL-пайплайна
Шаг extract: забрать из источников
Вызовите внешние API или запросите базу-источник. Сохраните сырые данные в промежуточное хранилище или передайте как payload пайплайна.
Шаг transform: применить бизнес-логику
Очистка, нормализация, обогащение и агрегация. Python с pandas; Node.js для JSON-трансформаций. На выходе — структурированные данные.
Шаг load: записать в приёмник
Upsert в хранилище данных, обновление аналитических таблиц, расчёт агрегатов. Идемпотентно по ID батча или диапазону дат.
Пример кода
Ночной ETL: API → transform → Postgres
Графовый пайплайн: узел cronTrigger ведёт в узлы-функции extract → transform → load, соединённые рёбрами, — без Airflow DAG. Extract забирает сырые данные, transform нормализует, load делает идемпотентный upsert; каждый узел ретраится независимо.
{ "schemaVersion": 1, "nodes": [ { "id": "cron", "kind": "cronTrigger", "name": "Nightly 02:00", "position": { "x": 0, "y": 0 }, "config": { "cron": "0 2 * * *", "timezone": "UTC" } }, { "id": "extract", "kind": "lambda", "name": "Extract daily report", "position": { "x": 240, "y": 0 }, "config": { "functionId": "etl-extract", "onError": "failPipeline" } }, { "id": "transform", "kind": "lambda", "name": "Transform records", "position": { "x": 480, "y": 0 }, "config": { "functionId": "etl-transform", "onError": "failPipeline" } }, { "id": "load", "kind": "lambda", "name": "Load to warehouse", "position": { "x": 720, "y": 0 }, "config": { "functionId": "etl-load", "onError": "failPipeline" } } ], "edges": [ { "id": "e1", "sourceNodeId": "cron", "targetNodeId": "extract", "sourceHandle": "default" }, { "id": "e2", "sourceNodeId": "extract", "targetNodeId": "transform", "sourceHandle": "success" }, { "id": "e3", "sourceNodeId": "transform", "targetNodeId": "load", "sourceHandle": "success" } ] }
export async function handler(event) { const date = new Date().toISOString().slice(0, 10); // YYYY-MM-DD const records = await externalApi.fetchDailyReport(date); // Store to object storage — pipeline passes storage key to next step const key = `etl/raw/${date}.json`; await storage.putJson(key, records); return { key, count: records.length, date }; }
export async function handler(event) { const { key, date } = event.previousOutput ?? {}; const raw = await storage.getJson(key); const transformed = raw.map((r) => ({ id: r.external_id, date, revenue: parseFloat(r.revenue_usd), region: r.region?.toLowerCase(), updatedAt: new Date().toISOString(), })); const outKey = `etl/transformed/${date}.json`; await storage.putJson(outKey, transformed); return { outKey, count: transformed.length, date }; }
export async function handler(event) { const { outKey, date } = event.previousOutput ?? {}; const rows = await storage.getJson(outKey); // Idempotent load — safe to retry without duplicating rows await db.analytics.upsertBatch(rows, { conflictKeys: ['id', 'date'] }); return { loaded: rows.length, date }; }
Когда подходит
Используйте serverless ночной ETL, когда
Когда это уместно
- У вас 1–10 ночных задач синхронизации данных, которые не оправдывают кластер Airflow
- Шаги ETL занимают 5–60 минут и нуждаются в ретраях на уровне шагов без перезапуска с нуля
Когда лучше выбрать другое
- Сложные DAG с сотнями задач, ветвлениями и зависимостями между DAG — на таком масштабе Airflow или Dagster оркестрируют лучше
Частые вопросы
Частые вопросы
Можно ли передавать данные между шагами пайплайна?
Да — шаги пайплайна возвращают структурированный вывод. Следующий шаг получает его как event.previousOutput. Большие payload (CSV, JSON-файлы) сохраняйте в объектное хранилище и передавайте между шагами ключ объекта.
Что делать с упавшим шагом load?
Настройте число ретраев и задержку на шаге load. Убедитесь, что шаг идемпотентен (upsert по стабильному ID или диапазону дат). Если ретраи исчерпаны, история выполнений фиксирует сбой и срабатывают алерты.