Проблема: dual-write
Типовой сценарий -- создать запись в БД и отправить событие в Kafka:
// НЕ делайте так
$db->beginTransaction();
$db->exec("INSERT INTO orders ...");
$db->commit();
$kafka->publish('orders.created', $event); // а если сеть моргнула?
Четыре возможных сценария при сбое:
| Сценарий | БД | Kafka | Последствие |
|---|---|---|---|
| Happy path | OK | OK | Консистентно |
| BD fail | FAIL | - | OK, ничего не произошло |
| Сеть умерла после commit | OK | FAIL | Заказ есть, события нет, склад не узнает |
| Сеть умерла после publish | ??? | OK | Событие послано, а commit? Падение перед commit = есть событие о несуществующем заказе |
Поменяйте порядок -- проблема остаётся, только зеркально. Это и есть dual-write problem: две независимые системы не дают атомарности.
Решение: Transactional Outbox
Идея: записывать событие в ту же БД, что и основные данные, в одной транзакции. Отдельный worker читает таблицу outbox и публикует в Kafka.
┌─────────────────────────────┐
│ API Handler │
│ │
│ BEGIN TRANSACTION │
│ INSERT INTO orders ... │
│ INSERT INTO outbox ... │
│ COMMIT │ <-- атомарно
└──────────┬──────────────────┘
│
│ (записано)
v
┌─────────────────────────────┐
│ PostgreSQL │
│ orders | outbox │
└──────────┬──────────────────┘
│
│ pg_notify / polling
v
┌─────────────────────────────┐
│ Outbox Publisher Worker │
│ SELECT FOR UPDATE SKIP │
│ Publish to Kafka │
│ UPDATE sent_at = NOW() │
└──────────┬──────────────────┘
│
v
Kafka
Гарантия: at-least-once -- событие будет опубликовано, возможно несколько раз (если worker упал между publish и UPDATE). Консьюмеры должны быть идемпотентны (см. inbox pattern ниже).
Схема outbox таблицы
CREATE TABLE outbox (
id BIGSERIAL PRIMARY KEY,
aggregate_type TEXT NOT NULL, -- 'order', 'user', etc
aggregate_id TEXT NOT NULL, -- ID сущности
event_type TEXT NOT NULL, -- 'OrderCreated'
payload JSONB NOT NULL, -- тело события
headers JSONB, -- correlation_id, trace_id
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
sent_at TIMESTAMPTZ, -- NULL = pending
attempts INT NOT NULL DEFAULT 0,
error TEXT
);
-- Индекс для worker: берёт только pending, в порядке создания
CREATE INDEX idx_outbox_pending ON outbox (created_at)
WHERE sent_at IS NULL;
-- По aggregate для отладки
CREATE INDEX idx_outbox_agg ON outbox (aggregate_type, aggregate_id);
Важно: FIFO per aggregate -- для одного aggregate_id события должны публиковаться в порядке вставки. На Kafka решается через partition key = aggregate_id.
Producer в handler
<?php
declare(strict_types=1);
namespace App\Outbox;
use Doctrine\DBAL\Connection;
/**
* Outbox writer. Call within the same transaction as your main INSERT/UPDATE.
*/
final class OutboxWriter
{
public function __construct(private readonly Connection $conn) {}
/**
* @param array<string, mixed> $payload
* @param array<string, mixed> $headers
*/
public function append(
string $aggregateType,
string $aggregateId,
string $eventType,
array $payload,
array $headers = [],
): void {
$this->conn->insert('outbox', [
'aggregate_type' => $aggregateType,
'aggregate_id' => $aggregateId,
'event_type' => $eventType,
'payload' => json_encode($payload, JSON_THROW_ON_ERROR),
'headers' => json_encode($headers, JSON_THROW_ON_ERROR),
]);
}
}
final class CreateOrderHandler
{
public function __construct(
private readonly Connection $conn,
private readonly OutboxWriter $outbox,
) {}
public function handle(CreateOrderCommand $cmd): string
{
$this->conn->beginTransaction();
try {
$orderId = bin2hex(random_bytes(16));
// Main business write
$this->conn->insert('orders', [
'id' => $orderId,
'user_id' => $cmd->userId,
'total_cents' => $cmd->totalCents,
'status' => 'created',
]);
// Outbox write in the same transaction -- atomic!
$this->outbox->append(
aggregateType: 'order',
aggregateId: $orderId,
eventType: 'OrderCreated',
payload: [
'orderId' => $orderId,
'userId' => $cmd->userId,
'totalCents' => $cmd->totalCents,
],
headers: [
'correlationId' => $cmd->correlationId,
'occurredAt' => (new \DateTimeImmutable())->format(DATE_RFC3339_EXTENDED),
],
);
$this->conn->commit();
return $orderId;
} catch (\Throwable $e) {
$this->conn->rollBack();
throw $e;
}
}
}
<?php
declare(strict_types=1);
namespace App\Outbox;
use Doctrine\DBAL\Connection;
use Psr\Log\LoggerInterface;
use RdKafka\Producer;
/**
* Outbox publisher. Polls pending rows, publishes to Kafka, marks sent.
* Runs as a long-lived worker (1 or more replicas, coordinate via SKIP LOCKED).
*/
final class OutboxPublisher
{
private const BATCH_SIZE = 100;
private const POLL_INTERVAL_MS = 200;
public function __construct(
private readonly Connection $conn,
private readonly Producer $kafka,
private readonly LoggerInterface $log,
) {}
public function run(): never
{
while (true) {
$processed = $this->tick();
if ($processed === 0) {
usleep(self::POLL_INTERVAL_MS * 1000);
}
}
}
private function tick(): int
{
// SKIP LOCKED allows N parallel workers without duplication
$sql = <<<'SQL'
SELECT id, aggregate_type, aggregate_id, event_type, payload, headers, attempts
FROM outbox
WHERE sent_at IS NULL
ORDER BY id ASC
LIMIT :limit
FOR UPDATE SKIP LOCKED
SQL;
$this->conn->beginTransaction();
try {
$rows = $this->conn->fetchAllAssociative($sql, ['limit' => self::BATCH_SIZE]);
foreach ($rows as $row) {
$this->publish($row);
}
$this->conn->commit();
return count($rows);
} catch (\Throwable $e) {
$this->conn->rollBack();
$this->log->error('outbox tick failed', ['error' => $e->getMessage()]);
return 0;
}
}
/**
* @param array<string, mixed> $row
*/
private function publish(array $row): void
{
$topic = $this->kafka->newTopic($this->topicFor($row['aggregate_type']));
// Partition key = aggregate_id to preserve ordering
$topic->producev(
partition: RD_KAFKA_PARTITION_UA,
msgflags: 0,
payload: $row['payload'],
key: $row['aggregate_id'],
headers: [
'event_type' => $row['event_type'],
'event_id' => (string) $row['id'],
],
);
$this->kafka->flush(5000);
$this->conn->update(
'outbox',
['sent_at' => (new \DateTimeImmutable())->format('Y-m-d H:i:s')],
['id' => $row['id']],
);
}
private function topicFor(string $aggregateType): string
{
return match ($aggregateType) {
'order' => 'orders.events',
'user' => 'users.events',
default => throw new \InvalidArgumentException("unknown aggregate: $aggregateType"),
};
}
}
package outbox
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"log/slog"
"time"
"github.com/segmentio/kafka-go"
)
const batchSize = 100
type Event struct {
ID int64
AggregateType string
AggregateID string
EventType string
Payload json.RawMessage
Headers json.RawMessage
}
type Publisher struct {
db *sql.DB
writer *kafka.Writer
log *slog.Logger
topicFn func(aggregateType string) string
}
// Run polls the outbox and publishes pending events.
// Graceful shutdown via ctx.
func (p *Publisher) Run(ctx context.Context) error {
ticker := time.NewTicker(200 * time.Millisecond)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
if err := p.tick(ctx); err != nil {
p.log.Error("tick failed", "err", err)
}
}
}
}
func (p *Publisher) tick(ctx context.Context) error {
tx, err := p.db.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelReadCommitted})
if err != nil {
return err
}
defer tx.Rollback()
rows, err := tx.QueryContext(ctx, `
SELECT id, aggregate_type, aggregate_id, event_type, payload, headers
FROM outbox
WHERE sent_at IS NULL
ORDER BY id ASC
LIMIT $1
FOR UPDATE SKIP LOCKED`, batchSize)
if err != nil {
return fmt.Errorf("select: %w", err)
}
var events []Event
for rows.Next() {
var e Event
if err := rows.Scan(&e.ID, &e.AggregateType, &e.AggregateID, &e.EventType, &e.Payload, &e.Headers); err != nil {
rows.Close()
return err
}
events = append(events, e)
}
rows.Close()
if len(events) == 0 {
return nil
}
// Batch publish; partition key = aggregate_id for FIFO per aggregate
msgs := make([]kafka.Message, 0, len(events))
for _, e := range events {
msgs = append(msgs, kafka.Message{
Topic: p.topicFn(e.AggregateType),
Key: []byte(e.AggregateID),
Value: e.Payload,
Headers: []kafka.Header{
{Key: "event_type", Value: []byte(e.EventType)},
{Key: "event_id", Value: []byte(fmt.Sprintf("%d", e.ID))},
},
})
}
if err := p.writer.WriteMessages(ctx, msgs...); err != nil {
return fmt.Errorf("kafka write: %w", err)
}
// Mark sent
ids := make([]int64, len(events))
for i, e := range events {
ids[i] = e.ID
}
if _, err := tx.ExecContext(ctx,
`UPDATE outbox SET sent_at = NOW() WHERE id = ANY($1)`, ids); err != nil {
return fmt.Errorf("update: %w", err)
}
return tx.Commit()
}
Inbox pattern для дедупликации
Worker может опубликовать событие, но упасть до UPDATE sent_at = событие придёт в Kafka дважды. Consumer должен уметь это пережить.
CREATE TABLE inbox (
event_id TEXT PRIMARY KEY, -- из headers (event_id publisher'а)
topic TEXT NOT NULL,
processed_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
<?php
declare(strict_types=1);
final class InventoryConsumer
{
public function __construct(
private readonly Connection $conn,
private readonly InventoryService $service,
) {}
public function handle(KafkaMessage $msg): void
{
$eventId = $msg->headers['event_id'];
$this->conn->beginTransaction();
try {
// Atomic: insert inbox marker + business logic in one TX
$inserted = $this->conn->executeStatement(
'INSERT INTO inbox (event_id, topic) VALUES (:id, :t) ON CONFLICT DO NOTHING',
['id' => $eventId, 't' => $msg->topic],
);
if ($inserted === 0) {
// Already processed -- idempotent skip
$this->conn->rollBack();
return;
}
$payload = json_decode($msg->value, true, flags: JSON_THROW_ON_ERROR);
$this->service->reserve($payload['orderId']);
$this->conn->commit();
} catch (\Throwable $e) {
$this->conn->rollBack();
throw $e;
}
}
}
Если важна латентность -- можно разбудить worker триггером:
CREATE FUNCTION outbox_notify() RETURNS TRIGGER AS $$
BEGIN
PERFORM pg_notify('outbox_new', NEW.id::text);
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER outbox_notify_trigger
AFTER INSERT ON outbox
FOR EACH ROW EXECUTE FUNCTION outbox_notify();
Worker делает LISTEN outbox_new, при каждом сигнале -- tick. Fallback на polling раз в 5 секунд на случай пропущенного notify (они best-effort).
CDC как альтернатива: Debezium
Вместо своего publisher'а можно читать WAL (Write-Ahead Log) PostgreSQL через Debezium. Он стримит изменения любой таблицы в Kafka автоматически.
| Аспект | Outbox + Worker | Debezium CDC |
|---|---|---|
| Код | Ваш + миграция | Конфиг Debezium |
| Гибкость events | Полный контроль payload | Маппинг row -> event |
| Операционная нагрузка | Worker мониторить | Kafka Connect кластер |
| Latency | 100-500ms | 10-100ms |
| Схема событий | Бизнес-осмысленная | Отражает БД |
Outbox лучше, если нужны бизнес-события с осмысленной схемой. CDC -- для репликации БД в поисковые/аналитические системы без написания кода.
Cleanup outbox
Outbox растёт линейно. Варианты уборки:
-- Keep 7 days history for audit/replay
DELETE FROM outbox
WHERE sent_at < NOW() - INTERVAL '7 days';
-- Or partition by day and DROP old partitions (faster)
Партицирование по дню -- быстрее, чем DELETE миллионов строк.
Мониторинг
| Метрика | Назначение | Алерт |
|---|---|---|
outbox_pending (gauge) |
Количество неопубликованных | > 1000 |
outbox_lag_seconds |
max(now - created_at) для pending | > 60 |
outbox_published_total |
counter | - |
outbox_publish_errors_total |
counter | rate > 0.01/s |
outbox_worker_tick_duration |
histogram | p95 > 1s |
-- Lag query
SELECT EXTRACT(EPOCH FROM (NOW() - MIN(created_at))) AS lag_seconds
FROM outbox
WHERE sent_at IS NULL;
Pitfalls
- Забыли COMMIT outbox вместе с бизнес-данными -- обычный dual-write. Обязательно одна транзакция.
- FIFO нарушается при N workers без шардинга -- order с id=1 опубликован worker B, id=2 -- worker A раньше. Решение: shard by
hash(aggregate_id) % N. - Consumer не идемпотентен -- at-least-once превратится в
обработано N раз. Inbox pattern обязателен. - Payload слишком большой (>1MB) -- Kafka откажется. Храните blob отдельно, в outbox -- ссылку.
- outbox таблица распухла -- ленивые DELETE / партиции обязательны.
- Worker публикует, но не UPDATE (упал) -- дубли. Консьюмеры обязаны быть идемпотентны.
- Schema evolution: изменили структуру payload -- consumers старой версии сломаются. Используйте schema registry (Avro/Protobuf).
Выводы
- Outbox решает dual-write проблему через атомарную запись события в ту же БД
- Гарантия: at-least-once. Consumers обязаны быть идемпотентны (inbox pattern)
- FOR UPDATE SKIP LOCKED позволяет масштабировать publisher'ы без дублей
- FIFO per aggregate достигается через partition key = aggregate_id
- CDC (Debezium) -- альтернатива без кода, но с менее осмысленной схемой событий
- Мониторинг outbox_lag и outbox_pending критичен
- Очистка таблицы через партицирование, не DELETE