gitaspen docs

Обмен сообщениями между сервисами: Kafka

Как довести обмен от «сервисы дёргают друг друга запросами» до состояния «работает брокер, темы созданы с известными настройками, события публикуются без потерь, неразобранное складывается отдельно и разбирается». Документ описывает контракт сообщений и постановку Kafka в контейнерах рядом с приложением.

45 минут

Как довести обмен от «сервисы дёргают друг друга запросами» до состояния «работает брокер, темы созданы с известными настройками, события публикуются без потерь, неразобранное складывается отдельно и разбирается». Документ описывает контракт сообщений и постановку Kafka в контейнерах рядом с приложением.

Исходное состояние: сервер с Docker, приложение во внутренней сети Compose, база данных сервиса работает, брокера нет.

Что нужно до начала:

УсловиеПроверкаОжидается
Docker и Compose работаютdocker compose versionDocker Compose version v2.…
внутренняя сеть приложения существуетdocker network lsсеть приложения в списке
база сервиса принимает соединения (нужна под таблицу исходящих и отметки об обработанном)docker compose exec db pg_isready -U appuser -d appdbaccepting connections
свободная память под брокерfree -gне меньше 2 ГБ свободно
место под журналы темdf -h /var/lib/dockerзапас не меньше расчётного (см. шаг 2)
часы сервера синхронизированыtimedatectl show -p NTPSynchronized --valueyes

Часы важны потому, что срок хранения тем и отметки времени в сообщениях считаются по системным часам: расхождение между серверами делает сравнение отметок бессмысленным, а хранение — непредсказуемым.

Место в цепочке

Откуда пришлиЭтот документКуда ведёт
Docker, сетевой контур, база данныхконтракт сообщений, брокер, темы, публикация без потерь, разбор неразобранногозапуск приложения с брокером в составе, затем наблюдение: задержка потребителя и рост тем

Что предыдущее звено обязано обеспечить:

  • внутреннюю сеть Compose — брокер общается с сервисами по имени сервиса, портов на хосте не публикует. Правило «база данных, очередь — без публикации» задано в сетевом контуре;
  • базу с отдельной схемой сервиса — в ней живут таблица исходящих (шаг 5) и отметки об обработанном (шаг 6). Обе таблицы принадлежат сервису и создаются его миграциями;
  • именованные тома под данные — журналы тем переживают пересоздание контейнера только в томе.

Что этот документ оставляет следующему: адрес брокера в переменных окружения сервисов; темы с закреплённым числом разделов и сроком хранения; по группе потребителя на сервис — имена групп нужны наблюдению, чтобы измерять задержку; тему для неразобранного, ненулевое наполнение которой является поводом для оповещения.

Порядок принципиален. Автосоздание тем выключено (шаг 1), поэтому темы создаются до первого запуска потребителей (шаг 2). Обратный порядок даёт тему с настройками по умолчанию — один раздел, одна копия — и число разделов после этого можно только увеличить, а срок хранения придётся исправлять на уже накопленных данных.


Когда очередь, а когда запрос

ПризнакЗапросОчередь
нужен ответ немедленноданет
отправитель ждёт результатданет
получателей может быть нескольконетда
допустимо выполнить чуть позженетда
получатель может быть недоступенвызов упадётсообщение дождётся

Очередь развязывает сервисы по времени и доступности: отправитель не знает, кто и когда прочитает. Плата за это — отсутствие немедленного ответа и необходимость учитывать повторную доставку.

Форма сообщения

Сообщение состоит из трёх частей: заголовки, тело и ключ.

Заголовки — служебные поля, одинаковые для всех сообщений:

json
{
  "message_type": "order_created",
  "entity_id": "A-1024",
  "message_id": "018f3a2c-6f21-7c6a-9a10-2b4f7d5e91c3",
  "timestamp": "2026-01-15T12:34:56Z"
}

Тело — полезные данные, вложенные в одно поле:

json
{ "payload": { "items": 3, "total": 1500 } }

Ключ — идентификатор сущности, к которой относится сообщение.

Правила формы:

  • имена полей в едином стиле (snake_case), один стиль во всех темах;
  • время — в ISO 8601 и в UTC (Z), чтобы сообщения от разных серверов сравнивались напрямую;
  • тип сообщения — в заголовках, а не выводится из содержимого тела: потребитель решает по нему, разбирать ли сообщение вообще;
  • message_id — уникальный идентификатор самого сообщения, а не сущности. По нему потребитель отличает повтор от нового сообщения (шаг 6). Значение назначается один раз, при создании сообщения, и не меняется при повторной отправке.

В Kafka заголовки — это список пар «имя — байты». Значения передаются строками в UTF-8; типы в заголовках не сохраняются, число 3 придёт как "3".

Ключ и порядок

Ключ определяет, в какой раздел темы попадёт сообщение, а раздел гарантирует порядок. Поэтому ключом берут идентификатор сущности: все сообщения об одном заказе попадут в один раздел и будут обработаны по порядку.

Порядок сохраняется только внутри одного ключа. Общего порядка по всей теме нет, и рассчитывать на него нельзя.

Пустой ключ означает произвольное распределение и потерю порядка — допустимо лишь там, где порядок не важен.

Темы

Тему называют по её содержанию, а не по отправителю: сменится отправитель — имя останется верным. Разделение по назначению:

ТемаЧто несёт
командыуказание что-то сделать; один потребитель
событиясообщение о том, что уже произошло; потребителей может быть много
состоянияобновления состояния сущности
сигналы присутствияпериодические отметки «жив»

Событие описывает свершившийся факт и формулируется в прошедшем времени. Оно не предписывает получателю действий: подписчик сам решает, что делать. Это позволяет добавлять новых потребителей, не трогая отправителя.

Правила обработки

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

Порядок не гарантирован между ключами. Обработчик не должен предполагать, что сообщение о дочерней сущности придёт после сообщения о родительской.

Неизвестные поля игнорируются. Отправитель может добавить поле, не согласовывая это с каждым потребителем; потребитель читает только то, что ему нужно. Так добавление поля перестаёт быть ломающим изменением.

Ошибка обработки не теряет сообщение. Сообщение, которое не удалось обработать, уходит в отдельную тему для разбора, а не отбрасывается. Иначе сбой в обработчике означает потерю данных. Как это настраивается — шаг 4.

Сообщение не содержит секретов. В очереди оно живёт дольше запроса и доступно всем потребителям темы. Передаётся идентификатор, по которому получатель запрашивает данные сам.


Шаг 1. Брокер в контейнерах

Что именно ставится

Kafka версии 4 работает без отдельной службы хранения метаданных: роль, которую раньше исполнял внешний координатор, вынесена внутрь самой Kafka (режим KRaft). Метаданными управляют узлы с ролью controller; на одном узле обе роли — broker и controller — совмещаются в одном процессе. Отдельного контейнера-координатора в схеме нет.

