DATAREON Platform под нагрузкой: настройка высоконагруженной интеграции DATAREON + PostgreSQL + REST

Всем привет! Я Дмитрий Пономарев, разработчик ESB ИТ-интегратора «Белый код»! Сначала всё выглядело довольно просто: забираем данные по REST API, обрабатываем и сохраняем в PostgreSQL. Проблемы начались, когда объём таблицы вырос до десятков миллионов строк, а несколько миллионов записей пришлось регулярно проверять заново.
В одном из проектов мы столкнулись именно с такой задачей. DATAREON Platform должен был получать данные через REST API, сохранять их в PostgreSQL, хранить историю состояний и каждый день повторно проверять актуальные записи. На небольших объёмах схема работала нормально, но после роста данных стало ясно: последовательная обработка и вычисление актуальной версии записи «на лету» больше не укладываются в рабочее окно.
В этой статье разберём, как мы перестроили интеграцию: вынесли актуальность записи в отдельный признак, добавили специализированные индексы PostgreSQL, распараллелили обработку через воркеры DATAREON и предусмотрели безопасное восстановление после сбоев. В итоге производительность промышленной актуализации удалось довести примерно до 1,1–1,2 млн объектов в час.
Задача
В рамках проекта нужно было организовать регулярную проверку кодов маркировки в системе «Честный знак».
В PostgreSQL хранится большая таблица с кодами идентификации товаров и информацией об их текущем состоянии. Новые коды регулярно поступают в таблицу, после чего DATAREON Platform должна обратиться во внешний REST API и получить по каждому коду актуальные данные.
Проверка выполняется в двух случаях:
для новых кодов, которые ещё ни разу не проверялись;
для уже обработанных кодов, состояние которых необходимо периодически актуализировать.
Во втором случае важно учитывать, что статус кода со временем может измениться. Например, товар был введён в оборот, а позже выбыл из него. Поэтому недостаточно один раз получить информацию и сохранить её в базе. Актуальные записи необходимо регулярно перепроверять через API.
DATAREON выбирает очередную пачку кодов из PostgreSQL, отправляет их во внешний REST API, разбирает ответ и обновляет данные в таблице. Для новых записей выполняется первичная проверка, а для уже обработанных записей выполняется повторная актуализация.

При этом объём данных оказался довольно большим: в таблице хранятся десятки миллионов записей, а несколько миллионов актуальных кодов необходимо регулярно перепроверять. Поэтому обычной последовательной обработки здесь уже недостаточно.
Для обращения к API также требуется авторизация. DATAREON получает служебный токен через отдельный endpoint, сохраняет его в Банке данных и использует при последующих запросах. Токен имеет ограниченный срок действия, поэтому интеграция должна контролировать validTo и получать новый токен после истечения текущего.
Почему обычная схема перестала справляться
Изначально для каждой записи нужно было понять, является ли она последней версией объекта. Для этого PostgreSQL искал максимальный lifecycle_num (номер жизненного цикла записи) или проверял отсутствие более новой строки. Логика корректная, но при большом объёме данных она заставляет базу выполнять огромное количество дополнительных индексных обращений.
На таблице порядка 90 млн строк выборка одной пачки из 1000 объектов могла занимать больше 50 секунд. При необходимости ежедневно обработать несколько миллионов записей математика становилась довольно грустной: даже стабильный и быстрый REST API уже не спасал, потому что узким местом была сама подготовка очереди.
Нам нужно было решить две задачи одновременно: быстро определять актуальные записи и безопасно раздать их нескольким параллельным воркерам.
Архитектура решения
После переработки поток разделили на два режима: первичную обработку новых записей и повторную актуализацию уже известных объектов. Оба режима используют одинаковый принцип пакетной обработки.

Диспетчер запускает 10 параллельных воркеров. Каждый воркер самостоятельно выбирает из PostgreSQL очередную пачку до 1000 записей. При выборке используется механизм блокировок PostgreSQL FOR UPDATE SKIP LOCKED: выбранные строки блокируются на время обработки, а другие воркеры пропускают уже занятые записи и забирают следующую свободную пачку. Благодаря этому несколько воркеров могут одновременно работать с одной таблицей, не обрабатывая одну и ту же запись дважды. После отбора выбранные записи переводятся в техническое состояние PROCESSING, затем их идентификаторы отправляются во внешний REST API. После получения ответа результат фиксируется в PostgreSQL.
Для новых записей выбираются записи в статусе NEW. Для повторной проверки используются только актуальные записи со статусом PROCESSED, которые ещё не проверялись в текущий день. В нашем случае повторная актуализация запускается ежедневно, поэтому одна и та же актуальная запись должна проверяться не чаще одного раза в сутки. После наступления следующего дня она снова может попасть в очередь на проверку. Записи, при обработке которых произошла ошибка, переводятся в статус ERROR и логируются отдельно.
В DATAREON параллельная обработка была реализована через отдельную схему-диспетчер.
В ней используются две переменные:
ThreadsCount — сколько воркеров нужно запустить одновременно;
Thread — сколько воркеров уже запущено.

Диспетчер работает по простому принципу: пока значение Thread меньше ThreadsCount, он запускает ещё один экземпляр основной схемы обработки данных.

