Transactional outbox и DLQ на практике
Надёжная доставка между двумя сервисами через брокер — тема, где почти любой ответ «в общих словах» звучит одинаково у всех кандидатов. Разница видна только на конкретном коде: что происходит в момент сбоя, шаг за шагом.
Проблема без outbox
Наивная реализация: закоммитить транзакцию создания заказа, потом опубликовать сообщение в Kafka.
tx.Commit(ctx) // заказ уже в БД
producer.Publish(ctx, msg) // ...а тут процесс может упастьМежду этими двумя строками — окно, в котором заказ существует в БД, но establishment-service никогда не узнает о нём. Переставить местами (сначала publish, потом commit) не решает проблему, а только переносит её: теперь можно опубликовать сообщение о заказе, которого в БД в итоге не оказалось (транзакция откатилась).
Решение: запись в outbox той же транзакцией
err := txManager.WithinTx(ctx, func(ctx context.Context) error {
order, err := orderRepo.Create(ctx, order)
if err != nil {
return err
}
payload, _ := json.Marshal(toOrderCreatedMessage(order))
return outboxWriter.Enqueue(ctx, domain.TopicOrdersNew, order.ID.String(), payload)
})Enqueue пишет строку в outbox_events через тот же q(ctx, pool), что и orderRepo.Create — обе операции живут в одной транзакции Postgres. Либо обе применяются, либо ни одна. Публикация в Kafka здесь вообще не участвует — она происходит позже, отдельным процессом.
Фоновый relay
func (r *Relay) tick(ctx context.Context) {
rows, _ := r.store.FetchUnpublished(ctx, r.batchSize)
for _, row := range rows {
if err := r.producer.Publish(ctx, row.Topic, row.Key, row.Payload); err != nil {
continue // ретрай следующим тиком, published_at не трогаем
}
if hook, ok := r.onPublished[row.Topic]; ok {
if err := hook(ctx, row); err != nil {
continue // тот же принцип: не пометили — хук вызовется снова
}
}
r.store.MarkPublished(ctx, row.ID)
}
}Три исхода на каждую строку:
- публикация упала → пропустить, следующий тик повторит с нуля;
- публикация прошла, но пост-хук (например,
MarkSentToEstablishment) упал → тоже пропустить, не помечатьpublished_at— иначе хук никогда не будет вызван повторно, а строка уже будет выглядеть «обработанной»; - оба шага успешны → пометить.
Ключевое решение — пост-хук выполняется после publish, и его сбой откатывает не саму публикацию (она уже произошла в Kafka), а только отметку published_at. Значит одно и то же сообщение может уйти в Kafka дважды, если хук падал между попытками — это осознанно допустимо, потому что consumer на другой стороне идемпотентен (см. «Пять разных идемпотентностей»).
Юнит-тест на этот цикл — не e2e, а таблично, с фейковыми Store/Producer, покрывает ровно эти три исхода плюс топик без зарегистрированного хука:
func TestRelay_Tick_PublishFails(t *testing.T) {
store := &fakeStore{rows: []outbox.Row{{ID: 1, Topic: "orders.new"}}}
producer := &fakeProducer{err: errors.New("kafka down")}
r := outbox.NewRelay(store, producer, nil, 10)
r.Tick(context.Background())
if store.markedPublished[1] {
t.Fatal("must not mark published when publish failed")
}
}DLQ: offset — это watermark, не индивидуальный ack
Ключевая вещь, которую легко сказать неправильно на собесе: Kafka-consumer не может «отложить одно сообщение и обработать следующее», как это можно сделать, например, с SQS-очередью с индивидуальными ack по сообщению.
Offset — это одно число на партицию: «всё до этой позиции обработано». Закоммитить offset сообщения N+1, оставив N необработанным, невозможно осмысленно — коммит N+1 неявно означает «N тоже подтверждено».
Отсюда — обязательное разделение ошибок на два класса до решения, коммитить offset или нет:
func process(ctx context.Context, msg kafka.Message, handler HandlerFunc) {
isDomainErr, err := handler(ctx, msg.Key, msg.Value)
if err == nil {
consumer.CommitMessages(ctx, msg)
return
}
if isDomainErr {
toDLQ(ctx, msg, err) // никогда не станет валидным само по себе
return // toDLQ коммитит offset сам, после успешной публикации в DLQ
}
// транзиентная: до 5 попыток с backoff, ПОТОМ тоже в DLQ
retryWithBackoff(ctx, msg, handler)
}Доменная ошибка (ErrInvalidTransition — переход не разрешён и не совпадает с текущим статусом) — ретраить бессмысленно, сообщение никогда не станет валидным сколько его ни пытайся обработать снова. Сразу в DLQ-топик (orders.status_updated.dlq), затем commit основного offset.
Транзиентная ошибка (БД временно недоступна) — до 5 попыток с экспоненциальным backoff внутри обработки этого же сообщения, не через повторный FetchMessage. Если все попытки исчерпаны — тоже в DLQ, а не бесконечный ретрай. Без этого одно «вечно падающее» сообщение блокирует партицию для всех последующих заказов навсегда.
func toDLQ(ctx context.Context, msg kafka.Message, cause error) {
dlqPayload := wrapWithReason(msg.Value, cause)
if err := dlqProducer.Publish(ctx, dlqTopic, msg.Key, dlqPayload); err != nil {
return // публикация в DLQ не удалась — НЕ коммитим offset основного топика,
// сообщение переобработается при рестарте, а не потеряется
}
consumer.CommitMessages(ctx, msg)
}Даже путь в DLQ не коммитит offset, пока публикация в сам DLQ-топик не подтверждена — иначе сбой Kafka именно в этот момент тихо терял бы сообщение вместо того, чтобы просто повторить попытку позже.
Проверено e2e, не только вручную
e2e/dlq_test.go публикует напрямую в orders.status_updated с хоста (через EXTERNAL://localhost:29092 listener, специально добавленный в docker-compose.yml для этого теста) невалидный переход sent_to_establishment → completed, минуя establishment-service:
func TestOrderStatusUpdated_InvalidTransitionGoesToDLQ(t *testing.T) {
// ... публикация невалидного перехода напрямую в Kafka ...
assertMessageInDLQ(t, dlqReader, orderID) // читает с kafka.FirstOffset —
// свежий consumer-group иначе стартовал бы с "сейчас" и пропустил
// уже опубликованное сообщение
assertOrderStatusUnchanged(t, orderID, "sent_to_establishment")
publishValidTransition(t, orderID, "accepted") // тот же ключ = та же партиция
waitForStatus(t, orderID, "accepted") // партиция не заблокирована
}Три вещи проверяются вместе: сообщение действительно попало в DLQ, заказ не сдвинулся с места, и — самое важное — партиция не заблокирована poison-message'ем: следующее валидное сообщение на том же ключе (значит на той же партиции) обрабатывается штатно.
Запусти сам: как ведёт себя relay на трёх исходах
Логика tick (публикация упала / хук упал / всё прошло) не требует ни Kafka, ни Postgres, чтобы её увидеть — вот та же машина состояний на fake-реализациях, целиком:
Вывод показывает ровно то, что описано выше прозой: A не опубликован (publish упал, ретрай следующим тиком), B фактически ушёл в Kafka, но не помечен published (хук упал — при следующем тике Publish вызовется ещё раз, это осознанно допустимо), только C дошёл до конца цикла.