УстановкаУзловРолиЧисло копий разделаmin.insync.replicas
один сервер, дев-контур1broker,controller в одном процессе11
требование к доступности33 контроллера, 3 брокера (можно совмещённые)32

Одна копия раздела означает: остановка узла — недоступность темы, потеря диска — потеря данных. Это допустимо там, где сообщение можно переиздать из таблицы исходящих (шаг 5), и недопустимо там, где очередь является единственным местом хранения факта.

Описание в Compose

Брокер добавляется в тот же compose.yml, что и сервисы, во внутреннюю сеть приложения:

yaml
services:
  broker:
    image: apache/kafka:4.0.0
    restart: unless-stopped
    environment:
      # --- роли и кворум ---
      KAFKA_NODE_ID: 1
      KAFKA_PROCESS_ROLES: broker,controller
      KAFKA_CONTROLLER_QUORUM_VOTERS: 1@broker:9093
      KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT

      # --- слушатели ---
      # LISTENERS — что брокер слушает внутри контейнера.
      # ADVERTISED_LISTENERS — адрес, который брокер возвращает клиенту при подключении.
      # Клиент после первого ответа ходит именно по этому адресу, поэтому адрес должен
      # разрешаться на стороне КЛИЕНТА, а не на стороне брокера.
      KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT

      # --- данные и хранение ---
      KAFKA_LOG_DIRS: /var/lib/kafka/data
      KAFKA_NUM_PARTITIONS: 3               # для тем, созданных без явного указания
      KAFKA_LOG_RETENTION_HOURS: 168        # 7 суток — значение по умолчанию для новых тем
      KAFKA_OFFSETS_RETENTION_MINUTES: 43200   # 30 суток: смещения простаивающих групп не пропадут

      # --- служебные темы на одном узле ---
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0

      # --- автосоздание тем выключено ---
      # Иначе первое обращение к несуществующей теме создаёт её с одним разделом,
      # одной копией и сроком хранения по умолчанию — и опечатка в имени темы
      # молча превращается в новую тему вместо ошибки.
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false"

      KAFKA_HEAP_OPTS: -Xmx1g -Xms1g
      CLUSTER_ID: ${KAFKA_CLUSTER_ID}       # см. ниже; задаётся один раз
    volumes:
      - kafka_data:/var/lib/kafka/data
    healthcheck:
      # Инструмент на JVM, запуск занимает несколько секунд — отсюда большой timeout.
      test: ["CMD-SHELL", "/opt/kafka/bin/kafka-cluster.sh cluster-id --bootstrap-server broker:9092 >/dev/null 2>&1"]
      interval: 15s
      timeout: 10s
      retries: 10
      start_period: 40s
    networks: [app_net]
    # ports не объявлены: брокер доступен только по имени `broker` внутри сети Compose

  app:
    # ...
    environment:
      KAFKA_BOOTSTRAP_SERVERS: broker:9092
    depends_on:
      broker:
        condition: service_healthy
    networks: [app_net]

volumes:
  kafka_data:

networks:
  app_net:

Разбор существенных мест:

  • Версия закреплена точным номером (4.0.0, не latest). Чтобы пересборка давала тот же образ, тег дополняют цифровым отпечатком: apache/kafka:4.0.0@sha256:….
  • Том именованный. Данные тем в томе переживают пересоздание контейнера; данные в слое контейнера исчезают вместе с ним.
  • ports отсутствуют — см. «Почему наружу не публикуется».
  • depends_on с условием service_healthy. Без него сервис стартует раньше брокера, первые публикации падают по таймауту подключения, а при включённом автосоздании тем ещё и создают темы с настройками по умолчанию.
  • CLUSTER_ID. Идентификатор кластера записывается в том при первом запуске. Если переменная не задана, образ генерирует случайное значение сам — тогда значение известно только по содержимому тома. Явное значение делает первый запуск воспроизводимым.

Идентификатор генерируется один раз и кладётся в файл окружения Compose:

bash
docker run --rm apache/kafka:4.0.0 /opt/kafka/bin/kafka-storage.sh random-uuid

Ожидается: строка из 22 символов, например MkU3OEVBNTcwNTJENDM2Qk. Значение вписывается в .env рядом с compose.yml: KAFKA_CLUSTER_ID=MkU3OEVBNTcwNTJENDM2Qk.

Права на каталог данных

Процесс в образе работает не от root, а от пользователя с идентификатором 1000. Пустой именованный том создаётся с владельцем root, поэтому при первом запуске брокер может не суметь записать в него метаданные. Каталог отдаётся владельцу один раз, до первого запуска:

bash
docker compose run --rm --user 0 --entrypoint sh broker \
  -c 'mkdir -p /var/lib/kafka/data && chown -R 1000:1000 /var/lib/kafka/data'

Запуск и проверка

bash
docker compose up -d broker
docker compose exec broker /opt/kafka/bin/kafka-cluster.sh cluster-id --bootstrap-server broker:9092

Ожидается: Cluster ID: MkU3OEVBNTcwNTJENDM2Qk — то же значение, что в .env. Другой идентификатор означает, что том был пересоздан и данные прежних тем потеряны.

bash
docker compose exec broker /opt/kafka/bin/kafka-broker-api-versions.sh --bootstrap-server broker:9092 \
  | head -1

