← назад к разделу

Когда сервис получает сигнал остановки (SIGTERM), есть риск потерять сообщения: потребитель мог получить событие из Kafka, но не успеть его обработать и зафиксировать offset — после перезапуска Kafka отдаст это сообщение снова. Продюсер может потерять несколько сообщений, которые ещё не успели уйти на broker.

В Spring Boot эти детали скрыты в ConcurrentMessageListenerContainer: контейнер сам дожимает текущую пачку и коммитит offset. В Go ту же работу нужно написать руками — это не сложнее, но требует понять правильный порядок действий.

Как работает остановка потребителя

Потребитель в kafka-go читает сообщения в бесконечном цикле через FetchMessage. Чтобы остановить этот цикл чисто, используют context.Context: когда на SIGTERM отменяется контекст, FetchMessage возвращает ошибку context.Canceled — и это нормальный выход, а не сбой.

Важная деталь про коммит offset. Автокоммит у Reader в группе есть, но только в ReadMessage: он фиксирует offset в момент выдачи сообщения, до всякой обработки. Поэтому берут пару FetchMessage + CommitMessages: offset уходит на брокер явно и только после успешной обработки.

func (c *OrderConsumer) Run(ctx context.Context) error {
    defer c.reader.Close()
    for {
        msg, err := c.reader.FetchMessage(ctx)
        if err != nil {
            if errors.Is(err, context.Canceled) || errors.Is(err, io.EOF) {
                return nil // нормальный выход при остановке
            }
            return fmt.Errorf("fetch message: %w", err)
        }

        workCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), c.handleTimeout)
        err = c.handle(workCtx, msg)
        if err == nil {
            err = c.reader.CommitMessages(workCtx, msg)
        }
        cancel()
        if err != nil {
            return fmt.Errorf("order event %d@%d: %w", msg.Partition, msg.Offset, err)
        }
    }
}

Отмену видит только FetchMessage: пока идёт обработка, цикл в нём не стоит, и сигнал будет замечен на следующем витке — после того как текущее сообщение обработано и offset зафиксирован. Чтобы так и было, handle и CommitMessages получают workCtx: он унаследовал значения ctx через context.WithoutCancel (Go 1.21), но не его отмену, и ограничен своим таймаутом. Передайте сюда тот же отменяемый ctx — и SIGTERM посреди транзакции откатит её и сорвёт коммит offset: сообщение вернётся после старта, а в логе выката появится ошибка.

defer c.reader.Close() тоже не формальность. Close явно выходит из группы, и брокер сразу отдаёт партиции оставшимся участникам. Без него брокер ждёт SessionTimeout — по умолчанию 30 секунд, — прежде чем признать участника мёртвым, и всё это время партиции старого пода никто не читает: новый экземпляр после выката полминуты молчит.

Ошибки context.Canceled и io.EOF — это нормальный выход потребителя при остановке, их не нужно логировать как ошибку. Иначе alert-канал будет шуметь при каждом деплое.

Регистрация горутины в WaitGroup

Чтобы основной процесс дождался завершения потребителя перед закрытием пула соединений, горутину регистрируют в sync.WaitGroup:

consumerCtx, cancelConsumer := context.WithCancel(ctx)

var wg sync.WaitGroup
wg.Add(1)
go func() {
    defer wg.Done()
    if err := orderConsumer.Run(consumerCtx); err != nil {
        slog.ErrorContext(ctx, "order consumer exited with error", "error", err)
    }
}()

wg.Add(1) стоит до запуска горутины — так нет гонки между стартом горутины и вызовом wg.Wait() в shutdown.

Что делать, если сообщение придёт повторно

Если SIGTERM пришёл после успешного handle, но до CommitMessages, то при следующем запуске сервиса Kafka отдаст это сообщение ещё раз. Поэтому обработчик должен быть идемпотентным — повторная обработка одного и того же события не должна приводить к двойному эффекту.

Стандартный приём — таблица уже обработанных событий:

func (c *OrderConsumer) handle(ctx context.Context, msg kafka.Message) error {
    var event OrderConfirmedEvent
    if err := json.Unmarshal(msg.Value, &event); err != nil {
        return fmt.Errorf("unmarshal order event: %w", err)
    }

    tx, err := c.pool.Begin(ctx)
    if err != nil {
        return fmt.Errorf("begin tx: %w", err)
    }
    defer tx.Rollback(ctx)

    q := c.queries.WithTx(tx)

    inserted, err := q.InsertProcessedEvent(ctx, db.InsertProcessedEventParams{
        EventID:       event.ID,
        ConsumerGroup: "billing-confirmations",
    })
    if err != nil {
        return fmt.Errorf("dedup check: %w", err)
    }
    if !inserted {
        return nil // уже обработали раньше
    }

    if err := c.orders.RecordConfirmation(ctx, event.OrderID, event.TotalAmount); err != nil {
        return fmt.Errorf("record confirmation: %w", err)
    }

    return tx.Commit(ctx)
}

InsertProcessedEvent использует ON CONFLICT (event_id, consumer_group) DO UPDATE ... RETURNING (xmax = 0) AS inserted. При первой вставке возвращает true, при повторе — false. Важно, что проверка дубля и сама бизнес-логика находятся в одной транзакции: если RecordConfirmation упадёт, транзакция откатится и событие будет переобработано при следующем запуске.

Почему нельзя делать долгий HTTP-запрос внутри handle

Распространённая ошибка — вызвать из обработчика внешний HTTP-сервис с несколькими retry:

// Опасный вариант — так делать не нужно
func (c *ProductConsumer) handle(ctx context.Context, msg kafka.Message) error {
    var event ProductPriceChangedEvent
    _ = json.Unmarshal(msg.Value, &event)

    // HTTP + retry → может занять 20-30 секунд
    if err := c.catalogClient.UpdatePrice(ctx, event.ProductID, event.NewPrice); err != nil {
        return fmt.Errorf("update price: %w", err)
    }
    return nil
}

При SIGTERM контекст отменяется и HTTP-клиент получает context.Canceled. Обработка окажется частичной, а shutdown зависнет на время таймаута.

Правильное решение — писать только в базу и класть событие в outbox-таблицу. Отдельная горутина-relay заберёт его и отправит HTTP уже независимо:

func (c *ProductConsumer) handle(ctx context.Context, msg kafka.Message) error {
    var event ProductPriceChangedEvent
    if err := json.Unmarshal(msg.Value, &event); err != nil {
        return fmt.Errorf("unmarshal product event: %w", err)
    }

    return sqlcTx(ctx, c.pool, func(q *db.Queries) error {
        if err := q.UpdateProductPrice(ctx, db.UpdateProductPriceParams{
            ProductID: event.ProductID,
            Price:     event.NewPrice,
        }); err != nil {
            return fmt.Errorf("update product price: %w", err)
        }
        return q.InsertOutboxEvent(ctx, db.InsertOutboxEventParams{
            EventType: "ProductPriceChanged",
            Payload:   msg.Value,
        })
    })
}

Так handle завершается за несколько миллисекунд, а relay работает на своём бюджете.

Остановка продюсера: зачем нужен writer.Close()

kafka.Writer копит сообщения в пачки. По умолчанию WriteMessages синхронный: он возвращается, когда пачка ушла и брокер её подтвердил — при RequiredAcks: kafka.RequireAll, как в примере ниже; с RequireNone, который стоит по умолчанию, подтверждения никто не ждёт. Такое сообщение уже не потеряется. Терять есть что в двух случаях: writer с Async: true, где вызов возвращается сразу, а пачка уходит в фоне, и вызовы WriteMessages, которые в момент выхода ещё не вернулись. Если процесс завершится без Close(), эти сообщения в Kafka не попадут.

type OrderPublisher struct {
    writer *kafka.Writer
}

func NewOrderPublisher(brokers []string) *OrderPublisher {
    return &OrderPublisher{
        writer: &kafka.Writer{
            Addr:         kafka.TCP(brokers...),
            Topic:        "orders.confirmed",
            Balancer:     &kafka.Hash{},
            RequiredAcks: kafka.RequireAll,
            BatchTimeout: 5 * time.Millisecond,
        },
    }
}

