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
Инкрементальные пайплайны с состоянием курсора
Пайплайны Inquir по расписанию запускаются как serverless-функции с историей выполнений. Обработчик синхронизации читает последнюю успешную отметку (из базы или другого хранилища состояния), забирает только записи, обновлённые после этого курсора, делает идемпотентный upsert и сохраняет новую отметку.
Если прогон упал, ретрай идёт с той же позиции курсора, а не с начала. В истории выполнений виден каждый прогон: диапазон отметок, число синхронизированных записей и сбои.
Что вы получаете
Из чего состоит синхронизация по расписанию
Инкрементальный курсор
Храните курсор updatedAt или порядковый номер для каждой задачи синхронизации. Забирайте только дельту с последнего успеха — экономно по квоте API и быстро по времени.
Идемпотентный upsert
Делайте upsert по стабильному внешнему ID. Повторный прогон того же диапазона курсора даёт то же состояние базы — безопасно ретраить при сбое.
История выполнений по прогонам
Каждый запуск по расписанию создаёт запись выполнения: диапазон курсора, число записей, длительность, успех или сбой. Чтобы ответить «запускалось ли?», SSH не нужен.
Настраиваемые cron-выражения
Каждые 15 минут, раз в час, раз в сутки или своё расписание — с проверкой при сохранении. История хранит все прогоны, даже когда выражение менялось.
Что дальше
Паттерн пайплайна инкрементальной синхронизации
Прочитать курсор из хранилища состояния
Получите последнюю успешную отметку из базы или хранилища состояния. Для первого запуска задайте разумное окно «в прошлое».
Забрать дельту и сделать upsert
Запросите у внешнего API записи, обновлённые после курсора. Делайте upsert по внешнему ID батчами. Пагинацию обрабатывайте внутри шага.
Сохранить новый курсор
При успехе сохраните новый курсор (максимальный updatedAt из батча). Верните статистику для истории выполнений.
Пример кода
Инкрементальная синхронизация с курсором-отметкой
Графовый пайплайн: узел cronTrigger ведёт в узел-функцию — тот же формат, что в консоли и документации. Обработчик синхронизации читает курсор-отметку, забирает дельту, делает идемпотентный upsert и сдвигает курсор только после успешного батча.
{ "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" } ] }
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, сравнение с продакшеном, применение дельты.