Ожидается: строка вида broker:9092 (id: 1 rack: null) -> ( — брокер отвечает и объявляет себя под тем адресом, по которому к нему пойдут клиенты. Если здесь стоит другое имя, клиенты из соседних контейнеров подключатся к первому адресу, получат этот и уйдут в таймаут.

bash
docker compose ps --format 'table {{.Service}}\t{{.Status}}\t{{.Ports}}'

Ожидается: broker в состоянии Up (healthy), в колонке портов — 9092/tcp без стрелки 0.0.0.0:…->. Стрелка означает публикацию на хосте.

Почему наружу не публикуется

Порт брокера не публикуется ни на все интерфейсы, ни «временно, чтобы посмотреть». Причины:

  • в этой конфигурации нет проверки подлинности. Протокол PLAINTEXT не спрашивает у клиента ничего: любой, кто дотянулся до порта, читает все темы и пишет в любую из них;
  • перед брокером нельзя поставить обратный прокси. Kafka — не HTTP: клиент обязан соединяться с каждым узлом напрямую по адресу из advertised.listeners. Схема «один nginx на входе», описанная в сетевом контуре, к брокеру неприменима, и открытый порт не прикрыт ничем;
  • правила ufw не действуют на опубликованные Docker порты: свои правила Docker добавляет в iptables раньше, поэтому порт остаётся открытым снаружи даже при запрещающем правиле (см. Docker, шаг 4).

Проверка с другой машины:

bash
nc -z -w5 example.com 9092; echo $?

Ожидается: ненулевой код (соединение не установлено). Код 0 означает, что брокер доступен из интернета — публикацию нужно убрать и считать данные скомпрометированными.

Проверка на самом сервере:

bash
sudo ss -ltnp | grep 9092

Ожидается: пустой вывод.

Когда сервисы на разных серверах

Брокер должен быть доступен сервисам на других серверах — но по частной сети, а не по публичному адресу. Контейнер брокера помещается в сетевое пространство клиента частной сети тем же приёмом, что и nginx микросервиса в сетевом контуре, а объявляемый адрес меняется на адрес частной сети:

yaml
  broker:
    image: apache/kafka:4.0.0
    network_mode: "service:vpn_client"     # общее сетевое пространство с клиентом частной сети
    environment:
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://10.0.0.20:9092   # адрес в частной сети
      # остальное без изменений

Сервисы на других серверах получают KAFKA_BOOTSTRAP_SERVERS=10.0.0.20:9092. Адрес в advertised.listeners обязан совпадать с адресом, по которому клиент фактически дотягивается до брокера: клиент подключается по адресу из bootstrap, получает в ответ объявленный адрес и дальше работает только с ним.


Шаг 2. Темы

Именование

Имя темы состоит из домена и события в прошедшем времени, разделённых точкой:

Код
<домен>.<что произошло>
ПримерЧто несёт
orders.createdзаказ создан
orders.cancelledзаказ отменён
payments.requestedзапрошено списание
payments.completedсписание завершено
orders.created.deadнеразобранное из orders.created (шаг 4)

Правила:

  • строчные буквы, разделитель между словами — дефис (orders.payment-method-changed). Точка разделяет уровни, поэтому внутри уровня она не используется;
  • имя не содержит ни имени отправителя, ни имени потребителя, ни номера версии сервиса;
  • имя темы не переименовывается. Переименование означает новую тему: накопленные сообщения и смещения групп остаются у старой;
  • в имени допустимы [a-zA-Z0-9._-], длина до 249 символов. Точка и подчёркивание в именах одновременно не используются: во внутренних метриках оба знака приводятся к подчёркиванию, и orders.created с orders_created в метриках сольются.

Число разделов и копий

Разделы задают предел параллелизма: сообщения одного раздела обрабатываются строго по очереди, а разных разделов — независимо. Число копий потребителя в одной группе, работающих одновременно, не может превышать число разделов; лишние копии простаивают.

Что определяетКак выбирать
нижняя границаожидаемое число одновременно работающих копий потребителя
верхняя границакаждый раздел — отдельные файлы и отдельный поток восстановления; тысячи разделов на узел увеличивают время перезапуска
отправная точка3 раздела на тему

Число разделов можно увеличить, но нельзя уменьшить. Увеличение меняет распределение ключей по разделам: сообщения с одним ключом, отправленные до и после изменения, могут попасть в разные разделы, и порядок между ними теряется. Поэтому число выбирают с запасом сразу.

Копии раздела (replication-factor) задают, на скольких узлах хранится каждый раздел. На одном узле возможно только значение 1. На кластере из трёх узлов берут 3 копии и min.insync.replicas=2: при acks=all запись подтверждается только после записи на две копии, и остановка одного узла не останавливает публикацию.

Срок хранения

Kafka удаляет старые сообщения по сроку, а не после прочтения: прочитанное сообщение остаётся в теме и доступно другому потребителю.

Параметр темыЧто задаётЗначение по умолчанию
retention.msсколько хранить по времени604800000 (7 суток)
retention.bytesпредел размера одного раздела; -1 — без предела-1
segment.bytesразмер отрезка журнала1073741824 (1 ГБ)
cleanup.policydelete — удалять по сроку; compact — оставлять последнее значение на ключdelete

Удаление происходит отрезками: активный отрезок не удаляется, пока не закроется. Поэтому фактический объём темы превышает расчётный на размер одного активного отрезка на раздел, а тема с малым потоком может хранить сообщения дольше retention.ms.

Расчёт места: поток × средний размер × срок × число копий раздела. Для 10 сообщений в секунду по 2 КБ, срока 7 суток и одной копии это 10 × 2048 × 604800 ≈ 12 ГБ на тему.

Темы состояний (по таблице в разделе «Темы») настраиваются на уплотнение — cleanup.policy=compact. Уплотнение оставляет по каждому ключу последнее значение, поэтому тема перестаёт расти вместе с числом обновлений и остаётся полным снимком состояния. Требования: ключ обязателен (сообщения без ключа не уплотняются), удаление сущности выражается сообщением с тем же ключом и пустым телом.

Создание и просмотр

bash
docker compose exec broker /opt/kafka/bin/kafka-topics.sh --bootstrap-server broker:9092 \
  --create --topic orders.created \
  --partitions 3 --replication-factor 1 \
  --config retention.ms=604800000 \
  --config cleanup.policy=delete

Ожидается: Created topic orders.created.

Повторный запуск той же команды завершается ошибкой Topic 'orders.created' already exists — это ожидаемо и означает, что тема на месте. Чтобы команду можно было выполнять повторно, добавляют --if-not-exists.

bash
docker compose exec broker /opt/kafka/bin/kafka-topics.sh --bootstrap-server broker:9092 \
  --describe --topic orders.created
Ожидается:
Topic: orders.created	TopicId: 8Kx1nQ2sTZ6bYw0pL3aRfg	PartitionCount: 3	ReplicationFactor: 1	Configs: cleanup.policy=delete,retention.ms=604800000
	Topic: orders.created	Partition: 0	Leader: 1	Replicas: 1	Isr: 1	Elr:	LastKnownElr:
	Topic: orders.created	Partition: 1	Leader: 1	Replicas: 1	Isr: 1	Elr:	LastKnownElr:
	Topic: orders.created	Partition: 2	Leader: 1	Replicas: 1	Isr: 1	Elr:	LastKnownElr:

Проверяется три значения: PartitionCount и ReplicationFactor совпадают с заданными, а в Configs присутствует retention.ms. Пустой Configs означает, что тема создана без явных настроек и живёт на значениях брокера по умолчанию.

Список всех тем:

bash
docker compose exec broker /opt/kafka/bin/kafka-topics.sh --bootstrap-server broker:9092 --list

Ожидается: имена созданных тем. Служебные темы (__consumer_offsets и подобные) в выводе не показываются.

Изменение настроек существующей темы:

bash
# срок хранения
docker compose exec broker /opt/kafka/bin/kafka-configs.sh --bootstrap-server broker:9092 \
  --entity-type topics --entity-name orders.created \
  --alter --add-config retention.ms=2592000000

# число разделов (только в сторону увеличения)
docker compose exec broker /opt/kafka/bin/kafka-topics.sh --bootstrap-server broker:9092 \
  --alter --topic orders.created --partitions 6

Проверка: сообщение публикуется и читается

bash
docker compose exec -T broker /opt/kafka/bin/kafka-console-producer.sh \
  --bootstrap-server broker:9092 --topic orders.created \
  --property parse.key=true --property key.separator=: <<'EOF'
A-1024:{"payload":{"items":3,"total":1500}}
EOF

Ожидается: команда завершается без вывода.

bash
docker compose exec broker /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server broker:9092 --topic orders.created \
  --from-beginning --max-messages 1 \
  --property print.key=true --property print.headers=true
Ожидается:
NO_HEADERS	A-1024	{"payload":{"items":3,"total":1500}}
Processed a total of 1 messages

NO_HEADERS здесь корректен: консольный отправитель заголовков не ставит. Сообщения, отправленные приложением, покажут заголовки в виде message_type:order_created,entity_id:A-1024,….


Шаг 3. Потребитель: группа, смещения, повторы

Группа — имя, под которым сервис читает тему. Брокер распределяет разделы между копиями одной группы и хранит для группы позицию чтения (смещение) по каждому разделу. Две разные группы читают одну тему независимо и получают каждое сообщение обе.

Правила:

  • одно имя группы на сервис, одинаковое для всех его копий. Имя задаётся явно и не выводится из имени хоста, идентификатора контейнера или случайного значения: при перезапуске такое имя меняется, группа считается новой и тема перечитывается с начала;
  • имя группы совпадает с именем сервиса: billing-service, notifications-service;
  • число одновременно работающих копий в группе не превышает число разделов темы;
  • если одну тему обрабатывают два разных сервиса — это две группы, и так и задумано. Если два экземпляра одного сервиса оказались в разных группах — работа выполняется дважды (см. «Типичные отказы»).

Смещение хранится в брокере, а не в процессе потребителя. Поэтому перезапуск и переезд потребителя на другой сервер не теряют позицию, а остановка потребителя не приводит к пропуску сообщений: они дождутся в теме, пока не истечёт срок хранения.

Автоматическая фиксация смещения выключается: она двигает позицию по расписанию, независимо от того, завершилась ли обработка. Отказ обработчика между автоматическими фиксациями означает пропущенное сообщение. Смещение двигается вручную — после обработки.

python
# consumer.py — потребитель одной темы
import json
import logging
import os
import time
from confluent_kafka import Consumer

TOPIC = "orders.created"
GROUP = "billing-service"      # одинаково для всех копий сервиса
MAX_ATTEMPTS = 3               # попытки обработки одного сообщения
BACKOFF_BASE = 0.5             # секунды; пауза удваивается с каждой попыткой

consumer = Consumer({
    "bootstrap.servers": os.environ["KAFKA_BOOTSTRAP_SERVERS"],
    "group.id": GROUP,
    "enable.auto.commit": False,      # смещение двигаем сами, после обработки
    "auto.offset.reset": "earliest",  # новая группа читает тему с начала, а не с конца
    "max.poll.interval.ms": 300000,   # предел на обработку одного сообщения — 5 минут
})
consumer.subscribe([TOPIC])

while True:
    msg = consumer.poll(1.0)
    if msg is None:
        continue
    if msg.error():
        logging.error("ошибка чтения: %s", msg.error())
        continue

    for attempt in range(MAX_ATTEMPTS):
        try:
            handle(json.loads(msg.value()))
            break
        except Exception as exc:
            if attempt + 1 < MAX_ATTEMPTS:
                time.sleep(BACKOFF_BASE * 2 ** attempt)
                continue
            to_dead_letter(msg, exc)   # попытки исчерпаны — см. шаг 4

    consumer.commit(msg)   # фиксируется offset + 1: следующее чтение начнётся со следующего

Три свойства этого цикла:

  • повторы конечны. Число попыток ограничено, пауза между ними растёт. Без ограничения неразбираемое сообщение повторяется бесконечно и блокирует свой раздел;
  • повтор происходит на месте. Сообщение не откладывается в конец очереди, поэтому порядок внутри раздела сохраняется. Плата — раздел стоит на время повторов (при значениях выше — до 1,5 с);
  • смещение двигается в обоих исходах — и после успешной обработки, и после отправки в тему для неразобранного. Не двигать смещение при отказе означает вечно перечитывать одно и то же сообщение.

max.poll.interval.ms — предел времени между обращениями к брокеру. Обработчик, который работает дольше, приводит к исключению копии из группы и перераспределению разделов; сообщение при этом достанется другой копии и будет обработано повторно.

Проверка: потребитель встал на нужное смещение.
docker compose exec broker /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server broker:9092 \
  --describe --group billing-service
Ожидается:
GROUP            TOPIC           PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG  CONSUMER-ID          HOST         CLIENT-ID
billing-service  orders.created  0          14              14              0    rdkafka-1a2b…        /10.0.1.14   rdkafka
billing-service  orders.created  1          9               9               0    rdkafka-1a2b…        /10.0.1.14   rdkafka
billing-service  orders.created  2          11              11              0    rdkafka-1a2b…        /10.0.1.14   rdkafka

Что читается из вывода:

КолонкаЗначение
CURRENT-OFFSETдокуда группа дочитала
LOG-END-OFFSETсколько всего записано в раздел
LAGразность — сколько не прочитано
CONSUMER-ID пустой или -к разделу никто не подключён: потребитель остановлен или копий меньше, чем разделов

Растущий LAG при живом потребителе — см. «Типичные отказы».


Шаг 4. Тема для неразобранного

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

Создание

Имя — имя исходной темы с суффиксом .dead. Одного раздела достаточно: тема читается вручную, порядок в ней значения не имеет. Срок хранения больше, чем у исходной темы, иначе накопленное исчезнет раньше, чем до него дойдут руки.

bash
docker compose exec broker /opt/kafka/bin/kafka-topics.sh --bootstrap-server broker:9092 \
  --create --topic orders.created.dead \
  --partitions 1 --replication-factor 1 \
  --config retention.ms=2592000000        # 30 суток

Ожидается: Created topic orders.created.dead.

Тема создаётся вместе с исходной, а не после первого отказа: при выключенном автосоздании тем попытка записать в несуществующую тему завершится ошибкой, и сообщение будет потеряно именно в тот момент, когда его требовалось сохранить.

Кто и что туда пишет

python
# продолжение consumer.py
from datetime import datetime, timezone
from confluent_kafka import Producer

producer = Producer({
    "bootstrap.servers": os.environ["KAFKA_BOOTSTRAP_SERVERS"],
    "enable.idempotence": True,
    "acks": "all",
})
DEAD_TOPIC = f"{TOPIC}.dead"


def to_dead_letter(msg, exc: Exception) -> None:
    """Переложить неразобранное сообщение в отдельную тему, не меняя тела."""
    headers = list(msg.headers() or [])
    headers += [
        ("dead-origin-topic",     msg.topic().encode()),
        ("dead-origin-partition", str(msg.partition()).encode()),
        ("dead-origin-offset",    str(msg.offset()).encode()),
        ("dead-consumer-group",   GROUP.encode()),
        ("dead-attempts",         str(MAX_ATTEMPTS).encode()),
        ("dead-error",            f"{type(exc).__name__}: {exc}"[:500].encode()),
        ("dead-at", datetime.now(timezone.utc).isoformat(timespec="seconds").encode()),
    ]
    producer.produce(DEAD_TOPIC, key=msg.key(), value=msg.value(), headers=headers)
    producer.flush(10)      # ждём подтверждения: без него смещение сдвинется раньше записи
    logging.error("сообщение отправлено в %s: %s", DEAD_TOPIC, exc)

Существенное:

  • тело не меняется. В тему кладутся исходные байты, чтобы сообщение можно было переиграть как есть. Разбор причины хранится в заголовках, а не подмешивается в тело;
  • ключ сохраняется. При возврате в исходную тему ключ определит раздел, а значит и порядок;
  • flush до фиксации смещения. Сначала подтверждение от брокера, потом commit. Обратный порядок теряет сообщение, если процесс остановится между ними;
  • если запись в тему для неразобранного не удалась, смещение двигать нельзя: сообщение остаётся в исходной теме и будет прочитано снова.

Заголовки: dead-origin-topic, dead-origin-partition, dead-origin-offset дают точное место исходного сообщения; dead-consumer-group — какой сервис не справился (одну тему могут читать несколько групп, и не справиться могла любая); dead-attempts и dead-error — сколько попыток сделано и на чём.

Как оттуда разбирают

Чтение с заголовками:

bash
docker compose exec broker /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server broker:9092 --topic orders.created.dead \
  --from-beginning --timeout-ms 5000 \
  --property print.headers=true --property print.key=true

Ожидается: строки вида

Код
dead-origin-topic:orders.created,dead-origin-partition:1,dead-origin-offset:57,dead-consumer-group:billing-service,dead-attempts:3,dead-error:KeyError: 'total',dead-at:2026-01-15T12:40:11+00:00	A-1024	{"payload":{"items":3}}
Processed a total of 4 messages

Сколько накопилось:

bash
docker compose exec broker /opt/kafka/bin/kafka-get-offsets.sh \
  --bootstrap-server broker:9092 --topic orders.created.dead

Ожидается: orders.created.dead:0:4 — в разделе 0 записано 4 сообщения.

Разбор ведётся по dead-error: сообщения группируются по причине, а не разбираются поштучно — одинаковая ошибка на сотне сообщений означает один дефект.

Исходов три:

ИсходКогдаЧто делать
дефект в обработчикетело сообщения корректно, обработчик его не принимаетисправить обработчик, выкатить, переиграть сообщения в исходную тему
дефект в отправителетело сообщения не соответствует контрактуисправить отправителя; для накопленного — переиздать события из таблицы исходящих (шаг 5) или признать невосстановимым
восстановление невозможносущность уже удалена, событие потеряло смыслоставить в теме до истечения срока хранения; решение записать

Переигрывание возвращает сообщение в исходную тему, а не обрабатывает его прямо из темы для неразобранного. Иначе обработка раздваивается на два пути, и второй остаётся без тестов и без наблюдения.

python
# replay.py — вернуть неразобранное в исходную тему
import os
from confluent_kafka import Consumer, Producer

SOURCE = "orders.created"
DEAD = f"{SOURCE}.dead"
SERVICE_HEADERS = {b"dead-origin-topic", b"dead-origin-partition", b"dead-origin-offset",
                   b"dead-consumer-group", b"dead-attempts", b"dead-error", b"dead-at"}

bootstrap = os.environ["KAFKA_BOOTSTRAP_SERVERS"]
consumer = Consumer({
    "bootstrap.servers": bootstrap,
    "group.id": "replay-tool",         # своя группа: не мешает рабочим потребителям
    "enable.auto.commit": False,
    "auto.offset.reset": "earliest",
})
producer = Producer({"bootstrap.servers": bootstrap, "enable.idempotence": True, "acks": "all"})
consumer.subscribe([DEAD])

moved = 0
while True:
    msg = consumer.poll(5.0)
    if msg is None:
        break
    if msg.error():
        continue
    # служебные заголовки снимаются: при повторном отказе они будут добавлены заново,
    # иначе в сообщении накопятся несколько разных `dead-error`
    headers = [(k, v) for k, v in (msg.headers() or []) if k.encode() not in SERVICE_HEADERS]
    producer.produce(SOURCE, key=msg.key(), value=msg.value(), headers=headers)
    moved += 1

producer.flush(30)
consumer.commit()      # смещение в теме неразобранного двигается только после отправки
print(f"переиграно сообщений: {moved}")

Переигрывать следует после выката исправления, иначе сообщения вернутся обратно тем же путём.

Что делать с накопившимся. Постоянно непустая тема для неразобранного — не «разберём позже», а работающий контракт, который перестал соблюдаться. Порог оповещения ставится на любое ненулевое число сообщений в ней (см. наблюдение); чем дольше разбор откладывается, тем вероятнее, что сообщения станут невосстановимыми: сущности удалены, суммы пересчитаны, срок хранения исходной темы истёк.


Шаг 5. Публикация после фиксации

Событие публикуется после того, как изменение зафиксировано в базе. Публикация внутри открытой транзакции создаёт расхождение: при откате событие уже ушло, и потребители узнают о том, чего не произошло.

Обратный порядок — сначала запись, потом публикация — оставляет другой риск: процесс может завершиться между ними, и событие не уйдёт. Где такая потеря недопустима, событие записывают в ту же транзакцию (отдельной таблицей) и публикуют отдельным шагом, который повторяется до успеха.

СпособЧто теряетсяКогда применим
публикация внутри транзакциисобытие уходит и при откатене применяется
запись, затем публикация из того же кодасобытие не уходит, если процесс упал между нимипотеря события допустима: уведомление, метрика
таблица исходящих и отдельный публикаторничего; возможен повторпотеря события недопустима

Третий способ описан ниже. Он даёт доставку «хотя бы один раз»: событие не пропадёт, но может уйти дважды — на этом и построено требование идемпотентности потребителя (шаг 6).

Схема таблицы

Таблица создаётся миграцией сервиса, в его схеме (см. база данных, шаг 3), рядом с предметными таблицами — она обязана попадать в ту же транзакцию.

sql
CREATE TABLE outbox (
    id           bigint      GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    topic        text        NOT NULL,
    message_key  text,
    headers      jsonb       NOT NULL DEFAULT '{}'::jsonb,
    payload      jsonb       NOT NULL,
    status       text        NOT NULL DEFAULT 'pending',
    attempts     integer     NOT NULL DEFAULT 0,
    last_error   text,
    created_at   timestamptz NOT NULL DEFAULT now(),
    available_at timestamptz NOT NULL DEFAULT now(),
    sent_at      timestamptz,
    CONSTRAINT outbox_status_check CHECK (status IN ('pending', 'sent', 'failed'))
);

-- Частичный индекс: в него попадают только неотправленные записи.
-- Таблица растёт, индекс остаётся размером с очередь на отправку.
CREATE INDEX outbox_pending_idx ON outbox (available_at, id) WHERE status = 'pending';
ПолеНазначение
idпорядок отправки; возрастает вместе с порядком записи
topicкуда отправлять; хранится в записи, а не выводится из типа события в коде публикатора
message_keyключ сообщения — идентификатор сущности (см. «Ключ и порядок»)
headersслужебные поля сообщения, включая message_id
payloadтело сообщения
statusстадия: см. таблицу ниже
attemptsсколько раз отправка не удалась; по нему считается пауза до следующей попытки
last_errorтекст последней ошибки отправки
available_atмомент, раньше которого запись не берут; отодвигается при отказе
sent_atкогда брокер подтвердил запись; по нему чистится таблица
СтатусЧто означаетКто ставит
pendingзаписано в транзакции с изменением, ещё не отправлено или отправка не удаласькод сервиса при записи; публикатор при откате попытки
sentброкер подтвердил запись во все требуемые копии разделапубликатор
failedпопытки исчерпаны; требуется вмешательствопубликатор

Запись события в той же транзакции, что и само изменение:

sql
BEGIN;

INSERT INTO orders (order_number, total, status)
VALUES ('A-1024', 1500, 'created');

INSERT INTO outbox (topic, message_key, headers, payload)
VALUES ('orders.created', 'A-1024',
        '{"message_type":"order_created","entity_id":"A-1024",
          "message_id":"018f3a2c-6f21-7c6a-9a10-2b4f7d5e91c3",
          "timestamp":"2026-01-15T12:34:56Z"}',
        '{"payload":{"items":3,"total":1500}}');

