Перейти к основному содержимому

Работа с Kafka

В Go-сервисах Ensi Kafka используется для обмена событиями между сервисами. Типичные сценарии такие: мастер-система опубликовала оффер или изменила бренд, маркетинговый сервис отправил команду пересчитать скидки, ваш сервис должен подхватить изменение и обновить локальный кэш или индекс. Вы либо публикуете сообщение в топик, либо подписываетесь на топик и обрабатываете входящий поток, либо делаете и то и другое.

Правила именования топиков и контуров описаны в Kafka Design Guide. Ниже — как писать и читать сообщения в сервисе на Goravel.

Код разделён по слоям так же, как остальная инфраструктура сервиса.

  • Общий клиент, registry хэндлеров, supervisor консьюмеров и продюсер лежат в app/adapters/kafka.
  • Обработчики конкретных фич живут рядом с модулем, в app/modules/{module}/kafka.
  • Список «какой ключ топика обслуживает какой хэндлер» собирается в bootstrap/kafka.go.
  • Параметры брокера, карта ключей на полные имена топиков и список слушаемых ключей по умолчанию задаются в config/kafka.go и переменных окружения.

Топики: короткий ключ вместо полного имени

В прикладном коде почти никогда не пишут настоящее имя топика вроде local.catalog.fact.offers.1. Такое имя зависит от контура и домена и легко разъезжается между окружениями. Вместо этого в config/kafka.go, в секции topics, для каждого топика заводят короткий ключ — offers, brands, examples и так далее — и сопоставляют ему полное имя. Полное имя обычно собирают из KAFKA_CONTOUR (local, dev, prod и т.п.) и соглашения об именовании платформы.

И продюсер, и консьюмер принимают именно этот короткий ключ. При отправке или подписке адаптер сам резолвит его в реальное имя топика на брокере. Поэтому новый топик сначала добавляют в kafka.topics, и только потом начинают ссылаться на ключ из кода, из kafka.listen и из констант модулей. Если ключа нет в карте, запись или чтение закончатся ошибкой конфигурации, а не «тихим» уходом в несуществующий топик.

Список listen в том же конфиге отвечает на вопрос, какие ключи слушаются по умолчанию: их подхватывает встроенный runner при старте web-процесса и команда kafka:consume, если её вызвали без аргументов. Отдельно задают offset_reset — для новых consumer group обычно оставляют beginning, чтобы группа не «пропустила» историю, пока вы явно не решите иное.

Для локальной разработки, когда топика на брокере ещё нет, в .env.example часто включают KAFKA_AUTO_CREATE_TOPICS=true: клиент может создать топик при первом обращении. В общих контурах stage/prod топики обычно уже заведены инфраструктурой, и этот флаг лучше держать выключенным, чтобы случайно не плодить лишние топики с дефолтными настройками партиций.

Подключение к брокеру задаётся через переменные окружения.

  • KAFKA_BROKER_LIST — список брокеров через запятую; без него консьюмер просто не стартует.
  • KAFKA_CONSUMER_GROUP_ID задаёт группу консьюмера; если переменная пустая, берётся APP_NAME.
  • KAFKA_CONSUMER_ENABLED включает встроенный консьюмер вместе с HTTP-приложением.
  • Для защищённых кластеров используют KAFKA_SECURITY_PROTOCOL и набор KAFKA_SASL_*.
  • Сжатие исходящих сообщений настраивают через KAFKA_PRODUCER_COMPRESSION_CODEC (часто snappy).

Как отправить сообщение

Публикация делается через фасад, из хэндлера, artisan-команды или джобы очереди — из любого места, где есть контекст запроса или фоновой задачи:

err := facades.Kafka().Produce(ctx, "topic-key", payload)

После контекста передаётся ключ топика из kafka.topics, затем тело сообщения. Как сериализуется тело, зависит от типа значения. Если передать структуру или map, адаптер закодирует её в JSON. Если передать уже готовые []byte или строку, байты уйдут в брокер как есть.

Ключ сообщения Kafka (partition key) задавать не обязательно. Он нужен, когда сообщения с одним и тем же бизнес-идентификатором должны попадать в одну партицию и обрабатываться строго по порядку — например, все события одного оффера:

err := facades.Kafka().Produce(
ctx,
"offers",
payload,
appkafka.WithKey(strconv.Itoa(offerID)),
)

Как читать сообщения

Входящие сообщения обрабатывает хэндлер — обычная структура с методом Handle(ctx context.Context, msg appkafka.Message) error. Хэндлеры конкретной предметной области кладут в app/modules/{module}/kafka/, рядом с actions и остальным кодом модуля, а не в adapters: adapters дают транспорт, модуль знает, что делать с полезной нагрузкой.

Минимальный вид хэндлера такой:

type OfferHandler struct{}

func (h *OfferHandler) Handle(ctx context.Context, msg appkafka.Message) error {
var payload OfferChanged
if err := msg.JSON(&payload); err != nil {
return err
}
// дальше — бизнес-логика, чаще всего через action или sync модели
return nil
}

Метод msg.JSON разбирает тело как JSON в переданную структуру. Если нужен сырой вид, байты доступны в msg.Value, ключ сообщения — в msg.Key. Поля Topic, TopicKey, номер партиции и offset полезны для логов и диагностики: по ним видно, какое именно сообщение взяли в работу.

