Привет, Хабр!
В этой статье я хочу подробно разобрать практический пример инкрементальной синхронизации данных между двумя базами PostgreSQL с использованием FESB.
Материал получился достаточно объёмным, поскольку я решил показать не только общую идею, но и весь путь реализации: от подготовки таблиц и триггеров до настройки интеграционного процесса, первоначальной загрузки и обработки последующих изменений.
Статья в первую очередь будет полезна разработчикам, интеграторам и архитекторам, которые работают с ESB, ETL и реляционными базами данных. Даже если вы не используете FESB, описанный подход с контрольной точкой и инкрементальной загрузкой можно адаптировать для других интеграционных платформ.
Надеюсь, каждый читатель сможет найти в материале что-то полезное для своих задач — готовое решение, отдельный SQL-приём или идею для построения собственного интеграционного процесса.
1. Теоретическая часть
1.1 Задача синхронизации данных
При интеграции информационных систем часто требуется регулярно переносить данные из одной базы данных в другую. Например, необходимо синхронизировать справочник сотрудников между операционной базой и отдельной базой, используемой интеграционным или аналитическим контуром.
Самый простой способ решить такую задачу — при каждом запуске полностью читать исходную таблицу и заново передавать все записи в систему-получатель. Однако при увеличении объёма данных такой подход становится неэффективным:
возрастает нагрузка на исходную базу данных;
увеличивается объём передаваемых данных;
растёт продолжительность интеграционного процесса;
чаще возникают блокировки и конкуренция за ресурсы;
повторная обработка неизменившихся данных не приносит практической пользы.
Поэтому вместо полной загрузки обычно используют инкрементальную загрузку — обработку только тех записей, которые были добавлены или изменены после последней успешно сохраненной контрольной точки
В этой статье будет рассмотрен вариант инкрементальной загрузки данных из PostgreSQL в PostgreSQL с использованием FESB. Для определения изменений применяются временные поля Dt_Ins и Dt_Upd, а состояние последней обработки сохраняется в отдельной таблице cdc_tracking.
1.2 Что такое CDC
Change Data Capture (CDC) — это общий подход к выявлению изменений в источнике данных и передаче этих изменений в другие системы.
Изменением может считаться:
добавление новой записи;
изменение существующей записи;
удаление записи;
иногда — изменение состояния или статуса объекта.
Существует несколько способов реализации CDC:
чтение журнала транзакций базы данных;
использование логического декодирования PostgreSQL и WAL;
применение триггеров;
ведение отдельной таблицы изменений;
использование временных отметок;
сравнение текущего состояния с предыдущим снимком данных.
В рамках данного материала используется наиболее распространенный прикладной вариант — определение изменений по временной отметке. Он не читает журнал транзакций PostgreSQL, а опирается на служебные поля, которые хранят дату создания и дату последнего изменения записи.
Поэтому корректнее воспринимать описанный подход как инкрементальную загрузку по timestamp или CDC на основе watermark, а не как полноценный потоковый CDC на уровне журнала транзакций.
1.3 Преимущества подхода
Инкрементальная загрузка по временной отметке имеет несколько важных преимуществ.

Простота реализации
Для работы метода не требуется отдельная брокерная система, сложная инфраструктура или чтение журнала транзакций. Достаточно:
добавить служебные поля в таблицу;
настроить триггеры;
создать таблицу контрольных точек;
реализовать периодический SQL-запрос.
Низкий порог внедрения
Подход можно использовать даже в проектах, где исходная система не предоставляет готового механизма CDC. Если есть возможность изменить структуру таблицы или добавить триггеры, решение обычно внедряется достаточно быстро.
Снижение нагрузки
При последующих запусках обрабатывается только изменившаяся часть данных. Это уменьшает:
количество строк, прочитанных из источника;
объём данных, переданных через интеграционный контур;
число операций записи в целевой системе;
время выполнения процесса.
Независимость от формата данных
Метод работает с обычными строками таблицы и не требует преобразования записей в специальные события. Это удобно для SQL-интеграций и ETL-процессов.
Идемпотентная запись
При использовании INSERT ... ON CONFLICT повторная доставка записи не приводит к созданию дубликата. Это особенно важно, поскольку интеграционные процессы обычно работают по модели «как минимум одна доставка» и могут повторять обработку после ошибок.
Наглядность
Состояние синхронизации хранится в обычной таблице. Администратор или разработчик может быстро проверить текущую контрольную точку обычным SQL-запросом
2. Архитектура решения

