Никакой интеграции нет. Есть реплика чужого состояния и спор о том, когда оно изменилось
Никакой интеграции нет. Есть реплика чужого состояния и спор о том, когда оно изменилось

В каждой компании такая есть. Система, которая владеет чем-то важным — обычно деньгами, — которую тебе не разрешено трогать, чей календарь релизов принадлежит другому отделу и чья модель данных была утверждена до того, как тебя наняли. ERP. Биллинг. Мейнфрейм с прослойкой, которая скрапит его же экраны.

У меня это была 1С. Уточню сразу, потому что дальше это важно: конкретно 1С в этой истории — наименее интересная деталь. Подставьте SAP, NetSuite или двадцатилетнее приложение на Oracle Forms — ниже не изменится ни строчки.

Я делал это дважды, с противоположных сторон. Один раз — Kafka-консьюмерами к ERP, которую трогать не давали. И один раз — терабайтом банковских данных, переехавшим с Oracle на PostgreSQL без даунтайма: семь десятков таблиц по кредитам и займам, миграция фазами, Kafka-пайплайн держит обе системы согласованными, пока деньги продолжают ходить, и всё это устаканивается на 500–3000 RPS на чтение и 50–300 на запись. Вторая задача — та же самая, только в приличном костюме: на время миграции у тебя работают две системы, которые обязаны сходиться, и расхождение принадлежит тебе. Всё, что делает 1С особенной в этом сюжете, — это то, что она очень хороша ровно в той роли, для которой её ставили, и очень не приспособлена к той, в которой её хочет видеть бэкенд.

Ошибка, которую я сделал — и, по-моему, делают почти все, — состояла в слове интеграция. Слово намекает на две системы, встречающиеся где-то посередине. Происходит не это. Происходит вот что:

Ты строишь реплику чужого состояния, и любой баг, который у тебя когда-либо будет, — это спор о том, когда это состояние изменилось.

Не во что изменилось. Когда. Как только я начал читать инциденты так, они перестали быть списком несвязанных проблем и стали одной проблемой в разных костюмах.

Дисклеймер: проект под NDA. Названия, точные объёмы и часть деталей изменены. Порядки величин, механизмы и сами инциденты — настоящие.

У источника нет событий. У него есть состояние

Это первое, что надо усвоить, и именно об этом тихо врёт любая схема архитектуры.

Такие системы проектировались как системы учёта. Они хранят текущую правду и делают это отлично. Чего у них нет — потока «вот это изменилось вот в этот момент». Вместо него обычно предлагают:

  • таблицу, которую можно читать (если повезло — представление);

  • периодическую выгрузку, которую кладут на шару;

  • OData/SOAP-эндпоинт, отдающий текущее состояние того, о чём спросили;

  • и если очень повезло — таблицу изменений, которую кто-то завёл под другой проект в 2016-м.

Ничто из этого не является событиями. Поэтому первое, что ты строишь, — компонент, который изготавливает события, регулярно разглядывая состояние и описывая разницу. И первая архитектурная ошибка — нарисовать этот компонент квадратиком «источник» и дальше рассуждать обо всём ниже по течению так, будто это настоящий поток событий.

Это не поток. Это догадка, сделанная с интервалом, о том, что случилось между двумя наблюдениями. Каждое свойство, которое ты хочешь от потока, — полнота, порядок, exactly-once — придётся восстанавливать руками на этой границе, потому что на этой границе они и были потеряны.

Баг 1: опрос по updated_at теряет строки и никогда об этом не говорит

Вот поллер, который пишут первым. Я тоже написал такой.

SELECT * FROM documents
WHERE  updated_at > :last_seen
ORDER  BY updated_at;

Дальше ставим last_seen в максимальный прочитанный updated_at и повторяем. Очевидно, эффективно, ложится на индекс — и молча теряет данные. Механизм:

Строка, которую поллер по updated_at не увидит никогда
Строка, которую поллер по updated_at не увидит никогда

Строка A получила отметку 10:00, потому что тогда её транзакция началась, а видимой для читателей стала в 10:07, потому что тогда она закоммитилась. Поллер отработал между. Он сдвинул водяной знак за отметку A раньше, чем A вообще начала существовать с точки зрения любого читающего.

Ничего не падает. Ничего не уезжает в DLQ. Строку просто никогда не увидят, и две системы навсегда расходятся во мнении о документе, в наличии которого учётная система абсолютно уверена.

Вот здесь стоит сказать прямо: я не знаю, во сколько строк это обошлось, и это ровно та самая проблема. Ничто в пайплайне не могло мне этого сообщить. Всплыло потому, что человек сверил итоги двух систем и обнаружил расхождение, — а это не стратегия мониторинга, это её отсутствие.

Три лечения, по возрастанию того, сколько тебе позволено:

Окно перекрытия. Читать updated_at > last_seen - 15 минут и сделать запись идемпотентной. Каждый цикл переобрабатываешь немного лишнего и перестаёшь терять строки, чья транзакция была короче окна. Дёшево и проблему не решает — покупает запас. Размер окна берётся из самой длинной транзакции, которую источник реально выполняет, а не из числа, которое кажется безопасным. И помни, что закрывающий месяц пакетник, крутящийся два часа, пробьёт любое такое окно навылет.

Водяной знак по последовательности, а не по часам. Если у источника есть монотонный идентификатор, присваиваемый на коммите, — используй его. Часы — это значение, которое источник записал; номер коммита — значение, которое источник упорядочил. Это разные по природе вещи, и сравнивать безопасно только второе.

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

И четвёртое, что лечением не является, но обязательно: часы источника — не твои часы. Если 1С крутится на хосте, чьё время отстаёт на две минуты от того, по которому живёт поллер, у тебя двухминутное окно потери данных зашито ещё до того, как выполнилась хоть одна твоя строка кода. Проверь. Это проверка в одну команду, и я ни разу не видел её в раннбуке интеграции.

Баг 2: ты выбрал не тот ключ партиционирования

Как только изготовленные события уезжают в Kafka, ты получаешь ровно одну гарантию порядка: сообщения с одинаковым ключом попадают в одну партицию, а партиция упорядочена. Всё. Между партициями порядка нет никакого и не было.

Поэтому вот это — баг:

producer.send(new ProducerRecord<>("erp.documents", event.messageId(), payload));

Уникальный ключ на сообщение размазывает жизненный цикл одного документа по всем партициям топика. Его created и его updated теперь гоняются наперегонки, и консьюмер применит их в том порядке, в каком брокеры их отдадут. Обычно разрыв — миллисекунды, и ты этого не видишь никогда. В закрытие месяца, когда 1С пакетом переписывает документы за день, а консьюмер-группа под нагрузкой ребалансится, видишь часто.

Ключ — это бизнес-сущность:

producer.send(new ProducerRecord<>("erp.documents", event.documentId(), payload));

Два следствия, о которых лучше знать до выбора, потому что кусаются они позже.

Число партиций — часть контракта. Отображение — hash(key) % partitions. Увеличил число партиций — переехали все ключи: документ, чья история лежала в партиции 3, начинает приходить в партицию 7, а его старые сообщения всё ещё стоят в очереди в 3. Гарантия порядка в этот момент не деградирует плавно — она просто отсутствует для всего, что в полёте. Масштабирование топика — это миграция, а не правка конфига, и относиться к ней надо как к миграции.

Горячий ключ — горячая партиция. Если 30% документов принадлежат одному контрагенту, а ты сключевался по контрагенту, потому что так «естественнее», — ты построил систему, чья пропускная способность упирается в один поток консьюмера, и наращивание группы не поможет.

Баг 3: exactly-once — свойство твоей записи, а не галочка в конфиге

Транзакции в Kafka настоящие и делают ровно то, что обещают, — но обещают они меньше, чем в них читают. Они дают атомарность для consume-transform-produce внутри Kafka. В тот момент, когда побочный эффект консьюмера — это INSERT в PostgreSQL, вызов внешнего API или запись файла, транзакция туда не распространяется. Две разные системы, общего коммита нет.

И это нормально, потому что он и не нужен. Нужна идемпотентная запись:

INSERT INTO documents (source_id, version, amount, status, raw)
VALUES (:source_id, :version, :amount, :status, :raw)
ON CONFLICT (source_id) DO UPDATE
SET amount  = EXCLUDED.amount,
    status  = EXCLUDED.status,
    raw     = EXCLUDED.raw,
    version = EXCLUDED.version
WHERE documents.version <= EXCLUDED.version;   -- никогда не класть старую версию поверх новой

Работают две вещи. ON CONFLICT по натуральному ключу источника делает передоставку бесплатной: at-least-once превращается в фактически-однократно без всякой координации. А условие в WHERE делает безопасной доставку не по порядку — что важно, потому что баг 2 полностью не решается: ребаланс может отдать партицию новому консьюмеру, пока у старого ещё летит сообщение.

Вот это условие и выкидывают чаще всего, и именно оно спасает. Без него переотправленное старое сообщение затирает новое, и реплика молча уезжает назад во времени.

Баг 4: удаление — тоже не событие

Удалённая строка ничего тебе не пришлёт. Как и строка, которую заархивировали, перенесли в другую таблицу при закрытии года или у которой поменяли статус процессом, не трогающим updated_at.

Отсутствие нельзя доставить. Это жёсткий предел всего подхода: реплика на опросе может сойтись с тем, что существует, и никогда не узнает сама о том, что существовать перестало.

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

Сверка — это и есть интеграция. Поток — просто оптимизация

Ночная джоба. Задаёт обеим системам один и тот же агрегатный вопрос и сравнивает ответы:

-- с обеих сторон, одинаковой формы
SELECT doc_date, account_id, count(*) AS n, sum(amount) AS total
FROM   documents
WHERE  doc_date BETWEEN :from AND :to
GROUP  BY doc_date, account_id;

Агрегаты, а не строки. Построчное сравнение дорожает сразу, как только суточный объём становится интересным, и, что хуже, выдаёт такой длинный дифф, что его никто не читает. Количества и суммы по дню и счёту дают горсть разошедшихся корзин, и каждая корзина — точечное перечитывание маленького среза.

Три правила, решающие, будет она полезной или просто шумной. Это дизайн, который я готов защищать и который в следующий раз построил бы первым:

  1. Она должна уметь чинить, а не только жаловаться. Расхождение запускает полное перечитывание среза за этот день через тот же идемпотентный путь записи. Сверка, которая только заводит тикеты, будет замьючена в течение месяца.

  2. Она должна перепроверять закрывающееся окно, а не только вчера. Бухгалтерия проводит задним числом. Документ, проведённый сегодня, может относиться к прошлому месяцу. Ретроспектива в неделю, а не в сутки, — это разница между сверкой, которая видит документы задним числом, и сверкой, которая структурно не способна их увидеть.

  3. Её вывод — метрика, а не строчка в логе. «Корзин разошлось за ночь» на дашборде, с алертом на тренд, а не на каждое срабатывание. Стабильные одна-две — норма. Сорок — инцидент, и интересен именно скачок.

Я бы сказал даже жёстче, потому что это то, что я сделал бы иначе, начиная заново:

Стройте сверку раньше потока. Это компонент, который говорит, работает ли вообще всё остальное; единственный, который переживает любой редизайн пайплайна; и тот, который всегда планируют последним и режут первым.

Потоковая интеграция без сверки — не источник правды. Это слух с хорошей латентностью.

Баг 5: DLQ, у которого нет хозяина

Ты его напишешь. Сообщение не распарсилось или нарушило констрейнт, и вместо того, чтобы навсегда заблокировать партицию, ты отодвигаешь его в сторону и едешь дальше. Правильно.

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

DLQ без хозяина — это удаление с лишними шагами и худшей совестью. Если ты не собираешься строить алерт и путь переигрывания, честная альтернатива — останавливать консьюмер и будить человека: это хотя бы падает громко. Выбирай осознанно; не позволяй DLQ стать местом, где решение не принимают.

Схема — это то, что бухгалтер сделал во вторник

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

С такой системой не бывает согласования схемы. Новое поле появляется, потому что его настроил пользователь в финансах. Энум отращивает значение, потому что для одного региона завели новый вид документа. Код, который всегда был числовым, обзаводится буквой — в релизе, о котором ты узнал из графика ошибок.

Политика устойчивости, пережившая столкновение с этой реальностью, — два правила, которые выглядят противоречиво:

Никогда не падать на неизвестном поле. Сохранить, залогировать один раз, ехать дальше. Падение здесь означает, что одна безобидная настройка в отделе, с которым ты никогда не разговаривал, кладёт твой пайплайн.

Всегда падать на неизвестном значении в поле, по которому ты ветвишься. Если у document_type четыре известных значения и приезжает пятое — стоп. Не проваливаться в default, который посчитает его самым частым типом. Падение здесь — испорченное утро; молчаливый дефолт — месяцы тихо неверных чисел, которые кто-нибудь однажды найдёт в отчёте. Та же форма отказа, что у скрапера, который возвращает валидные, правдоподобные и неверные данные.

И под обоими правилами: храни сырой payload в каждой строке. Колонка jsonb, нетронутое сообщение, навсегда. Это стоит диска и это единственная причина, по которой ты сможешь перевывести три месяца истории в день, когда выяснится, что маппинг был неверен с июня. Все интеграции, где эта колонка была, оказывались восстановимыми. Те, что парсили на входе и хранили только результат, — нет.

Одним предложением

Если сжать всё выше до чего-то, что стоит запомнить:

Ты не владеешь источником, значит не владеешь временем. Всё, что ты строишь, — попытка установить порядок, которого источник никогда не обещал, и сверка — единственный компонент, способный сказать, получилось ли.

Поток даёт латентность. Сверка даёт корректность. Команды стабильно строят их в этом порядке и стабильно обнаруживают, что второе было нужно первым.


Вот здесь я жду несогласия. Я утверждаю, что сверка должна быть раньше потока, — а это значит выкатить ночной пакетник раньше, чем то, ради чего проект вообще профинансировали. В большинстве компаний это тяжёлый разговор, и «сверку добавим во второй фазе» — фраза, которую я и слышал, и сам произносил.

Кто-нибудь реально выигрывал этот спор заранее? Или сверка достаётся всем так же, как досталась мне, — после того, как человек нашёл пропавшие строки, сверяя итоги в экселе?