Сценарий · Inquir Compute

Serverless ночной ETL без Airflow и Kubernetes

Запускайте ночной ETL — извлечение из API и баз, трансформацию на Node.js или Python, загрузку в хранилище или витрину — как шаги serverless-пайплайна с cron-триггером, ретраями на уровне отдельных шагов и историей выполнений. Без Airflow DAG, за которым нужно следить, и без кластера Kubernetes под воркеры.

Обновлено: 2026-06-28

Суть ответа

Serverless ночной ETL без Airflow и Kubernetes. Оформите извлечение, преобразование и загрузку как отдельные функции или шаги с явным порядком выполнения — тогда сбой на загрузке ретраится точечно, без перезапуска тяжёлого извлечения из источников.

Когда подходит и когда нет

  • Контур до примерно десяти последовательных шагов без сложных циклов и динамически собираемого графа задач
  • Нужны ретраи отдельных шагов при сбое без развёртывания отдельного оркестратора вроде Airflow

На что обратить внимание

  • Один длинный cron-скрипт без разделения на шаги оставляет вас в подвешенном состоянии при частичном сбое: непонятно, на каком этапе остановились и можно ли безопасно продолжить.
  • Если загрузка в витрину не идемпотентна, ретрай прогона после сбоя или ручной перезапуск дублируют строки в аналитической таблице.

Варианты оркестрации ETL и их компромиссы

  • Airflow: мощная оркестрация DAG, но нужен Kubernetes или managed-сервис
  • Prefect / Dagster: лучший DX, но всё равно отдельный деплой
  • crontab: просто, но без истории прогонов, ретраев и зависимостей между шагами
  • Lambda + EventBridge: возможно, но цепочка шагов и наблюдаемость требуют серьёзной обвязки

Airflow, Luigi, Prefect и аналоги оправданы при больших графах задач и сложных зависимостях. Если же у вас ночной цикл из трёх–пяти понятных шагов, заводить под него отдельный кластер оркестратора, отдельный деплой и отдельную систему мониторинга часто невыгодно: растёт операционная нагрузка, а не качество данных.

Где ломаются «простые» сценарии ETL

Один длинный cron-скрипт без разделения на шаги оставляет вас в подвешенном состоянии при частичном сбое: непонятно, на каком этапе остановились и можно ли безопасно продолжить.

Если загрузка в витрину не идемпотентна, ретрай прогона после сбоя или ручной перезапуск дублируют строки в аналитической таблице.

Ночной ETL как цепочка шагов пайплайна

Оформите извлечение, преобразование и загрузку как отдельные функции или шаги с явным порядком выполнения — тогда сбой на загрузке ретраится точечно, без перезапуска тяжёлого извлечения из источников.

В консоли видна история прогонов: какой шаг завершился с ошибкой, сколько строк прошло через каждый этап и когда состоялся следующий запуск по расписанию.

Что нужно для устойчивого ночного ETL

Изоляция шагов и точечные ретраи

Извлечение, преобразование и загрузка — отдельные шаги пайплайна; при сбое платформа ретраит только упавший шаг, а не весь контур.

Идемпотентная загрузка

Стратегия «очистить партицию и вставить заново», MERGE или UPSERT по ключу — чтобы повторный прогон давал тот же итог в целевой таблице, что и первый успешный.

История прогонов

Каждый запуск по cron фиксируется с журналами по шагам: видно, сколько строк извлечено, преобразовано и загружено.

Зависимости между шагами

Преобразование стартует только после успешного извлечения, а загрузка — только после успешного преобразования или записи промежуточного результата в надёжное хранилище.

Как организовать ночной ETL-пайплайн

1

Извлечение (extract)

Запросить данные из источников (база, HTTP API, файлы) за нужный интервал — например, за прошедшие сутки.

2

Преобразование (transform)

Агрегация, очистка, обогащение справочниками. Результат удобно сохранить во временную таблицу или объект в хранилище.

3

Загрузка (load)

Идемпотентная запись в аналитическое хранилище: например, очистка партиции и вставка, MERGE или UPSERT по бизнес-ключу.

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

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

pipelines/nightly-etl.json (pipeline graph)
{
  "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 (step 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 (step 2 — receives step 1 output)
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 (step 3 — receives step 2 output)
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 };
}

Когда Inquir уместен для ETL

Когда это уместно

  • Контур до примерно десяти последовательных шагов без сложных циклов и динамически собираемого графа задач
  • Нужны ретраи отдельных шагов при сбое без развёртывания отдельного оркестратора вроде Airflow

Когда лучше выбрать другое

  • Десятки взаимозависимых задач, динамически собираемый граф и строгие SLA оркестрации — здесь специализированные системы (Airflow, Prefect и др.) обычно гибче

Вопросы и ответы

Как передать данные между шагами?

Обычно через объектное хранилище (ключ объекта передают в полезной нагрузке следующего шага) или через промежуточную таблицу в базе, если объём позволяет.

Можно ли параллельно извлекать из нескольких источников?

Да: используйте узел Parallel, чтобы разветвить поток на несколько независимых ветвей — например, одновременное чтение из разных API или баз, — а затем узел Merge, чтобы собрать результаты.