В рамках данного кейса реализуется процесс синхронизации таблицы employees между двумя базами данных PostgreSQL.
Компоненты системы:
Источник данных (fesb_source_db):
employees: Основная бизнес-таблица с данными о сотрудниках.cdc_tracking: Служебная таблица для отслеживания состояния синхронизации. Содержит уникальный ключ потока (integration_flow) и отметкуcheck_update.Триггеры и функции: Обеспечивают автоматическое заполнение полей
dt_insиdt_updпри вставке новых записей, а также изменение уже существующих.
Оркестратор ETL (FESB):
Управляет всем процессом синхронизации.
Выполняет бизнес-логику.
Целевая система (fesb_integration_db):
employees: Таблица-получатель. Может иметь схожую, но не обязательно идентичную структуру (например,dt_updможет иметь значение по умолчанию).
Логика потока данных:После инициации процесса в FESB выполняется следующая последовательность действий:
Производится обращение к таблице
cdc_trackingв базе-источнике для получения значения столбцаcheck_updateдля конкретного потока интеграции, в нашем случае, "employees".Возможны два сценария:
Сценарий 1: Возвращается пустое значение. Это означает, что ранее никаких обращений не производилось. В этом случае все записи из таблицы
employeesпомещаются в таблицу-получатель.Сценарий 2: Возвращается значение
check_update(не пустое) и название нужного потока. Это означает, что ранее производилось обращение к таблице-источнику. В этом случае из таблицыemployeesизвлекаются и отправляются в таблицу-получатель только те записи, у которых значениеdt_insилиdt_updбольше полученного значенияcheck_update.
В обоих сценариях, после успешной обработки пакета данных, в нем находятся значения столбцов
dt_insиdt_upd. Максимальное из этих значений помещается обратно в таблицуcdc_trackingв столбецcheck_update, обновляя контрольную точку для следующего запуска.
3. Подготовка базы данных: Скрипты таблиц и триггеров
3.1 База-источник (fesb_source_db)
Таблица employees (основная бизнес-таблица)
CREATE TABLE IF NOT EXISTS public.employees ( "EmployeeId" integer NOT NULL, "LastName" character varying(255) COLLATE pg_catalog."default" NOT NULL, "FirstName" character varying(255) COLLATE pg_catalog."default" NOT NULL, "Title" character varying(255) COLLATE pg_catalog."default", "ReportsTo" integer, "BirthDate" timestamp with time zone, "HireDate" timestamp with time zone, "Address" character varying(255) COLLATE pg_catalog."default", "City" character varying(255) COLLATE pg_catalog."default", "State" character varying(255) COLLATE pg_catalog."default", "Country" character varying(255) COLLATE pg_catalog."default", "PostalCode" character varying(50) COLLATE pg_catalog."default", "Phone" character varying(50) COLLATE pg_catalog."default", "Fax" character varying(50) COLLATE pg_catalog."default", "Email" character varying(255) COLLATE pg_catalog."default", "Dt_Ins" timestamp with time zone DEFAULT CURRENT_TIMESTAMP, -- Автоматически при вставке "Dt_Upd" timestamp with time zone, -- Обновляется триггером CONSTRAINT employees_pkey PRIMARY KEY ("EmployeeId") );
Функции для таблицы employees
CREATE OR REPLACE FUNCTION public.update_dt_ins_on_insert() RETURNS trigger LANGUAGE 'plpgsql' COST 100 VOLATILE NOT LEAKPROOF AS $BODY$ BEGIN NEW."Dt_Ins" = CURRENT_TIMESTAMP; RETURN NEW; END; $BODY$; ALTER FUNCTION public.update_dt_ins_on_insert() OWNER TO postgres;
CREATE OR REPLACE FUNCTION public.update_dt_upd_on_update() RETURNS trigger LANGUAGE 'plpgsql' COST 100 VOLATILE NOT LEAKPROOF AS $BODY$ BEGIN IF (NEW.* IS DISTINCT FROM OLD.*) THEN NEW."Dt_Upd" = CURRENT_TIMESTAMP; END IF; RETURN NEW; END; $BODY$; ALTER FUNCTION public.update_dt_upd_on_update() OWNER TO postgres;
Триггеры для таблицы employees
Предполагается, что в базе существуют функции update_dt_ins_on_insert() и update_dt_upd_on_update().
-- Триггер на вставку (INSERT) CREATE OR REPLACE TRIGGER trg_employees_before_insert BEFORE INSERT ON public.employees FOR EACH ROW EXECUTE FUNCTION public.update_dt_ins_on_insert(); -- Триггер на обновление (UPDATE) CREATE OR REPLACE TRIGGER trg_employees_before_update BEFORE UPDATE ON public.employees FOR EACH ROW EXECUTE FUNCTION public.update_dt_upd_on_update();
Таблица cdc_tracking (контрольная точка)
CREATE TABLE IF NOT EXISTS public.cdc_tracking ( integration_flow character varying(100) COLLATE pg_catalog."default" NOT NULL, check_update timestamp with time zone DEFAULT CURRENT_TIMESTAMP, CONSTRAINT cdc_tracking_pkey PRIMARY KEY (integration_flow) );
3.2 База-получатель (fesb_integration_db)
Таблица employees
Структура может быть адаптирована. Обратите внимание, что Dt_Upd здесь имеет значение по умолчанию.
CREATE TABLE IF NOT EXISTS public.employees ( "EmployeeId" integer NOT NULL, "LastName" character varying(255) COLLATE pg_catalog."default" NOT NULL, "FirstName" character varying(255) COLLATE pg_catalog."default" NOT NULL, "Title" character varying(255) COLLATE pg_catalog."default", "ReportsTo" integer, "BirthDate" timestamp with time zone, "HireDate" timestamp with time zone, "Address" character varying(255) COLLATE pg_catalog."default", "City" character varying(255) COLLATE pg_catalog."default", "State" character varying(255) COLLATE pg_catalog."default", "Country" character varying(255) COLLATE pg_catalog."default", "PostalCode" character varying(50) COLLATE pg_catalog."default", "Phone" character varying(50) COLLATE pg_catalog."default", "Fax" character varying(50) COLLATE pg_catalog."default", "Email" character varying(255) COLLATE pg_catalog."default", "Dt_Ins" timestamp with time zone DEFAULT CURRENT_TIMESTAMP, "Dt_Upd" timestamp with time zone DEFAULT CURRENT_TIMESTAMP, CONSTRAINT employees_pkey PRIMARY KEY ("EmployeeId") );
4. Практика
4.1 Добавление библиотеки
Для начала необходимо добавить библиотеку для работы с Базой Данных PostgreSQL. Для этого необходимо перейти в "Настройки" → "Управление библиотеками"