COMMIT;

При откате транзакции откатываются обе вставки: события о несостоявшемся изменении не остаётся.

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

РазмещениеКогдаЧто учесть
фоновая задача в процессе сервисаодин сервис, небольшой поток событийостанавливается вместе с сервисом; при N копиях сервиса работают N публикаторов — блокировка обязательна
отдельный процесс (свой контейнер)поток большой либо публикация не должна конкурировать с обработкой запросовтот же образ, другая команда запуска; масштабируется и перезапускается отдельно

В обоих случаях порядок один: взял → отправил → пометил.

  1. Взял

    Транзакция открывается, из таблицы выбирается пачка записей pending со сроком available_at <= now(), с блокировкой FOR UPDATE SKIP LOCKED.

  2. Отправил

    Записи уходят в брокер, публикатор дожидается подтверждения. Транзакция всё это время открыта — блокировка держится, вторая копия эти строки не увидит.

  3. Пометил

    Подтверждённые получают status='sent', неудавшиеся — увеличенный attempts и отодвинутый available_at. Транзакция закрывается, блокировка снимается.

Почему именно такой порядок: пометка до отправки теряет событие навсегда, если процесс остановится между пометкой и отправкой. Пометка после отправки в том же случае даёт повтор — запись останется pending и уйдёт второй раз. Повтор поглощается идемпотентным потребителем (шаг 6), потеря — нет.