func (p *OrderPublisher) Publish(ctx context.Context, event OrderConfirmedEvent) error {
    payload, err := json.Marshal(event)
    if err != nil {
        return fmt.Errorf("marshal order confirmed: %w", err)
    }
    return p.writer.WriteMessages(ctx, kafka.Message{
        Key:   []byte(event.OrderID),
        Value: payload,
    })
}

func (p *OrderPublisher) Close() error {
    return p.writer.Close()
}

Balancer: &kafka.Hash{} стоит здесь не ради остановки, а ради порядка: балансировщик по умолчанию раскладывает сообщения по партициям по кругу и ключ не смотрит, так что события одного заказа разъехались бы по партициям. Hash кладёт один ключ в одну партицию; если партиция должна совпадать с той, что выбрал бы Java-продюсер, берут kafka.Murmur2Balancer.

writer.Close() блокируется до тех пор, пока все накопленные сообщения не отправлены на broker. После этого он возвращается — и только тогда можно закрывать остальные ресурсы.

Правильный порядок shutdown

Порядок остановки компонентов имеет значение. Нельзя закрыть пул соединений до того, как потребитель завершил обработку — следующий Begin в handle получит ошибку closed pool, сообщение не закоммитится и вернётся после старта.

shutdownFns := []func(){
    func() { appState.SetNotReady() },   // перестаём принимать трафик
    func() { cancelConsumer() },          // сигнализируем потребителю остановиться
    func() { wg.Wait() },                // ждём выхода из цикла FetchMessage
    func() { srv.Shutdown(shutCtx) },    // HTTP: дожидаемся текущих запросов
    func() {
        if err := orderPublisher.Close(); err != nil {
            slog.ErrorContext(ctx, "kafka writer close", "error", err)
        }
    },                                    // все, кто публикует, уже остановлены
    func() { pool.Close() },             // БД — самый последний
}

Логика проста: сначала останавливаем то, что принимает новую работу, затем дожидаемся завершения текущей, и только потом закрываем разделяемые ресурсы. Writer — тоже разделяемый: в него пишут и HTTP-обработчики, и горутины, поэтому закрывать его можно только после wg.Wait() и srv.Shutdown, иначе последний запрос в обработке получит ошибку записи в закрытый writer.

Частые ошибки

Автокоммит через ReadMessage. В группе ReadMessage фиксирует offset в момент выдачи сообщения — до того, как оно обработано. При SIGTERM часть сообщений будет считаться обработанной, хотя обработчик до них не добрался. Нужна пара FetchMessage + CommitMessages после каждого успешного handle. CommitInterval — другая история: он лишь копит уже подтверждённые offset и шлёт их брокеру пачкой; при штатном reader.Close() накопленное досылается, а вот SIGKILL до отправки обернётся повтором нескольких сообщений.

Запустить горутину без WaitGroup. Если не дождаться завершения потребителя перед pool.Close(), горутина может попытаться взять соединение из уже закрытого пула.

Пропустить writer.Close() или reader.Close(). Без закрытия writer теряются сообщения асинхронного режима и незавершённые вызовы; без закрытия reader брокер полминуты ждёт мёртвого участника группы, и партиции стоят.

Логировать context.Canceled как ошибку. Это нормальный способ завершения потребителя — не нужно включать в alert.

Коротко

  • Потребитель останавливается через отмену context.Context — FetchMessage возвращает context.Canceled, и это нормальный выход.
  • ReadMessage в группе коммитит offset до обработки; берите FetchMessage + CommitMessages после каждого успешного handle, а обработку ведите на context.WithoutCancel(ctx) с таймаутом.
  • Горутина-потребитель регистрируется в sync.WaitGroup — shutdown ждёт её через wg.Wait() перед закрытием пула.
  • Если SIGTERM пришёл между handle и CommitMessages, сообщение придёт повторно — обработчик должен быть идемпотентным.
  • Долгие HTTP-запросы с retry из handle опасны: при отмене контекста обработка будет частичной. Используй outbox-паттерн.
  • writer.Close() дожидается незавершённых записей и пачек асинхронного режима, reader.Close() выходит из группы сразу — без него партиции ждут SessionTimeout.
  • Порядок shutdown: отмена контекста → wg.Wait() → HTTP drain → writer.Close() → pool.Close().

Что почитать дальше