Сценарий · ETL

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-бэкенда.

Шаги пайплайна как лёгкий оркестратор 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-пайплайна

1

Шаг extract: забрать из источников

Вызовите внешние API или запросите базу-источник. Сохраните сырые данные в промежуточное хранилище или передайте как payload пайплайна.

2

Шаг transform: применить бизнес-логику

Очистка, нормализация, обогащение и агрегация. Python с pandas; Node.js для JSON-трансформаций. На выходе — структурированные данные.

3

Шаг load: записать в приёмник

Upsert в хранилище данных, обновление аналитических таблиц, расчёт агрегатов. Идемпотентно по ID батча или диапазону дат.

Ночной ETL: API → transform → Postgres

Графовый пайплайн: узел cronTrigger ведёт в узлы-функции extract → transform → load, соединённые рёбрами, — без Airflow DAG. Extract забирает сырые данные, transform нормализует, load делает идемпотентный upsert; каждый узел ретраится независимо.

pipelines/nightly-etl.json (граф пайплайна)
{
  "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" }
  ]
}
jobs/etl-extract.mjs (шаг 1)
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 };
}
jobs/etl-transform.mjs (шаг 2 — получает вывод шага 1)
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 };
}
jobs/etl-load.mjs (шаг 3 — получает вывод шага 2)
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 или диапазону дат). Если ретраи исчерпаны, история выполнений фиксирует сбой и срабатывают алерты.