Блокировка. FOR UPDATE SKIP LOCKED заставляет вторую копию публикатора пропустить уже занятые строки и взять следующие.

Как выбиратьЧто происходит при двух копиях
без FOR UPDATEобе копии берут одни и те же строки и отправляют их дважды
FOR UPDATE без SKIP LOCKEDвторая копия ждёт освобождения строк; дублей нет, но работа идёт последовательно
FOR UPDATE SKIP LOCKEDкопии берут непересекающиеся пачки и работают параллельно

Пример

python
# outbox_publisher.py — отправка записей из таблицы исходящих
import json
import logging
import os
import signal
import time

import psycopg
from confluent_kafka import Producer

BATCH = 100          # записей за проход
IDLE_PAUSE = 0.5     # пауза, когда отправлять нечего
MAX_ATTEMPTS = 10    # после этого запись помечается failed и требует вмешательства
FLUSH_TIMEOUT = 30   # ожидание подтверждений от брокера

producer = Producer({
    "bootstrap.servers": os.environ["KAFKA_BOOTSTRAP_SERVERS"],
    "enable.idempotence": True,   # брокер отбрасывает повтор одной и той же записи от продюсера
    "acks": "all",                # подтверждение после записи во все синхронные копии раздела
    "linger.ms": 5,
})

