Inquir Compute · архитектура пайплайнов

Долгие serverless-задачи: цепочки шагов, fan-out и ретраи на шаг

Долгие serverless-задачи — не одна безграничная функция, а цепочка шагов, у каждого из которых свой бюджет таймаута (5 с по умолчанию, настраивается до 24 часов), политика ретраев и чекпоинт выхода. Когда падает шаг 7 из 10, вы повторяете шаг 7, а не всю задачу. Разветвляйтесь в параллельные ветки, собирайте результаты, ставьте паузу на подтверждение человеком и возобновляйте с места остановки.

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

  • Ретраи на шаг: заново запускается только упавший шаг, сделанная работа сохранена в чекпоинтах
  • Fan-out + merge: параллельные ветки сходятся в одном шаге агрегации
  • Паузы human gate: пайплайн ждёт подтверждения перед продолжением
  • Каждый шаг: свой таймаут (5 с по умолчанию, до 24 ч) в изолированном контейнере — цепочка ради чекпоинтов, а не из-за потолка

Суть ответа

Долгие serverless-задачи: цепочки шагов, fan-out и ретраи на шаг. Пайплайн — это граф шагов. Каждый шаг — вызов serverless-функции со своим таймаутом, политикой ретраев (maxAttempts, backoffMs, стратегия fixed или exponential) и записью выхода. Завершённые шаги не перезапускаются при падении следующего — пайплайн возобновляется с упавшего шага, а выходы предыдущих доступны как контекст.

Когда подходит

  • Работа раскладывается на независимые шаги, где ретраи на шаг стоят добавленной структуры
  • Fan-out: N параллельных подзадач, которые сходятся в один агрегированный результат
  • Паузы на подтверждение человеком или ожидание внешних событий посреди воркфлоу

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

  • Одна фоновая функция без декомпозиции на шаги означает ретраи «всё или ничего». Временная сетевая ошибка при записи в базу в конце 45-минутной трансформации перезапускает всю 45-минутную трансформацию.
  • Fan-out — обработка N записей параллельно — в одной функции требует N горутин или цепочек промисов внутри одного контейнера, без видимости, какие подзадачи успели до таймаута.

Почему долгая работа не помещается в одну функцию

Одна serverless-функция, каким бы длинным ни был её таймаут, не может надёжно выполнять многочасовой ETL, параллельное обогащение данных или воркфлоу с подтверждением человеком посередине. Одна функция — один домен сбоя: если шаг 7 из 10 падает на 50-й минуте, вы перезапускаете с нулевой.

Нарезать работу самостоятельно — вызывать Lambda рекурсивно, делить на батчи SQS — значит строить движок воркфлоу в коде приложения: отслеживать, какие чанки завершились, повторять нужные, агрегировать результаты. Именно эту проблему решают пайплайны.

Почему одношаговые async-задачи ломаются под сложностью

Одна фоновая функция без декомпозиции на шаги означает ретраи «всё или ничего». Временная сетевая ошибка при записи в базу в конце 45-минутной трансформации перезапускает всю 45-минутную трансформацию.

Fan-out — обработка N записей параллельно — в одной функции требует N горутин или цепочек промисов внутри одного контейнера, без видимости, какие подзадачи успели до таймаута.

Шаги пайплайна как независимые единицы выполнения

Пайплайн — это граф шагов. Каждый шаг — вызов serverless-функции со своим таймаутом, политикой ретраев (maxAttempts, backoffMs, стратегия fixed или exponential) и записью выхода. Завершённые шаги не перезапускаются при падении следующего — пайплайн возобновляется с упавшего шага, а выходы предыдущих доступны как контекст.

Fan-out — полноценный тип узла: параллельный узел порождает N ветвей одновременно; узел merge ждёт все ветви перед продолжением. Узлы human gate ставят пайплайн на паузу до внешнего события подтверждения. Всё это — конфигурация пайплайна, а не код приложения.

Паттерны многошаговой архитектуры пайплайна

Цепочка шагов с ретраями на шаг

Цепляйте шаги через dependsOn. У каждого шага свои число попыток и backoff (фиксированный или экспоненциальный). При падении шага 7 повторяется только шаг 7 — шаги 1–6 остаются в чекпоинтах.

Fan-out и merge

Параллельный узел порождает N ветвей одновременно. Узел merge агрегирует их выходы. Подходит для N-стороннего обогащения через API, параллельного ресайза изображений или массовой рассылки — без управления параллелизмом в коде приложения.

Пауза human gate и возобновление

Узел humanGate ставит пайплайн на паузу. Внешнее событие (вебхук, ручной триггер) возобновляет его. Используйте перед чувствительными действиями: рассылка 100 000 пользователям, списание, изменение продакшен-данных.

Долгая работа через композицию шагов

У каждого шага свой бюджет таймаута (5 с по умолчанию, до 24 часов). Многочасовой ETL всё равно декомпозируйте: извлечение (шаг 1), трансформация батча A (шаг 2), трансформация батча B (шаг 3), загрузка (шаг 4). Цепляйте их: платформа управляет последовательностью, а сбой перезапускает один шаг, а не всю ночь.