После перехода вы попадете в интерфейс загрузки библиотек, нам необходимо поместить сюда соответствующую библиотеку, которую можно получить в открытом доступе в интернете.
В статье используется следующая библиотека для работы с PostgreSQL

4.2 Создание брокера
После того, как библиотека была успешно загружена, нам необходимо создать домен брокера для настройки дальнейшей бизнес-логики.
Переходим в "Интеграционный брокер" → "Домены брокера" → Создать домен

Имя домену можно указать любое, главное, чтобы в будущей работе не заблудиться среди множества доменов.
Как только мы создали домен брокера, нам необходимо добавить настройки подключения к базам источника и получателя.
Переходим в раздел домена "Дополнительная конфигурация"

И добавляем конфигурации наших БД, настройки подключения стандартные, пример:
Тип объекта | Источник данных |
Тип | Обычный |
Имя | Имя Вашей Базы Данных (fesb_integration_db/fesb_source_db) |
Класс JDBC драйвера | org.postgresql.Driver |
URL базы данных | jdbc:postgresql://localhost:5432/fesb_integration_db |
Имя пользователя | Имя Вашего пользователя По умолчанию: postgres |
Пароль | Пароль Вашего пользователя По умолчанию: postgres |
Делаем 2 конфигурации в данном разделе, для базы-источника и базы-получателя.
4.3 СОПС Информация по потоку есть
Что такое СОПС
СОПС - Схема Обработки Потока Сообщений.
Если кратко - это пошаговый сценарий, который описывает, что именно шина FESB должна сделать с поступившим сообщением.
Как только все приготовления были завершены, мы можем перейти к настройке СОПС.
Для этого переходим в раздел домена брокера "Список СОПС" и создаем его

