Сценарий · Inquir Compute

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

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

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

Суть ответа

Синхронизация данных по расписанию с инкрементальным курсором. Пайплайн по расписанию читает сохранённую отметку из базы или другого хранилища состояния, запрашивает у источника только строки с updated_at новее этой отметки, выполняет идемпотентный UPSERT по бизнес-идентификатору и сдвигает отметку только после успешной фиксации батча.

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

  • Источник поддерживает фильтрацию по updated_at или другому курсору
  • Нужна регулярная синхронизация дельты, потому что полный перезапрос слишком долгий или дорогой по квоте

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

  • Сдвигать отметку last_synced_at до того, как транзакция с данными зафиксирована, опасно: при ошибке записи курсор уже «уехал вперёд», часть строк не попала в целевую базу, а расхождение заметят поздно.
  • Без идемпотентной записи — UPSERT по первичному ключу или явной дедупликации — повторный прогон после сбоя или повторная доставка события создают дубликаты.

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

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

Если на каждом прогоне вытягивать полный снимок данных, синхронизация становится медленной и дорогой и перегружает источник; с ростом объёма один прогон может длиться дольше, чем интервал между запусками cron.

Типичные ошибки при инкрементальной синхронизации

Сдвигать отметку last_synced_at до того, как транзакция с данными зафиксирована, опасно: при ошибке записи курсор уже «уехал вперёд», часть строк не попала в целевую базу, а расхождение заметят поздно.

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

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

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

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

Из чего складывается устойчивая инкрементальная синхронизация

Курсор состояния

Храните отметку последней успешной синхронизации (last_synced_at или эквивалент) в базе или key-value и запрашивайте у источника только строки новее этой отметки — экономнее по квоте API и быстрее по времени.

Идемпотентная вставка или обновление (UPSERT)

Используйте конструкцию вида INSERT … ON CONFLICT (id) DO UPDATE … или эквивалент в вашей СУБД, чтобы повторный прогон с теми же строками не создавал дубликатов.

Атомарный сдвиг курсора

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

Расписание cron в пайплайне

Часовой или суточный прогон задаётся cron-выражением, которое валидируется при сохранении; в консоли остаются история запусков и политика ретраев при временных сбоях.

Как устроен инкрементальный прогон синхронизации

1

Прочитать курсор

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

2

Запросить дельту

Инкрементальный запрос вида WHERE updated_at > last_synced_at ORDER BY updated_at LIMIT N забирает у источника только изменения после курсора.

3

Upsert и сдвиг курсора

Сделайте UPSERT батча и обновите last_synced_at = MAX(updated_at) в одной транзакции.

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

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

pipelines/sync-crm-contacts.json (pipeline graph)
{
  "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 (pipeline step)
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 };
}

Когда выбирать инкрементальную синхронизацию

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

  • Источник поддерживает фильтрацию по updated_at или другому курсору
  • Нужна регулярная синхронизация дельты, потому что полный перезапрос слишком долгий или дорогой по квоте

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

  • Источник не поддерживает фильтрацию по времени — понадобится полная выгрузка со сравнением по хэшу

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

Что если прогон упал на середине?

Следующий запуск по расписанию снова прочитает прежнюю отметку last_synced_at, и те же строки запросятся повторно. При идемпотентном UPSERT дубликаты не появятся.

Где хранить состояние курсора?

Обычно в отдельной таблице sync_state или в key-value-хранилище — по одной записи на каждый источник и тип синхронизации.