С точки зрения консьюмера сообщение считается обработанным только когда Handle вернул nil. Если метод вернул ошибку, offset за это сообщение не сдвигают. После перезапуска консьюмера (или после повторной выдачи в рамках политики ретраев) то же сообщение придёт снова. Поэтому обработку стоит писать идемпотентной: повторный create не должен плодить дубли, повторный update не должен портить уже согласованное состояние, повторный delete при отсутствующей записи обычно считается успешным исходом.

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

Регистрация хэндлера

Мало написать Handle: хэндлер нужно зарегистрировать в registry. Иначе при старте консьюмер откажется слушать топик — для ключа из listen не окажется обработчика. Регистрация сосредоточена в bootstrap/kafka.go: туда передают ключ топика и экземпляр хэндлера.

registry.Register(brandskafka.TopicKey, brandskafka.NewHandler())
registry.Register(offerskafka.TopicKey, offerskafka.NewHandler())

События изменения моделей (model events)

Многие топики Ensi несут не произвольный JSON «на усмотрение автора», а стандартизованное событие изменения модели. В теле есть тип события (create, update или delete), список полей, которые изменились (dirty), и полный снимок сущности после события (attributes). Такой конверт уже описан в адаптере (ModelEvent, DecodeModelEvent), и модульные хэндлеры обычно декодируют его в типизированный payload:

{
"event": "update",
"dirty": ["name", "updated_at"],
"attributes": { }
}

Типичный ход обработки такой. Хэндлер декодирует envelope, открывает транзакцию ORM и ищет локальную запись по внешнему идентификатору мастер-системы (brand_id, offer_id и т.п.). Для create/update при отсутствии строки создаёт новую, заполняет поля из attributes и сохраняет; для delete удаляет, если запись ещё есть. Часто после успешного сохранения помечают связанные сущности к переиндексации или другой отложенной работе. На update многие хэндлеры смотрят на пересечение dirty с «интересными» полями: если пришло только то, что локально не используется, сохранение и побочные эффекты можно пропустить.

Повторяющийся каркас find / fill / save / delete в сервисах часто выносят в общий sync-action, чтобы каждый топик описывал только маппинг своей сущности. Как устроены модели и транзакции — в работе с базой данных. Как писать component-тесты на такие хэндлеры (вызов Handle с синтетическим Message, без живого брокера) — в автотестах.

Как запускают консьюмеры

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

Первый — вместе с web-приложением. Если выставлены KAFKA_CONSUMER_ENABLED=true и непустой KAFKA_BROKER_LIST, при старте сервиса поднимается встроенный runner. Он подписывается на ключи из kafka.listen, обычно одним потоком. Так удобно локально держать «всё в одном процессе»: подняли сервис — и HTTP, и консьюмер уже работают.

Второй способ — artisan-команда kafka:consume. Её запускают рядом с приложением или вместо встроенного runner:

./artisan kafka:consume
./artisan kafka:consume examples
./artisan kafka:consume offers brands --threads=3

Без аргументов команда слушает те же ключи, что указаны в kafka.listen. Аргументы позволяют сузить набор до конкретных ключей из kafka.topics — это нужно при запуске консюмеров в stage/prod окружении, чтобы в разных k8s подах слушались разные топики. Флаг --threads задаёт, сколько горутин-консюмеров будет создано для каждого топика. Важно понимать - если мы запускаем три пода и по две горутины в каждом, то всего будет 6 консюмеров.

Идентификатор группы берётся из KAFKA_CONSUMER_GROUP_ID, а если он не задан — из APP_NAME. Несколько процессов с одним и тем же group id делят партиции между собой и не дублируют обработку одного сообщения внутри группы. Процессы с разными group id читают одни и те же топики независимо: каждый ведёт свой offset.

При старте в stdout обычно печатается список ключей и соответствующих им полных имён топиков — по нему сразу видно, не ошиблись ли в конфиге. По каждому обработанному сообщению пишется прогресс в духе RUNNING / DONE / FAIL вместе с ключом топика и offset.

Остановить консьюмер можно Ctrl+C локально или SIGTERM в Kubernetes. Чтобы быстро проверить связку end-to-end на учебной фиче, в одном терминале запускают kafka:consume, в другом — команду публикации в тестовый топик: консьюмер должен взять сообщение и завершить Handle без ошибки.

Ошибки и повторная обработка

Если Handle вернул ошибку, в прогресс-логе это отразится как FAIL, а offset за текущее сообщение не сдвинется вперёд. После рестарта консьюмера (или когда группа снова получит это сообщение) обработка начнётся с него же. Уже успешно обработанные сообщения и работа соседних партиций из-за одной ошибки не откатываются и не «пропадают»: страдает только прогресс на конкретной партиции, где сорвался хэндлер.

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

Остановка без обрыва текущего сообщения

При получении сигнала остановки (SIGINT или SIGTERM) консьюмер не бросает текущее сообщение на середине. Он даёт уже начатому Handle завершиться, но новые сообщения из брокера больше не берёт. Так проще сохранить согласованность: либо обработка дошла до конца и offset можно сдвинуть, либо процесс убили уже после graceful-фазы — и тогда сообщение обработается снова после старта нового пода.

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