Пока что просто создаем СОПС с именем "Информация по потоку есть" и идем дальше. Мы к нему еще вернемся.
Чтобы не ругалось на отсутствие настроек в СОПС, создайте точку входа "Ссылка на СОПС" и добавьте любую обработку, например, лог.
4.4 СОПС Информация по потоку отсутствует
Пока что просто создаем СОПС с именем "Информация по потоку отсутствует" и идем дальше. Мы к нему еще вернемся.
Чтобы не ругалось на отсутствие настроек в СОПС, создайте точку входа "Ссылка на СОПС" и добавьте любую обработку, например, лог.
4.5 СОПС Получение информации по потоку - первый и основной СОПС на запуск
Шаг "Запуск СОПСа" (Таймер - Точка входа)
Первым шагом будет таймер, в котором прописывается, как часто мы должны запускать основной СОПС по проверке даты последнего обращения к таблице cdc_tracking.

Соответственно, по настройкам видно, что СОПС будет запускаться каждые 30 секунд до тех пор, пока не будет остановлен.
Шаг "SQL"
Следующим шагом нам необходимо обратиться к базе данных "fesb_source_db", а именно к таблице cdc_tracking.

Код запроса:
SELECT * FROM cdc_tracking WHERE integration_flow = 'employees'
В данном запросе мы обращаемся к таблице cdc_tracking и пытаемся получить запись, у которой integration_flow = 'employees'
Шаг "Логирование"
Чтобы проще было понимать следующий шаг "Фильтр", посмотрим тело сообщения, которое приходит нам ответом после выполнения запроса в SQL.
Сообщение логирования:
${body}
После настройки этого шага мы можем запустить СОПС и посмотреть тело сообщения в "Протокол работы".
В результате, мы увидим следующий лог:
[PostgreSQL_Отметка_По_Дате_Изменения/Получение информации по потоку] - [{integration_flow=employees, check_update=2026-02-16 16:47:09.878447}]
Все что в "[]" - это и есть наше тело сообщения, полученного от SQL
Шаг "Фильтр"

JSONPath
JSONPath — это язык запросов для анализа и извлечения данных из JSON‑структур. Если провести аналогию, то это как XPath для XML или SQL‑запросы для таблиц, но только для JSON‑объектов.
Давайте разберем наше выражение в jsonpath:
$[?(@.integration_flow == 'employees')]
$(Корень): Обозначает начало документа (весь ваш JSON‑объект, который мы увидели в логе выше)
[...](Оператор индекса/скрипта): Указывает на то, что мы применяем операцию к текущему уровню данных.
?(...)(Фильтр): Это команда поиска. Она говорит: «Пройдись по элементам текущего уровня данных и оставь только те, что соответствуют условию в скобках».
@(Текущий объект): Ссылка на элемент, который проверяется прямо сейчас.
.integration_flow: Обращение к конкретному полю внутри объекта.
== 'employees': Оператор сравнения. Мы ищем точное совпадение строкового значения.Подводя итог, наше выражение ищет по всему массиву значение: integration_flow=employees
На данном шаге строится условие. Это условие будет проверять, существует ли изначально в таблице cdc_tracking интеграционный поток под названием "employees".
С этого момент бизнес-логика разбивается на 2 части:
В таблице cdc_tracking интеграционный поток под названием "employees" присутствует
В таблице cdc_tracking интеграционный поток под названием "employees" отсутствует
Настройки фильтра
Основное условие

Это условие проверяет наличие интересующего нас интеграционного потока в полученном теле сообщения после SQL
Текст условия:
$[?(@.integration_flow == 'employees')]
Если условие выполняется, то есть мы находим запись с нужным интеграционным потоком, то мы переходим в созданный ранее СОПС "Информация по потоку есть".
Переход осуществляется при помощи шага "Ссылка на СОПС"

Если условие НЕ выполняется, то есть мы не находим запись с нужным интеграционным потоком, то мы переходим в созданный ранее СОПС "Информация по потоку отсутствует".
Переход осуществляется при помощи шага "Ссылка на СОПС"

