Обмен сообщениями между сервисами: Kafka
Как довести обмен от «сервисы дёргают друг друга запросами» до состояния «работает брокер, темы созданы с известными настройками, события публикуются без потерь, неразобранное складывается отдельно и разбирается». Документ описывает контракт сообщений и постановку Kafka в контейнерах рядом с приложением.
45 минутКак довести обмен от «сервисы дёргают друг друга запросами» до состояния «работает брокер, темы созданы с известными настройками, события публикуются без потерь, неразобранное складывается отдельно и разбирается». Документ описывает контракт сообщений и постановку Kafka в контейнерах рядом с приложением.
Исходное состояние: сервер с Docker, приложение во внутренней сети Compose, база данных сервиса работает, брокера нет.
Что нужно до начала:
| Условие | Проверка | Ожидается |
|---|---|---|
| Docker и Compose работают | docker compose version | Docker Compose version v2.… |
| внутренняя сеть приложения существует | docker network ls | сеть приложения в списке |
| база сервиса принимает соединения (нужна под таблицу исходящих и отметки об обработанном) | docker compose exec db pg_isready -U appuser -d appdb | accepting connections |
| свободная память под брокер | free -g | не меньше 2 ГБ свободно |
| место под журналы тем | df -h /var/lib/docker | запас не меньше расчётного (см. шаг 2) |
| часы сервера синхронизированы | timedatectl show -p NTPSynchronized --value | yes |
Часы важны потому, что срок хранения тем и отметки времени в сообщениях считаются по системным часам: расхождение между серверами делает сравнение отметок бессмысленным, а хранение — непредсказуемым.
Место в цепочке
| Откуда пришли | Этот документ | Куда ведёт |
|---|---|---|
| Docker, сетевой контур, база данных | контракт сообщений, брокер, темы, публикация без потерь, разбор неразобранного | запуск приложения с брокером в составе, затем наблюдение: задержка потребителя и рост тем |
Что предыдущее звено обязано обеспечить:
- внутреннюю сеть Compose — брокер общается с сервисами по имени сервиса, портов на хосте не публикует. Правило «база данных, очередь — без публикации» задано в сетевом контуре;
- базу с отдельной схемой сервиса — в ней живут таблица исходящих (шаг 5) и отметки об обработанном (шаг 6). Обе таблицы принадлежат сервису и создаются его миграциями;
- именованные тома под данные — журналы тем переживают пересоздание контейнера только в томе.
Что этот документ оставляет следующему: адрес брокера в переменных окружения сервисов; темы с закреплённым числом разделов и сроком хранения; по группе потребителя на сервис — имена групп нужны наблюдению, чтобы измерять задержку; тему для неразобранного, ненулевое наполнение которой является поводом для оповещения.
Порядок принципиален. Автосоздание тем выключено (шаг 1), поэтому темы создаются до первого запуска потребителей (шаг 2). Обратный порядок даёт тему с настройками по умолчанию — один раздел, одна копия — и число разделов после этого можно только увеличить, а срок хранения придётся исправлять на уже накопленных данных.
Когда очередь, а когда запрос
| Признак | Запрос | Очередь |
|---|---|---|
| нужен ответ немедленно | да | нет |
| отправитель ждёт результат | да | нет |
| получателей может быть несколько | нет | да |
| допустимо выполнить чуть позже | нет | да |
| получатель может быть недоступен | вызов упадёт | сообщение дождётся |
Очередь развязывает сервисы по времени и доступности: отправитель не знает, кто и когда прочитает. Плата за это — отсутствие немедленного ответа и необходимость учитывать повторную доставку.
Форма сообщения
Сообщение состоит из трёх частей: заголовки, тело и ключ.
Заголовки — служебные поля, одинаковые для всех сообщений:
{
"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 } }Ключ — идентификатор сущности, к которой относится сообщение.
Правила формы:
- имена полей в едином стиле (
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 |
|---|---|---|---|---|
| один сервер, дев-контур | 1 | broker,controller в одном процессе | 1 | 1 |
| требование к доступности | 3 | 3 контроллера, 3 брокера (можно совмещённые) | 3 | 2 |
Одна копия раздела означает: остановка узла — недоступность темы, потеря диска — потеря данных. Это допустимо там, где сообщение можно переиздать из таблицы исходящих (шаг 5), и недопустимо там, где очередь является единственным местом хранения факта.
Описание в Compose
Брокер добавляется в тот же compose.yml, что и сервисы, во внутреннюю сеть приложения:
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:
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, поэтому при первом запуске брокер может не суметь
записать в него метаданные. Каталог отдаётся владельцу один раз, до первого запуска:
docker compose run --rm --user 0 --entrypoint sh broker \
-c 'mkdir -p /var/lib/kafka/data && chown -R 1000:1000 /var/lib/kafka/data'Запуск и проверка
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. Другой
идентификатор означает, что том был пересоздан и данные прежних тем потеряны.
docker compose exec broker /opt/kafka/bin/kafka-broker-api-versions.sh --bootstrap-server broker:9092 \
| head -1Ожидается: строка вида broker:9092 (id: 1 rack: null) -> ( — брокер отвечает и объявляет себя
под тем адресом, по которому к нему пойдут клиенты. Если здесь стоит другое имя, клиенты из соседних
контейнеров подключатся к первому адресу, получат этот и уйдут в таймаут.
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).
Проверка с другой машины:
nc -z -w5 example.com 9092; echo $?Ожидается: ненулевой код (соединение не установлено). Код 0 означает, что брокер доступен из
интернета — публикацию нужно убрать и считать данные скомпрометированными.
Проверка на самом сервере:
sudo ss -ltnp | grep 9092Ожидается: пустой вывод.
Когда сервисы на разных серверах
Брокер должен быть доступен сервисам на других серверах — но по частной сети, а не по публичному адресу. Контейнер брокера помещается в сетевое пространство клиента частной сети тем же приёмом, что и nginx микросервиса в сетевом контуре, а объявляемый адрес меняется на адрес частной сети:
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.policy | delete — удалять по сроку; compact — оставлять последнее значение на ключ | delete |
Удаление происходит отрезками: активный отрезок не удаляется, пока не закроется. Поэтому фактический
объём темы превышает расчётный на размер одного активного отрезка на раздел, а тема с малым потоком
может хранить сообщения дольше retention.ms.
Расчёт места: поток × средний размер × срок × число копий раздела. Для 10 сообщений в секунду по
2 КБ, срока 7 суток и одной копии это 10 × 2048 × 604800 ≈ 12 ГБ на тему.
Темы состояний (по таблице в разделе «Темы») настраиваются на уплотнение —
cleanup.policy=compact. Уплотнение оставляет по каждому ключу последнее значение, поэтому тема
перестаёт расти вместе с числом обновлений и остаётся полным снимком состояния. Требования: ключ
обязателен (сообщения без ключа не уплотняются), удаление сущности выражается сообщением с тем же
ключом и пустым телом.
Создание и просмотр
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.
docker compose exec broker /opt/kafka/bin/kafka-topics.sh --bootstrap-server broker:9092 \
--describe --topic orders.createdTopic: 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 означает, что тема создана без явных
настроек и живёт на значениях брокера по умолчанию.
Список всех тем:
docker compose exec broker /opt/kafka/bin/kafka-topics.sh --bootstrap-server broker:9092 --listОжидается: имена созданных тем. Служебные темы (__consumer_offsets и подобные) в выводе не
показываются.
Изменение настроек существующей темы:
# срок хранения
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Проверка: сообщение публикуется и читается
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Ожидается: команда завершается без вывода.
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=trueNO_HEADERS A-1024 {"payload":{"items":3,"total":1500}}
Processed a total of 1 messagesNO_HEADERS здесь корректен: консольный отправитель заголовков не ставит. Сообщения, отправленные
приложением, покажут заголовки в виде message_type:order_created,entity_id:A-1024,….
Шаг 3. Потребитель: группа, смещения, повторы
Группа — имя, под которым сервис читает тему. Брокер распределяет разделы между копиями одной группы и хранит для группы позицию чтения (смещение) по каждому разделу. Две разные группы читают одну тему независимо и получают каждое сообщение обе.
Правила:
- одно имя группы на сервис, одинаковое для всех его копий. Имя задаётся явно и не выводится из имени хоста, идентификатора контейнера или случайного значения: при перезапуске такое имя меняется, группа считается новой и тема перечитывается с начала;
- имя группы совпадает с именем сервиса:
billing-service,notifications-service; - число одновременно работающих копий в группе не превышает число разделов темы;
- если одну тему обрабатывают два разных сервиса — это две группы, и так и задумано. Если два экземпляра одного сервиса оказались в разных группах — работа выполняется дважды (см. «Типичные отказы»).
Смещение хранится в брокере, а не в процессе потребителя. Поэтому перезапуск и переезд потребителя на другой сервер не теряют позицию, а остановка потребителя не приводит к пропуску сообщений: они дождутся в теме, пока не истечёт срок хранения.
Автоматическая фиксация смещения выключается: она двигает позицию по расписанию, независимо от того, завершилась ли обработка. Отказ обработчика между автоматическими фиксациями означает пропущенное сообщение. Смещение двигается вручную — после обработки.
# 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-serviceGROUP 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. Одного раздела достаточно: тема читается вручную,
порядок в ней значения не имеет. Срок хранения больше, чем у исходной темы, иначе накопленное
исчезнет раньше, чем до него дойдут руки.
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.
Тема создаётся вместе с исходной, а не после первого отказа: при выключенном автосоздании тем попытка записать в несуществующую тему завершится ошибкой, и сообщение будет потеряно именно в тот момент, когда его требовалось сохранить.
Кто и что туда пишет
# продолжение 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 — сколько попыток
сделано и на чём.
Как оттуда разбирают
Чтение с заголовками:
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Сколько накопилось:
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) или признать невосстановимым |
| восстановление невозможно | сущность уже удалена, событие потеряло смысл | оставить в теме до истечения срока хранения; решение записать |
Переигрывание возвращает сообщение в исходную тему, а не обрабатывает его прямо из темы для неразобранного. Иначе обработка раздваивается на два пути, и второй остаётся без тестов и без наблюдения.
# 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), рядом с предметными таблицами — она обязана попадать в ту же транзакцию.
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 | попытки исчерпаны; требуется вмешательство | публикатор |
Запись события в той же транзакции, что и само изменение:
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 публикаторов — блокировка обязательна |
| отдельный процесс (свой контейнер) | поток большой либо публикация не должна конкурировать с обработкой запросов | тот же образ, другая команда запуска; масштабируется и перезапускается отдельно |
В обоих случаях порядок один: взял → отправил → пометил.
-
Взял
Транзакция открывается, из таблицы выбирается пачка записей
pendingсо срокомavailable_at <= now(), с блокировкойFOR UPDATE SKIP LOCKED. -
Отправил
Записи уходят в брокер, публикатор дожидается подтверждения. Транзакция всё это время открыта — блокировка держится, вторая копия эти строки не увидит.
-
Пометил
Подтверждённые получают
status='sent', неудавшиеся — увеличенныйattemptsи отодвинутыйavailable_at. Транзакция закрывается, блокировка снимается.
Почему именно такой порядок: пометка до отправки теряет событие навсегда, если процесс
остановится между пометкой и отправкой. Пометка после отправки в том же случае даёт повтор —
запись останется pending и уйдёт второй раз. Повтор поглощается идемпотентным потребителем
(шаг 6), потеря — нет.
Блокировка. FOR UPDATE SKIP LOCKED заставляет вторую копию публикатора пропустить уже занятые
строки и взять следующие.
| Как выбирать | Что происходит при двух копиях |
|---|---|
без FOR UPDATE | обе копии берут одни и те же строки и отправляют их дважды |
FOR UPDATE без SKIP LOCKED | вторая копия ждёт освобождения строк; дублей нет, но работа идёт последовательно |
FOR UPDATE SKIP LOCKED | копии берут непересекающиеся пачки и работают параллельно |
Пример
# 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.
Очистка таблицы
Отправленные записи не нужны, но удаляются не сразу: они нужны, чтобы восстановить события, если понадобится переиздание. Срок хранения берут равным сроку хранения тем.
DELETE FROM outbox WHERE status = 'sent' AND sent_at < now() - interval '7 days';Запрос ставится в суточное расписание. Без очистки таблица растёт неограниченно; частичный индекс при этом не растёт, поэтому отправка не замедляется — замедляются резервное копирование и обслуживание таблицы.
Шаг 6. Идемпотентность потребителя
Общее правило и способы его выполнения описаны в BMBP, раздел «Идемпотентность». Здесь — только то, что относится к очереди: где хранится отметка об обработанном сообщении и когда её чистят.
Отметка нужна там, где по состоянию сущности повтор не отличить: обработка не меняет состояние ( отправка уведомления), или меняет его накопительно (начисление, списание, счётчик). Там, где повтор виден по состоянию — заказ уже не в статусе «ждёт оплаты», файл уже создан — отдельная отметка не нужна, и проверка делается по состоянию.
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 | переживает переигрывание из темы для неразобранного: идентификатор остаётся прежним |
| тройка «тема, раздел, смещение» | доступна всегда, менять контракт не требуется; у переигранного сообщения будет другое смещение, и повтор не распознается |
Отметка ставится в той же транзакции, что и результат обработки:
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)Порядок «вставка отметки, затем работа» внутри одной транзакции: при откате откатывается и отметка, поэтому «отметка есть, работа не выполнена» невозможно. Отметка в отдельной транзакции даёт именно это расхождение.
Когда чистят. Отметка нужна ровно столько, сколько сообщение может быть доставлено повторно, то есть пока оно лежит в теме. Срок хранения отметок берут вдвое больше срока хранения темы — запас покрывает переигрывание и остановленного потребителя:
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. Проверка контура целиком
Проверка идёт снизу вверх; первый шаг, который не даёт ожидаемого результата, и есть место отказа.
# 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 $?| Шаг | Ожидается |
|---|---|
| 1 | Cluster ID: … — значение из .env |
| 2 | PartitionCount: 3, ReplicationFactor: 1, в Configs есть retention.ms |
| 3 | только sent; pending — единицы, failed отсутствует |
| 4 | orders.created:0:… по каждому разделу, сумма растёт после публикации |
| 5 | LAG равен 0 или уменьшается; CONSUMER-ID заполнен по всем разделам |
| 6 | orders.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 replicas | acks=all при числе синхронных копий меньше min.insync.replicas | проверить, что все узлы кластера работают; на одном узле min.insync.replicas должен быть равен 1 |
| копии потребителя постоянно переподключаются, сообщения обрабатываются повторно | обработка одного сообщения дольше max.poll.interval.ms | сократить работу в обработчике либо поднять max.poll.interval.ms; проверить, не ждёт ли обработчик внешний сервис без таймаута |
| действие выполнено дважды | обработчик не идемпотентен | проверка по message_id или по состоянию сущности (шаг 6) |
| сообщения обрабатываются не по порядку | разные ключи или пустой ключ | ключом брать идентификатор сущности |
| потребитель не видит сообщений | подписан на другую тему или другую группу | сверить имя темы и группы; kafka-topics.sh --list покажет опечатку в имени |
| потребитель падает на новом поле | строгий разбор сообщения | игнорировать неизвестные поля |
Сдвинуть смещение застрявшей группы (группа должна быть остановлена — при живых участниках команда откажет):
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 — отказы обработчиков и публикатора.
Откат и снятие
Остановить потребителя, не потеряв смещение. Смещения хранятся в брокере, а не в процессе, и остановка их не трогает:
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) |
| подписка потребителей | конфигурация сервисов | потребители продолжат опрашивать несуществующую тему и писать в журнал ошибки |
Поэтому снятие темы выполняется в порядке: остановить потребителей → убрать тему из конфигурации сервисов → удалить группу → удалить тему.
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 продолжит работу с
накопленными данными. Удаление данных — отдельное действие:
docker volume rm <проект>_kafka_dataОно необратимо и уничтожает все темы, сообщения и смещения. Восстановить их можно только повторной публикацией из таблиц исходящих — и только те события, которые в них ещё не очищены.
Откройте исходник документа по ссылке «Предложить правку» — там же видно, что и когда в нём менялось.