SELECT_BATCH = """
SELECT id, topic, message_key, headers, payload
  FROM outbox
 WHERE status = 'pending' AND available_at <= now()
 ORDER BY id
 LIMIT %s
 FOR UPDATE SKIP LOCKED
"""

MARK_FAILED = """
UPDATE outbox
   SET attempts     = attempts + 1,
       last_error   = %s,
       available_at = now() + make_interval(secs => least(300, power(2, attempts)::int)),
       status       = CASE WHEN attempts + 1 >= %s THEN 'failed' ELSE 'pending' END
 WHERE id = %s
"""


def publish_batch(conn) -> int:
    with conn.transaction():                       # блокировка держится до конца транзакции
        with conn.cursor() as cur:
            cur.execute(SELECT_BATCH, (BATCH,))
            rows = cur.fetchall()
            if not rows:
                return 0

            results: dict[int, str | None] = {}

            def on_delivery(err, _msg, row_id):
                results[row_id] = str(err) if err else None

            for row_id, topic, key, headers, payload in rows:
                producer.produce(
                    topic=topic,
                    key=key.encode() if key else None,
                    value=json.dumps(payload).encode(),
                    headers=[(k, str(v).encode()) for k, v in headers.items()],
                    on_delivery=lambda err, msg, rid=row_id: on_delivery(err, msg, rid),
                )

            producer.flush(FLUSH_TIMEOUT)           # после этого исход каждой записи известен

            sent = [rid for rid, err in results.items() if err is None]
            failed = {rid: err for rid, err in results.items() if err is not None}

            if sent:
                cur.execute(
                    "UPDATE outbox SET status = 'sent', sent_at = now() WHERE id = ANY(%s)",
                    (sent,),
                )
            for row_id, err in failed.items():
                cur.execute(MARK_FAILED, (err[:500], MAX_ATTEMPTS, row_id))

            # Записи без подтверждения за FLUSH_TIMEOUT остаются pending и будут взяты
            # следующим проходом: отсюда возможен повтор отправки.
            return len(rows)


running = True