5. В таблице cdc_tracking интеграционный поток под названием "employees" отсутствует
5.1 Переходим в СОПС "Информация по потоку отсутствует"
Шаг "Ссылка на СОПС"
На данном шаге точкой входа будет являться основной СОПС "Получение информации по потоку". Соответственно, каждый раз при не выполнении условия мы будем попадать в данный СОПС.
На практике же, мы перейдем в текущий СОПС только единожды, поскольку во все следующие раз у нас уже будет присутствовать запись по интеграционному потоку.

Шаг "SQL"

Максимально простой шаг. Поскольку в текущей ветке бизнес-логики мы понимаем, что еще ни разу не получали данные от базы-источника, значит нам необходимо выгрузить абсолютно все данные из таблицы employees.
Код запроса
SELECT * FROM employees
Шаг "Фильтр"
Хоть в текущем СОПС мы и подразумеваем, что необходимо забрать все записи. Логически лучше все-таки поставить проверку на пустое тело сообщения.
По такой логике, если таблица-источник еще никак не заполнена, СОПС дальше не пойдет по шагам.

Условие шага
${body} != '[]'
Если условие выполняется, то мы идем дальше.
Если условие не выполняется, мы логируем это событие

Нет записей для обновления по потоку employees. Текущий чекпоинт: ${header.last_sync}
И завершаем дальнейшую обработку.

Шаг "SQL"

После того, как запрос успешно выполнился, и мы получили все записи из таблицы employees-источника, нам необходимо их записать в таблицу employees-получателя.
Дополнительно, на текущем шаге необходимо добавить параметр "Групповое выполнение" и установить его значение в true. Это делается в самом низу шага

Групповое выполнение (batch):
Если в дополнительных параметрах выбрано
Групповое выполнение (batch), то интерпретация тела входящего сообщения немного меняется — вместо итератора параметров компонент ожидает итератор, содержащий итераторы параметров. Размер внешнего итератора определяет размер пакета.Результат запроса
Для операций
SELECTрезультат возвращается какList>(результат методаJdbcTemplate.queryForList()).Для операций
UPDATEтело сообщения остаетсяNULL, а результат операции доступен в заголовках.По умолчанию результат помещается в тело сообщения. Если установлен параметр
Записать результат в заголовок (outputHeader), результат будет записан в указанный заголовок. Это альтернатива использованию полного шаблона обогащения сообщения для добавления заголовков, он обеспечивает краткий синтаксис для запроса последовательности или другого небольшого значения в заголовок. Удобно использоватьoutputHeaderиoutputTypeвместе.
Код запроса
INSERT INTO public.employees ( "EmployeeId", "LastName", "FirstName", "Title", "ReportsTo", "BirthDate", "HireDate", "Address", "City", "State", "Country", "PostalCode", "Phone", "Fax", "Email", "Dt_Ins", "Dt_Upd" ) VALUES ( :?EmployeeId, :?LastName, :?FirstName, :?Title, :?ReportsTo, :?BirthDate, :?HireDate, :?Address, :?City, :?State, :?Country, :?PostalCode, :?Phone, :?Fax, :?Email, :?Dt_Ins, :?Dt_Upd ) ON CONFLICT ("EmployeeId") DO UPDATE SET "LastName" = EXCLUDED."LastName", "FirstName" = EXCLUDED."FirstName", "Title" = EXCLUDED."Title", "ReportsTo" = EXCLUDED."ReportsTo", "BirthDate" = EXCLUDED."BirthDate", "HireDate" = EXCLUDED."HireDate", "Address" = EXCLUDED."Address", "City" = EXCLUDED."City", "State" = EXCLUDED."State", "Country" = EXCLUDED."Country", "PostalCode" = EXCLUDED."PostalCode", "Phone" = EXCLUDED."Phone", "Fax" = EXCLUDED."Fax", "Email" = EXCLUDED."Email", "Dt_Upd" = EXCLUDED."Dt_Upd" WHERE ( public.employees."LastName", public.employees."FirstName", public.employees."Title", public.employees."ReportsTo", public.employees."BirthDate", public.employees."HireDate", public.employees."Address", public.employees."City", public.employees."State", public.employees."Country", public.employees."PostalCode", public.employees."Phone", public.employees."Fax", public.employees."Email" ) IS DISTINCT FROM ( EXCLUDED."LastName", EXCLUDED."FirstName", EXCLUDED."Title", EXCLUDED."ReportsTo", EXCLUDED."BirthDate", EXCLUDED."HireDate", EXCLUDED."Address", EXCLUDED."City", EXCLUDED."State", EXCLUDED."Country", EXCLUDED."PostalCode", EXCLUDED."Phone", EXCLUDED."Fax", EXCLUDED."Email" );
Объяснение синтаксиса запроса
Этот синтаксис не является стандартным для работы с PostgreSQL, он скорее модифицирован под FESB и под Apache Camel.
В стандартном JDBC используются порядковые знаки вопроса (
?). Camel позволяет использовать имена. Синтаксис:?ParameterNameсообщает Camel: «Найди в Body (если это Map) или в Headers значение с ключом 'ParameterName' и подставь его сюда».Map - это структура данных, которая хранит значения в формате «Ключ — Значение».
Шаг "Программная трансформация"

