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-скрипт без разделения на шаги оставляет вас в подвешенном состоянии при частичном сбое: непонятно, на каком этапе остановились и можно ли безопасно продолжить.
Если загрузка в витрину не идемпотентна, ретрай прогона после сбоя или ручной перезапуск дублируют строки в аналитической таблице.
Как Inquir помогает в этом сценарии
Ночной ETL как цепочка шагов пайплайна
Оформите извлечение, преобразование и загрузку как отдельные функции или шаги с явным порядком выполнения — тогда сбой на загрузке ретраится точечно, без перезапуска тяжёлого извлечения из источников.
В консоли видна история прогонов: какой шаг завершился с ошибкой, сколько строк прошло через каждый этап и когда состоялся следующий запуск по расписанию.
Что вы получаете на платформе
Что нужно для устойчивого ночного ETL
Изоляция шагов и точечные ретраи
Извлечение, преобразование и загрузка — отдельные шаги пайплайна; при сбое платформа ретраит только упавший шаг, а не весь контур.
Идемпотентная загрузка
Стратегия «очистить партицию и вставить заново», MERGE или UPSERT по ключу — чтобы повторный прогон давал тот же итог в целевой таблице, что и первый успешный.
История прогонов
Каждый запуск по cron фиксируется с журналами по шагам: видно, сколько строк извлечено, преобразовано и загружено.
Зависимости между шагами
Преобразование стартует только после успешного извлечения, а загрузка — только после успешного преобразования или записи промежуточного результата в надёжное хранилище.
Что сделать дальше, по шагам
Как организовать ночной ETL-пайплайн
Извлечение (extract)
Запросить данные из источников (база, HTTP API, файлы) за нужный интервал — например, за прошедшие сутки.
Преобразование (transform)
Агрегация, очистка, обогащение справочниками. Результат удобно сохранить во временную таблицу или объект в хранилище.
Загрузка (load)
Идемпотентная запись в аналитическое хранилище: например, очистка партиции и вставка, MERGE или UPSERT по бизнес-ключу.
Пример кода
Ночной 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 }; }
Когда подходит и когда нет
Когда Inquir уместен для ETL
Когда это уместно
- Контур до примерно десяти последовательных шагов без сложных циклов и динамически собираемого графа задач
- Нужны ретраи отдельных шагов при сбое без развёртывания отдельного оркестратора вроде Airflow
Когда лучше выбрать другое
- Десятки взаимозависимых задач, динамически собираемый граф и строгие SLA оркестрации — здесь специализированные системы (Airflow, Prefect и др.) обычно гибче
Вопросы и ответы
Вопросы и ответы
Как передать данные между шагами?
Обычно через объектное хранилище (ключ объекта передают в полезной нагрузке следующего шага) или через промежуточную таблицу в базе, если объём позволяет.
Можно ли параллельно извлекать из нескольких источников?
Да: используйте узел Parallel, чтобы разветвить поток на несколько независимых ветвей — например, одновременное чтение из разных API или баз, — а затем узел Merge, чтобы собрать результаты.