Шаг «Воркер» запускает основную схему, которая самостоятельно забирает свою пачку записей из PostgreSQL и обрабатывает её.
Шаг «Добавление воркера» увеличивает значение Thread на 1. Затем диспетчер снова проверяет условие. Если запущено меньше процессов, чем указано в ThreadsCount, запускается следующий воркер.
Убираем тяжёлый поиск актуальной записи
Ключевое изменение было простым: вместо того чтобы каждый раз вычислять последнюю версию объекта, мы начали хранить её явно.

Поле lifecycle_num при этом осталось. Оно отвечает за историю состояний, а is_current отвечает только на один вопрос: какую строку считать текущей прямо сейчас.
Для защиты добавили частичный уникальный индекс чтобы для одного объекта не могло существовать две PROCESSED-записи с is_current = true.
Отдельный индекс для Repeat
Повторная актуализация идёт по актуальным записям, которые не проверялись сегодня. Поэтому под этот сценарий сделали отдельный частичный индекс по полю updated_at. В результате PostgreSQL больше не перебирает историю каждого объекта. Он сразу работает с небольшой выборкой относительно всей таблицы актуальных строк.
Как сохраняем историю изменений
Ответ REST API может подтвердить текущее состояние объекта или вернуть новое. В первом случае мы просто обновляем техническое состояние и дату проверки. Во втором сохраняем предыдущую строку как историческую и создаём или активируем новое состояние.

Так можно быстро получить текущий статус и при этом не терять историю. Важный момент: lifecycle_num больше не используется для определения актуальности. Источник истины для текущей записи — is_current.
Подводный камень: общий REST-токен и 10 воркеров
После распараллеливания проявилась ещё одна проблема, связанная уже не с PostgreSQL, а с хранением токена доступа. Для обращения к REST API используется отдельный токен авторизации. DATAREON получает его через служебный endpoint getToken и в ответ API возвращает рабочий токен и срок его действия validTo. Полученный токен сохраняется в Банке данных DATAREON в отдельной служебной записи.
Перед отправкой очередного запроса воркер читает эту запись и проверяет validTo. Пока срок действия не истёк, используется уже сохранённый токен. После наступления validTo DATAREON повторно вызывает getToken, получает новый токен и новый срок действия, после чего обновляет ту же запись в Банке данных.
Изначально каждый воркер мог самостоятельно проверить наличие записи токена и при её отсутствии создать новую. При одновременном запуске нескольких воркеров возникала ситуация, что несколько процессов почти одновременно получали результат «запись не найдена» и каждый создавал собственную служебную запись токена в Банке данных. В результате в Банке данных появлялись дубли токенов.
Проблему решили отказом от динамического создания записи. Для токена заранее создаётся одна служебная запись с фиксированным EntityId. Воркеры больше не создают новые записи, а только читают существующую и при необходимости обновляют в ней DataValue и ValidTo и только по EntityId.
Таким образом, все параллельные процессы работают с одним экземпляром токена, а после его истечения обновляется та же самая запись, а не создаётся новая.
Скрытый текст
Правило: общий токен должен иметь singleton-семантику. Запись создаётся один раз, а воркеры имеют право только читать и обновлять её.
Результат

После перехода на is_current, частичные индексы и параллельную обработку через 10 воркеров промышленная скорость Repeat-потока составила порядка 1,1–1,2 млн объектов в час.
История состояний при этом сохранилась: старые версии записи остаются в таблице, а текущая версия определяется признаком is_current.
Дополнительно в PostgreSQL был создан частичный уникальный индекс, который гарантирует, что для одного объекта не может одновременно существовать две записи со статусом INTRODUCED, record_status = 'PROCESSED' и is_current = true.
Рост производительности получился не за счёт одного индекса. Основной эффект дала комбинация нескольких изменений:
хранение признака актуальности непосредственно в is_current;
частичные индексы под конкретные очереди обработки;
выборка пачек через FOR UPDATE SKIP LOCKED;
пакетная обработка по 1000 объектов;
параллельный запуск 10 воркеров.
Заключение
Высоконагруженная интеграция редко начинается как высоконагруженная. Сначала это обычный REST-запрос и INSERT в таблицу, а проблемы появляются позже — вместе с миллионами записей, историей состояний и несколькими параллельными процессами.
В нашем случае основными точками роста стали четыре решения: хранить актуальность явно, использовать частичные индексы под конкретные очереди, раздавать пачки через FOR UPDATE SKIP LOCKED и относиться к общим техническим сущностям вроде токена как к singleton-объектам.
В результате DATAREON остался оркестратором интеграции, PostgreSQL взял на себя эффективную работу с очередью и конкурентным доступом, а REST API — только свою прямую задачу: вернуть актуальное состояние объектов. Такая граница ответственности хорошо масштабируется и остаётся достаточно простой для сопровождения.
Если перед вами стоит похожая задача на DATAREON — массовая актуализация данных, параллельные воркеры, оптимизация PostgreSQL или нестандартная логика работы с REST API — такой подход можно адаптировать под конкретный интеграционный контур.
Если перед вами стоят задачи по развитию интеграционного контура на DATAREON — от доработки существующих потоков до реализации нестандартных сценариев, — обращайтесь за консультацией или проработкой решения.
KioskNews shows a cleaned-up reading view extracted from the publisher’s page — the original always lives on their site, not ours.