После успешной записи данных в таблицу-получатель нам необходимо определить, какая запись из полученных имела самое большое время dt_ins или dt_upd.
Код трансформации:
var result = []; if (typeof body !== 'undefined' && body !== null && Array.isArray(body) && body.length > 0) { var maxDate = null; for (var i = 0; i < body.length; i++) { var row = body[i]; result.push(row); var rowDate = row.Dt_Upd || row.Dt_Ins; if (rowDate) { if (!maxDate || new Date(rowDate) > new Date(maxDate)) { maxDate = rowDate; } } } if (typeof headers !== 'undefined' && maxDate) { headers.max_dt = maxDate; headers.new_checkpoint = maxDate; } } result;
Шаг "Установить заголовки"

На данном шаге нам необходимо установить заголовки максимального значения времени.
Хоть на предыдущем шаге в трансформации мы и выполнили присвоение, в большинстве low-code ESB решениях изменения заголовков внутри трансформаций не всегда автоматически пробрасываются в основной контекст выполнения.
Таким образом, можно считать. что в трансформации мы подготовили данные, а в шаге "Установить заголовки" мы уже регистрируем заголовок, чтобы он не потерялся в бизнес-логике.
Шаг "SQL"

Последний шаг в текущем СОПСе. Как только мы получили максимальное временное значение, мы можем записать в таблице cdc_tracking интеграционный поток с именем "employees", и добавить к нему максимальное значение времени, полученное ранее.
Код запроса:
INSERT INTO public.cdc_tracking (integration_flow, check_update) VALUES ('employees', :?new_checkpoint::timestamptz) ON CONFLICT (integration_flow) DO UPDATE SET check_update = EXCLUDED.check_update;
6. В таблице cdc_tracking интеграционный поток под названием "employees" присутствует
6.1 Переходим в СОПС "Информация по потоку есть"
Шаг "Ссылка на СОПС"

На данном шаге точкой входа будет являться основной СОПС "Получение информации по потоку". Соответственно, каждый раз при выполнении условия мы будем попадать в данный СОПС.
Шаг "Логирование"

Для простоты понимания дальнейших действий выведем один из логов, которые не были показаны в СОПС.
При помощи данного шага мы сможем увидеть, какой результат запроса пришел после обращения к таблице cdc_tracking в основном СОПС.
В следующем шаге мы уже обратимся к телу сообщения
Шаг "Установить заголовки"

На данном шаге мы обращаемся к исходному телу сообщения, которое получили после обращения к таблице cdc_tracking.
Нас интересует время последнего обращения к таблице, поэтому в заголовок мы помещаем значение check_update
Шаг "SQL"

После получения времени последнего обращения, мы можем получить из таблицы-источника employees только самые свежие записи, минуя те, которые уже забирали ранее.
Код запроса:
SELECT * FROM employees WHERE (:?last_sync::timestamp IS NULL) OR ("Dt_Ins" > :?last_sync::timestamp OR "Dt_Upd" > :?last_sync::timestamp)
Шаг "Фильтр"

Полностью аналогичный шаг тому, как мы делали в СОПС, где информация по потоку отсутствовала.
Дальнейшие шаги также будут идентичны тому, что мы делали ранее, но в рамках закрепления материала мы их повторим
Шаг "SQL"