def stop(*_):
    global running
    running = False


signal.signal(signal.SIGTERM, stop)
signal.signal(signal.SIGINT, stop)

with psycopg.connect(os.environ["DATABASE_URL"]) as conn:
    while running:
        try:
            taken = publish_batch(conn)
        except Exception:
            logging.exception("проход публикатора не удался")
            taken = 0
        if taken < BATCH:            # взяли меньше пачки — очередь разобрана, ждём
            time.sleep(IDLE_PAUSE)

producer.flush(FLUSH_TIMEOUT)

Существенное:

  • enable.idempotence и acks=all у отправителя. Первое отсекает повтор, возникший при внутреннем повторе отправки в клиенте; второе означает, что подтверждение приходит после записи во все синхронные копии раздела. Без acks=all подтверждение приходит от одного узла, и потеря узла теряет подтверждённое сообщение.
  • Пачка ограничена. Транзакция открыта всё время отправки; чем больше пачка, тем дольше она держит блокировки и соединение.
  • Пауза до следующей попытки растёт (power(2, attempts)) и ограничена 300 секундами. Без ограничения недоступный брокер превращается в цикл без пауз, который нагружает базу.
  • failed не отправляется автоматически. Запись в этом статусе означает, что дело не в доступности брокера: несуществующая тема, сообщение больше предельного размера, неверный формат. После устранения причины записи возвращают в работу: UPDATE outbox SET status='pending', attempts=0, available_at=now() WHERE status='failed'.
Проверка публикатора:
SELECT status, count(*), min(created_at) AS oldest FROM outbox GROUP BY status;

Ожидается: pending — единицы записей с недавним oldest, failed — ноль. Растущий pending со старым oldest означает, что публикатор не работает; ненулевой failed требует разбора по last_error.

Очистка таблицы

Отправленные записи не нужны, но удаляются не сразу: они нужны, чтобы восстановить события, если понадобится переиздание. Срок хранения берут равным сроку хранения тем.

sql
DELETE FROM outbox WHERE status = 'sent' AND sent_at < now() - interval '7 days';

Запрос ставится в суточное расписание. Без очистки таблица растёт неограниченно; частичный индекс при этом не растёт, поэтому отправка не замедляется — замедляются резервное копирование и обслуживание таблицы.


Шаг 6. Идемпотентность потребителя

Общее правило и способы его выполнения описаны в BMBP, раздел «Идемпотентность». Здесь — только то, что относится к очереди: где хранится отметка об обработанном сообщении и когда её чистят.

Отметка нужна там, где по состоянию сущности повтор не отличить: обработка не меняет состояние ( отправка уведомления), или меняет его накопительно (начисление, списание, счётчик). Там, где повтор виден по состоянию — заказ уже не в статусе «ждёт оплаты», файл уже создан — отдельная отметка не нужна, и проверка делается по состоянию.

sql
CREATE TABLE processed_messages (
    consumer_group text        NOT NULL,
    topic          text        NOT NULL,
    message_id     text        NOT NULL,
    processed_at   timestamptz NOT NULL DEFAULT now(),
    PRIMARY KEY (consumer_group, topic, message_id)
);

CREATE INDEX processed_messages_processed_at_idx ON processed_messages (processed_at);

Таблица живёт в схеме сервиса-потребителя, а не в общей: отметка принадлежит тому, кто обрабатывал. Имя группы входит в ключ, потому что одну тему могут читать несколько сервисов, и для каждого повтор считается отдельно.

Откуда берётся message_id:

ИсточникСвойства
заголовок message_idпереживает переигрывание из темы для неразобранного: идентификатор остаётся прежним
тройка «тема, раздел, смещение»доступна всегда, менять контракт не требуется; у переигранного сообщения будет другое смещение, и повтор не распознается

Отметка ставится в той же транзакции, что и результат обработки:

python
def handle(message_id: str, data: dict) -> None:
    with conn.transaction():
        rows = conn.execute(
            "INSERT INTO processed_messages (consumer_group, topic, message_id) "
            "VALUES (%s, %s, %s) ON CONFLICT DO NOTHING RETURNING 1",
            (GROUP, TOPIC, message_id),
        ).fetchall()
        if not rows:
            return          # сообщение уже обработано — повтор, выходим без работы
        apply_business_change(data)

Порядок «вставка отметки, затем работа» внутри одной транзакции: при откате откатывается и отметка, поэтому «отметка есть, работа не выполнена» невозможно. Отметка в отдельной транзакции даёт именно это расхождение.

Когда чистят. Отметка нужна ровно столько, сколько сообщение может быть доставлено повторно, то есть пока оно лежит в теме. Срок хранения отметок берут вдвое больше срока хранения темы — запас покрывает переигрывание и остановленного потребителя:

sql
DELETE FROM processed_messages WHERE processed_at < now() - interval '14 days';

Запрос ставится в суточное расписание. Без очистки таблица растёт вместе с общим числом сообщений за всё время, а её единственный индекс — вместе с ней.

Проверка:
SELECT count(*) FROM processed_messages WHERE processed_at < now() - interval '14 days';

Ожидается: 0. Ненулевое значение означает, что очистка не выполняется.


Шаг 7. Проверка контура целиком

Проверка идёт снизу вверх; первый шаг, который не даёт ожидаемого результата, и есть место отказа.

bash
# 1. брокер отвечает
docker compose exec broker /opt/kafka/bin/kafka-cluster.sh cluster-id --bootstrap-server broker:9092

# 2. темы созданы с нужными настройками
docker compose exec broker /opt/kafka/bin/kafka-topics.sh --bootstrap-server broker:9092 \
  --describe --topic orders.created

# 3. сервис записал событие в таблицу исходящих и публикатор его отправил
docker compose exec db psql -U appuser -d appdb -c \
  "SELECT status, count(*) FROM outbox GROUP BY status"

# 4. сообщение лежит в теме
docker compose exec broker /opt/kafka/bin/kafka-get-offsets.sh \
  --bootstrap-server broker:9092 --topic orders.created

# 5. потребитель его прочитал
docker compose exec broker /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server broker:9092 \
  --describe --group billing-service

# 6. неразобранного нет
docker compose exec broker /opt/kafka/bin/kafka-get-offsets.sh \
  --bootstrap-server broker:9092 --topic orders.created.dead

# 7. порт брокера снаружи не слушается (с другой машины)
nc -z -w5 example.com 9092; echo $?
ШагОжидается
1Cluster ID: … — значение из .env
2PartitionCount: 3, ReplicationFactor: 1, в Configs есть retention.ms
3только sent; pending — единицы, failed отсутствует
4orders.created:0:… по каждому разделу, сумма растёт после публикации
5LAG равен 0 или уменьшается; CONSUMER-ID заполнен по всем разделам
6orders.created.dead:0:0
7ненулевой код возврата