Как спроектировать многошаговый serverless-пайплайн

Сначала спроектируйте граф шагов, затем реализуйте каждый шаг как независимую функцию.

1

Декомпозировать работу на шаги

Нанесите задачу на граф: какие шаги последовательны (dependsOn), какие параллельны (узел parallel + merge), каким нужен human gate (узел humanGate). Размер шага выбирайте так, чтобы перезапуск после сбоя был дешёвым.

2

Реализовать и связать каждый шаг

Каждый шаг — serverless-функция, читающая event.previousOutput или event.stepResults для контекста предыдущих шагов. Настройте число попыток и backoff на шаг по характеру его сбоев.

3

Запустить из HTTP и наблюдать

Примите HTTP-запрос, вызовите global.durable.startNew(), верните 202. Следите за прогоном в истории выполнения: какой шаг работает, какой повторился, что вернул каждый.

Пайплайн HTTP → извлечение → fan-out → merge

HTTP-обработчик принимает задачу и возвращает 202. Оркестратор разветвляется в параллельные шаги обогащения, затем шаг merge агрегирует результаты. У каждого шага свой таймаут и политика ретраев.

api/start-enrichment.mjs (HTTP-вход)
export async function handler(event) {
  const { batchId, itemIds } = JSON.parse(event.body || '{}');
  if (!batchId || !itemIds?.length)
    return { statusCode: 400, body: JSON.stringify({ error: 'batchId and itemIds required' }) };
  const { instanceId: jobId } = await global.durable.startNew(
    'enrich-batch', undefined, { batchId, itemIds }
  );
  return { statusCode: 202, body: JSON.stringify({ jobId, status: 'started' }) };
}
steps/extract-items.mjs (шаг 1: последовательный)
export async function handler(event) {
  // step 1 — runs first, output feeds into fan-out branches
  const { batchId, itemIds } = event.payload ?? {};
  const items = await db.fetchItems(itemIds);        // may take several minutes
  await db.markBatchStarted(batchId);
  return { batchId, items };                          // passed to parallel branches via {{steps.extract.output}}
}
steps/enrich-item.mjs (шаг 2: параллельная ветка, на элемент)
export async function handler(event) {
  // Each parallel branch receives the upstream node's output as its event.
  // Node retry: maxAttempts=3, backoffMs=2000, strategy=exponential
  const enriched = await externalApi.enrich(event.item);  // retry isolates this network call
  return { itemId: event.item.id, enriched };
}
steps/merge-results.mjs (шаг 3: узел merge)
export async function handler(event) {
  // Runs after ALL parallel branches complete. A merge node hands the next step
  // { inputs: { <branchNodeId>: <branchOutput>, ... } } — one entry per branch.
  const branchOutputs = Object.values(event.inputs ?? {});
  const merged = branchOutputs.map(b => b.enriched).filter(Boolean);
  await db.saveBatchResults(merged);
  return { saved: merged.length, total: branchOutputs.length };
}

Используйте многошаговые пайплайны, когда…

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

  • Работа раскладывается на независимые шаги, где ретраи на шаг стоят добавленной структуры
  • Fan-out: N параллельных подзадач, которые сходятся в один агрегированный результат
  • Паузы на подтверждение человеком или ожидание внешних событий посреди воркфлоу

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

  • Простая последовательная работа, которую дёшево перезапустить целиком, — один шаг пайплайна без декомпозиции вполне подходит

Частые вопросы

Как долго может выполняться один шаг пайплайна?

Каждый шаг — вызов функции с настраиваемым таймаутом: 5 000 мс по умолчанию, до 24 часов (86 400 000 мс) при лимитах платформы по умолчанию. Долгую работу всё равно декомпозируйте на шаги — не из-за потолка, а чтобы сбой перезапускал один шаг и каждый шаг оставлял чекпоинт.

Как именно работают ретраи на шаг?

У каждого шага пайплайна своя конфигурация ретраев: maxAttempts, backoffMs и strategy (fixed или exponential). Когда шаг падает и исчерпывает попытки, пайплайн помечает его failed и прогон останавливается — ранее завершённые шаги не перезапускаются. Затем пайплайн можно переиграть с упавшего шага.

Может ли human gate ждать бесконечно?

Узел humanGate ставит пайплайн на паузу до прихода внешнего события через event API платформы. Жёсткого предела ожидания нет, но прогоны в статусе WAITING учитываются в лимитах воркспейса. Для compliance-воркфлоу с требованиями SLA задавайте таймауты ожидания.

Как ветви fan-out получают входные данные?

Параллельный узел передаёт выход предыдущего узла каждой исходящей ветви, и обработчик ветви читает его прямо как своё событие. Когда все ветви завершились, узел merge передаёт следующему шагу объект inputs с ключами по id узлов ветвей — перебирайте Object.values(event.inputs), чтобы прочитать выходы ветвей.