После получения самых актуальных записей из таблицы-источника, мы записываем их в целевую таблицу-получателя.
Не забываем добавить параметр "Групповое выполнение"

Код запроса (аналогичен предыдущему INSERT ... ON CONFLICT, поэтому не дублируем для краткости, но в реальности он такой же как на шаге 5.1.4).
Шаг "Программная трансформация"

На данном шаге мы также пробегаемся по всем актуальным записям, полученным ранее, и берем максимальное значение dt_ins и dt_upd
Код трансформации (аналогичен предыдущему, поэтому не дублируем).
Шаг "Установить заголовки"

Помещаем максимальное значение в заголовок для дальнейшего обновления таблицы cdc_tracking
Шаг "SQL"

Данный шаг немного отличается от предыдущего СОПС, поскольку в этой ветке нам уже не нужно создавать в таблице cdc_tracking интеграционный поток с именем employees.
Он уже был создан ранее, поэтому нам остается только обновить время последнего обращения к таблице
Код запроса
UPDATE public.cdc_tracking SET check_update = :?new_checkpoint::timestamptz WHERE integration_flow = 'employees'
7. Проверка работоспособности
Теперь осталось проверить, что наши настройки успешны и все работает так, как планировалось.
Для этого у нас должны быть пустые таблицы таблицы-получателя и таблицы cdc_tracking
Таблица-источник(employees) должна быть заполнена
Базовый запрос для наполнения таблицы-источника
INSERT INTO public.employees ("EmployeeId", "LastName", "FirstName", "Title", "ReportsTo", "BirthDate", "HireDate", "Address", "City", "State", "Country", "PostalCode", "Phone", "Fax", "Email") VALUES -- Руководитель (ReportsTo = NULL) (1, 'Adams', 'Andrew', 'General Manager', NULL, '1962-02-18 00:00:00+00', '2002-08-14 00:00:00+00', '11120 Jasper Ave NW', 'Edmonton', 'AB', 'Canada', 'T5K 2N1', '+1 (780) 428-9482', '+1 (780) 428-3457', 'andrew@chinookcorp.com')
7.1 Первый запуск
Запустим наш домен, который закончили настраивать.
В результате мы уйдем в ветку отсутствия записи в таблице cdc_tracking → Возьмем все записи из таблицы employees-источника → Запишем все в таблицу employees-получателя → Запишем в cdc_tracking наш поток и обновим время последнего обращения
Проверяем таблицу cdc_tracking

Все успешно записалось
Проверяем таблицу employees-получателя

Все успешно записалось
Дополнительно
Можете добавить различные шаги логов, чтобы в "Протоколе работы" наблюдать, как мы проходим через все шаги наших СОПСов
7.2 Второй запуск. Уже присутствует информация о потоке
Если мы не будем останавливать наш брокер, и не будем изменять записи в таблице-источнике, то мы будем переходить в СОПС "Информация по потоку есть".
Однако, поскольку свежих записей на текущий момент не появлялось, нам и нечего забирать, можно убедиться, поставив тот же лог после выполнения SQL запроса, где мы берем только свежие записи.

Теперь обновим какую-нибудь запись в таблице-источнике employees.
Например, у нас под EmployeeId в FirstName написано Andre, заменим на "Andrew"
После этого мы заметим, что значение для этой записи в столбце Dt_Upd изменилось и теперь выше, чем в dt_ins
Запустим наш домен повторно. Теперь мы уйдем в СОПС "Информация по потоку есть".
И если мы поставим лог после выполнения SQL запроса, где мы берем только свежие записи, то сможем увидеть, что из таблицы-источника была получена только одна запись

После выполнения СОПСа проверим таблицу cdc_tracking и убедимся, что время последнего обращения обновилось
Благодарность
На этом настройка интеграционного процесса завершена. Мы реализовали первоначальную загрузку данных, последующую передачу только новых и изменённых записей, а также сохранение контрольной точки для следующих запусков.
В результате получился простой механизм инкрементальной синхронизации PostgreSQL через FESB, который можно адаптировать под другие таблицы и интеграционные потоки.
Спасибо, что дочитали статью до конца! Надеюсь, материал оказался полезным и поможет при решении похожих задач. Буду рад вашим вопросам, замечаниям и предложениям в комментариях.