Типичные отказы

ПризнакПричинаЧто делать
LAG растёт, потребитель живойобработчик медленнее потока; часто — синхронный вызов внешнего сервиса внутри обработчикаувеличить число копий потребителя (не больше числа разделов), при упоре в число разделов — увеличить разделы; вынести медленный вызов из обработчика
LAG растёт, CONSUMER-ID пустойпотребитель остановлен или упалdocker compose ps, docker compose logs; проверить, что имя группы не менялось
одно сообщение обрабатывается бесконечно, раздел стоитнет предела числа попыток: обработчик падает, смещение не двигается, сообщение читается сноваограничить попытки (шаг 3), после исчерпания — в тему для неразобранного; для уже застрявшего сдвинуть смещение вручную (см. ниже)
группа перечитывает тему с начала после простоясмещения удалены по offsets.retention.minutes (по умолчанию 7 суток)поднять offsets.retention.minutes; восстановить позицию через --reset-offsets --to-datetime
группа читает с начала при первом запускеновая группа и auto.offset.reset=earliestэто ожидаемо; чтобы начать с текущего момента, до запуска задать смещение через --reset-offsets --to-latest
диск заполнен журналами тему темы нет retention.ms или он рассчитан без учёта потока; retention.bytes задан на раздел, а не на темунайти крупные темы (du внутри контейнера); снизить retention.ms; помнить, что активный отрезок не удаляется до закрытия
работа выполняется дважды, два потребителя читают однокопии одного сервиса оказались в разных группах: group.id собран из имени хоста, идентификатора контейнера или случайного значениязадать group.id явно и одинаково для всех копий; проверить kafka-consumer-groups.sh --list — лишние похожие имена видны сразу
тема существует, но с одним разделом и без срока хранениятема создана автоматически при первом обращении, до того как её создали рукамивыключить auto.create.topics.enable; число разделов увеличить (--alter --partitions), срок задать через kafka-configs.sh; уменьшить число разделов нельзя — только пересоздать тему
потребитель подключается, затем таймаутadvertised.listeners объявляет адрес, недоступный клиентупривести объявляемый адрес к тому, по которому клиент реально дотягивается (шаг 1)
брокер не стартует: отказ записи в каталог данныхтом создан с владельцем root, процесс работает от пользователя 1000сменить владельца тома (шаг 1)
брокер не стартует: несоответствие идентификатора кластераCLUSTER_ID изменён при существующем томевернуть прежнее значение; смена идентификатора требует пустого тома и означает потерю данных
отправка падает: Not enough in-sync replicasacks=all при числе синхронных копий меньше min.insync.replicasпроверить, что все узлы кластера работают; на одном узле min.insync.replicas должен быть равен 1
копии потребителя постоянно переподключаются, сообщения обрабатываются повторнообработка одного сообщения дольше max.poll.interval.msсократить работу в обработчике либо поднять max.poll.interval.ms; проверить, не ждёт ли обработчик внешний сервис без таймаута
действие выполнено дваждыобработчик не идемпотентенпроверка по message_id или по состоянию сущности (шаг 6)
сообщения обрабатываются не по порядкуразные ключи или пустой ключключом брать идентификатор сущности
потребитель не видит сообщенийподписан на другую тему или другую группусверить имя темы и группы; kafka-topics.sh --list покажет опечатку в имени
потребитель падает на новом полестрогий разбор сообщенияигнорировать неизвестные поля

Сдвинуть смещение застрявшей группы (группа должна быть остановлена — при живых участниках команда откажет):

bash
docker compose stop app
docker compose exec broker /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server broker:9092 \
  --group billing-service --topic orders.created \
  --reset-offsets --shift-by 1 --execute
docker compose start app

Доступные варианты: --shift-by N (сдвиг на N сообщений), --to-offset N, --to-earliest, --to-latest, --to-datetime 2026-01-15T00:00:00.000. Без --execute команда только показывает, что сделает.

Журналы: docker compose logs broker — запуск брокера, выборы контроллера, ошибки записи на диск; docker compose logs app — отказы обработчиков и публикатора.


Откат и снятие

Остановить потребителя, не потеряв смещение. Смещения хранятся в брокере, а не в процессе, и остановка их не трогает:

bash
docker compose stop app          # SIGTERM: потребитель успевает зафиксировать смещение
docker compose exec broker /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server broker:9092 \
  --describe --group billing-service

Ожидается: группа в выводе, CURRENT-OFFSET заполнен, CONSUMER-ID пуст. Это означает, что позиция сохранена и при запуске чтение продолжится с неё.

Условия сохранения позиции:

  • остановка штатная (stop, а не kill -9): при принудительном завершении не зафиксированные смещения теряются, и обработанные сообщения будут прочитаны повторно — их поглотит проверка идемпотентности;
  • простой короче offsets.retention.minutes. Дольше — смещения удаляются, и группа при запуске начнёт согласно auto.offset.reset;
  • сообщения дожидаются потребителя не дольше retention.ms темы. Простой длиннее срока хранения означает пропуск сообщений независимо от смещений.
Удалить тему:
docker compose exec broker /opt/kafka/bin/kafka-topics.sh --bootstrap-server broker:9092 \
  --delete --topic orders.created

Удаление возвращает управление сразу, файлы удаляются асинхронно. Подтверждение — тема исчезла из --list.

Что при удалении темы не удаляется:

ОстаётсяГдеПочему это важно
смещения групп, читавших темуслужебная тема брокерапересозданная тема с тем же именем будет читаться со старых смещений, а они больше не соответствуют содержимому
тема для неразобранногоотдельная темаудаляется отдельной командой
таблица исходящихбаза сервисазаписи pending продолжат уходить в удалённую тему; при выключенном автосоздании отправка будет падать
отметки об обработанномбаза сервисачистятся по расписанию (шаг 6)
подписка потребителейконфигурация сервисовпотребители продолжат опрашивать несуществующую тему и писать в журнал ошибки

Поэтому снятие темы выполняется в порядке: остановить потребителей → убрать тему из конфигурации сервисов → удалить группу → удалить тему.

bash
docker compose exec broker /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server broker:9092 \
  --delete --group billing-service

Команда работает только при отсутствии активных участников группы.

Снять брокер целиком:
docker compose stop broker      # данные в томе сохраняются
docker compose rm -f broker

Том с журналами тем при этом остаётся: повторный запуск с тем же CLUSTER_ID продолжит работу с накопленными данными. Удаление данных — отдельное действие:

bash
docker volume rm <проект>_kafka_data

Оно необратимо и уничтожает все темы, сообщения и смещения. Восстановить их можно только повторной публикацией из таблиц исходящих — и только те события, которые в них ещё не очищены.

Инструкция не помогла?

Откройте исходник документа по ссылке «Предложить правку» — там же видно, что и когда в нём менялось.