Проблема називається dual write: у транзакції ви змінюєте базу, а поряд треба сказати про це решті світу — покласти повідомлення в Kafka, RabbitMQ, SQS чи Redis. Це два різні сховища без спільної транзакції, і жоден порядок дій не рятує. Опублікуєте всередині транзакції — ROLLBACK відкотить рядок, а споживачі вже отримали подію про те, чого не сталося: гроші «списані», лист надісланий, склад зарезервував товар. Опублікуєте після COMMIT — між комітом і publish() лишається проміжок у кілька мілісекунд, і падіння процесу, OOM-kill контейнера чи розрив зʼєднання з брокером саме там означають, що стан змінився, а події не буде ніколи, і ніхто про це не дізнається. Двофазний коміт закрив би питання теоретично, але Kafka, SQS і Redis його не пропонують, а MySQL XA історично болісний у відновленні підвішених гілок — на практиці цей шлях не беруть.
Вихід у тому, щоб зробити атомарним не «база + брокер», а «база + база»: подію записують рядком у таблицю outbox тим самим зʼєднанням і тією самою транзакцією, що й зміну стану. Тепер відкат забирає обидва записи разом, а успішний COMMIT означає, що подія існує як факт, зафіксований на диску. Окремий процес-релей після коміту читає невідправлені рядки, публікує їх у брокер і позначає sent_at. Ключова деталь вибірки — SELECT ... FOR UPDATE SKIP LOCKED (PostgreSQL 9.5+, MySQL 8.0+): він дозволяє кільком релеям розбирати різні пачки й не дає двом процесам взяти ту саму сотню рядків. Чого робити не варто — тримати позначку id > :last зовні: значення послідовності видаються до коміту, тому рядок з меншим id може стати видимим пізніше за рядок з більшим, і подія тихо випаде з обробки.
Головне, чого outbox не дає — exactly-once. Релей може опублікувати повідомлення і впасти до UPDATE ... SET sent_at, після рестарту він опублікує його вдруге; сам брокер теж працює в режимі at-least-once. Тому дедуплікація не опція, а друга половина патерна: кожна подія несе message_id, згенерований ще в транзакції запису, а споживач в одній транзакції вставляє цей id у таблицю оброблених повідомлень з унікальним індексом і виконує корисну дію. Порушення унікальності = «вже робили», повідомлення підтверджують і йдуть далі. Це та сама ідея ключа ідемпотентності, що й у платежах, тільки застосована до споживача черги. Порядок теж не глобальний: у межах агрегату його дає ключ партиціювання (aggregate_id як ключ повідомлення в Kafka) або номер версії всередині події, який споживач порівнює з уже застосованим.
У payload кладуть уже серіалізований стан на момент транзакції, а не лише ідентифікатор. Різниця принципова: релей, який за id піде читати поточний рядок, віддасть стан на момент читання — після ще двох змін, а подія має описувати факт, що стався. Тонка подія з самим id теж має право на життя (персональні дані, великі документи), але тоді треба свідомо визнати, що споживач працює зі свіжішим станом, і закласти версіонування. Поряд з payload у рядку тримають тип події, ключ агрегату, occurred_at і версію схеми — цього достатньо, щоб через рік додати новий формат, не ламаючи старих споживачів.
У PHP-стеку є приємна деталь: те, що часто описують як «треба зробити outbox», у багатьох проєктах уже стоїть. Черга Laravel на драйвері database пише завдання в таблицю jobs, а Doctrine-транспорт Symfony Messenger (doctrine://default) — у messenger_messages. Якщо це те саме зʼєднання і диспетч відбувається всередині вашої транзакції, ви вже отримали транзакційний outbox без єдиного власного класу, а воркер грає роль релея. Плутати з цим dispatch()->afterCommit() (чи after_commit => true) і DispatchAfterCurrentBusMiddleware не можна: вони лише відкладають відправлення за коміт, лікуючи «job не бачить щойно створеної моделі», і вікно втрати між COMMIT і publish у них залишається.
Ціна патерна — затримка й експлуатація. Подія доходить не миттєво, а за час циклу релея, і це треба закласти в UX там, де користувач очікує реакції одразу. Таблиця росте: рядки або видаляють одразу після публікації, або позначають і чистять пакетно за розкладом, і в PostgreSQL частковий індекс WHERE sent_at IS NULL тут майже обовʼязковий, бо інакше індекс тягне за собою весь архів, а постійні UPDATE+DELETE роздувають таблицю. Обовʼязковий і моніторинг: алерт на вік найстарішого невідправленого рядка й на розмір черги — застряглий релей інакше виявляють через скаргу «клієнту не прийшов лист», коли минуло вже пів дня. Коли polling впирається в межу або систем-джерел стає багато, наступний крок — CDC: Debezium читає WAL/binlog і публікує рядки outbox без запитів до бази, ціною Kafka Connect, слотів реплікації та ще одного сервісу, за яким треба стежити.
// 1. Запис стану і події — одна транзакція, одне зʼєднання. Брокера тут немає.
DB::transaction(function () use ($order): void {
$order->markPaid();
$order->save();
DB::table('outbox')->insert([
'message_id' => (string) Str::ulid(), // ключ дедуплікації для споживача
'aggregate_id' => $order->id, // ключ партиції: порядок у межах замовлення
'type' => 'order.paid',
'payload' => json_encode([ // стан НА момент транзакції, не id
'order_id' => $order->id,
'amount' => $order->amount_cents,
'currency' => $order->currency,
]),
'occurred_at' => now(),
'sent_at' => null,
]);
}); // ROLLBACK відкотить і зміну, і подію — фантомних подій не буває
// 2. Релей: окремий процес, крутиться в циклі. Публікує вже після коміту.
final readonly class OutboxRelay
{
public function __construct(private Publisher $broker) {}
public function drainBatch(int $limit = 100): int
{
return DB::transaction(function () use ($limit): int {
// SKIP LOCKED: кілька релеїв беруть різні пачки (PostgreSQL 9.5+, MySQL 8.0+)
$rows = DB::table('outbox')
->whereNull('sent_at')
->orderBy('id')
->limit($limit)
->lockForUpdate()->skipLocked()
->get();
foreach ($rows as $row) {
// Падіння тут -> рядок лишиться невідправленим -> повтор.
// Падіння після publish, але до update -> дубль. Звідси at-least-once.
$this->broker->publish($row->type, $row->aggregate_id, $row->payload, $row->message_id);
DB::table('outbox')->where('id', $row->id)->update(['sent_at' => now()]);
}
return $rows->count();
});
}
}
Скажіть одним реченням, звідки проблема: «база й брокер — два сховища без спільної транзакції, тож атомарним може бути лише запис у ту саму базу». І одразу назвіть ціну: at-least-once, тому дедуплікація на споживачі — частина патерна, а не окрема опція.