Сценарий · синхронизация данных

Serverless-синхронизация данных по расписанию с инкрементальными курсорами

Автоматизируйте инкрементальную синхронизацию между внешними API и внутренними базами на serverless cron-пайплайнах: курсор-отметка последней успешной синхронизации, идемпотентный upsert, история выполнений по каждому прогону и ретраи при сбоях — без crontab на VPS и без постоянно работающего воркера.

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

Суть ответа

Serverless-синхронизация данных по расписанию с инкрементальными курсорами. Пайплайны Inquir по расписанию запускаются как serverless-функции с историей выполнений. Обработчик синхронизации читает последнюю успешную отметку (из базы или другого хранилища состояния), забирает только записи, обновлённые после этого курсора, делает идемпотентный upsert и сохраняет новую отметку.

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

  • Данные CRM, ERP или сторонних API, которые нужно периодически обновлять в вашей базе
  • Часовая или суточная синхронизация дельты, когда полная выгрузка слишком медленная или дорогая по квоте

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

  • Синхронизация всех записей на каждом прогоне работает на маленьких наборах данных. На 100 000+ записей полная синхронизация длится слишком долго и съедает слишком много квоты API. Внешние API начинают вас ограничивать, и синхронизация падает время от времени.
  • Без курсора-отметки каждая упавшая синхронизация начинается заново. Прогон, упавший на записи 95 000, стартует с нулевой — расточительно, медленно и с большой вероятностью того же сбоя.

Почему синхронизация по расписанию падает тихо

  • Crontab на VPS: упавшие синхронизации уходят в /dev/null или в почту root, которую никто не читает
  • Полная пересинхронизация на каждом прогоне: повторная выгрузка всех записей тратит квоту API, а время синхронизации растёт квадратично
  • Нет идемпотентности: сетевой сбой посреди синхронизации оставляет частичное состояние и дубли записей при ретрае
  • Нет истории прогонов: на вопрос «синхронизация ночью запускалась?» отвечают через SSH и поиск по логам

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

Почему cron-синхронизация без курсора не масштабируется

Синхронизация всех записей на каждом прогоне работает на маленьких наборах данных. На 100 000+ записей полная синхронизация длится слишком долго и съедает слишком много квоты API. Внешние API начинают вас ограничивать, и синхронизация падает время от времени.

Без курсора-отметки каждая упавшая синхронизация начинается заново. Прогон, упавший на записи 95 000, стартует с нулевой — расточительно, медленно и с большой вероятностью того же сбоя.

Инкрементальные пайплайны с состоянием курсора

Пайплайны Inquir по расписанию запускаются как serverless-функции с историей выполнений. Обработчик синхронизации читает последнюю успешную отметку (из базы или другого хранилища состояния), забирает только записи, обновлённые после этого курсора, делает идемпотентный upsert и сохраняет новую отметку.

Если прогон упал, ретрай идёт с той же позиции курсора, а не с начала. В истории выполнений виден каждый прогон: диапазон отметок, число синхронизированных записей и сбои.

Из чего состоит синхронизация по расписанию

Инкрементальный курсор

Храните курсор updatedAt или порядковый номер для каждой задачи синхронизации. Забирайте только дельту с последнего успеха — экономно по квоте API и быстро по времени.

Идемпотентный upsert

Делайте upsert по стабильному внешнему ID. Повторный прогон того же диапазона курсора даёт то же состояние базы — безопасно ретраить при сбое.

История выполнений по прогонам

Каждый запуск по расписанию создаёт запись выполнения: диапазон курсора, число записей, длительность, успех или сбой. Чтобы ответить «запускалось ли?», SSH не нужен.

Настраиваемые cron-выражения

Каждые 15 минут, раз в час, раз в сутки или своё расписание — с проверкой при сохранении. История хранит все прогоны, даже когда выражение менялось.

Паттерн пайплайна инкрементальной синхронизации

1

Прочитать курсор из хранилища состояния

Получите последнюю успешную отметку из базы или хранилища состояния. Для первого запуска задайте разумное окно «в прошлое».

2

Забрать дельту и сделать upsert

Запросите у внешнего API записи, обновлённые после курсора. Делайте upsert по внешнему ID батчами. Пагинацию обрабатывайте внутри шага.

3

Сохранить новый курсор

При успехе сохраните новый курсор (максимальный updatedAt из батча). Верните статистику для истории выполнений.

Инкрементальная синхронизация с курсором-отметкой

Графовый пайплайн: узел cronTrigger ведёт в узел-функцию — тот же формат, что в консоли и документации. Обработчик синхронизации читает курсор-отметку, забирает дельту, делает идемпотентный upsert и сдвигает курсор только после успешного батча.

pipelines/sync-crm-contacts.json (граф пайплайна)
{
  "schemaVersion": 1,
  "nodes": [
    { "id": "cron", "kind": "cronTrigger", "name": "Hourly", "position": { "x": 0, "y": 0 }, "config": { "cron": "0 * * * *", "timezone": "UTC" } },
    { "id": "sync", "kind": "lambda", "name": "Sync CRM contacts", "position": { "x": 240, "y": 0 }, "config": { "functionId": "sync-crm-contacts", "onError": "failPipeline" } }
  ],
  "edges": [
    { "id": "e1", "sourceNodeId": "cron", "targetNodeId": "sync", "sourceHandle": "default" }
  ]
}
jobs/sync-crm-contacts.mjs (шаг пайплайна)
export async function handler(event) {
  // Read watermark — default to 24h ago on first run
  const cursor = await db.syncCursors.get('crm-contacts') ?? new Date(Date.now() - 86_400_000).toISOString();
  let page = 0, synced = 0, newCursor = cursor;
  do {
    const { contacts, nextPage } = await crm.fetchContacts({ updatedAfter: cursor, page });
    if (contacts.length === 0) break;
    await db.contacts.upsertBatch(contacts); // upsert by contacts[].externalId
    synced += contacts.length;
    newCursor = contacts.at(-1)?.updatedAt ?? newCursor;
    page = nextPage;
  } while (page);
  await db.syncCursors.set('crm-contacts', newCursor);
  return { synced, cursor, newCursor };
}

Используйте синхронизацию по расписанию для

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

  • Данные CRM, ERP или сторонних API, которые нужно периодически обновлять в вашей базе
  • Часовая или суточная синхронизация дельты, когда полная выгрузка слишком медленная или дорогая по квоте

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

  • Синхронизация в реальном времени с задержкой меньше минуты — для событийных обновлений используйте вебхуки вместо cron

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

Как учитывать лимиты внешнего API в задачах синхронизации?

Добавьте паузу между страницами в цикле синхронизации; ловите ответы 429 и ждите перед повторным запросом страницы. Настройте ретраи шага пайплайна с backoff для временных срабатываний лимита.

Что если внешний API не поддерживает инкрементальные запросы?

Забирайте все записи и сравнивайте с базой по хэшу или updatedAt. Для больших наборов используйте staging-таблицу: полный импорт в staging, сравнение с продакшеном, применение дельты.