Когда Kafka нужна интеграции, а когда достаточно очереди или HTTP
Интеграция начинается с факта: заказ создан, заявка согласована, статус доставки изменился, платёж отклонён. Источник этого факта должен сообщить о нём другим системам. Самый прямой вариант — HTTP-вызов из системы-источника в систему-получатель. Он удобен, когда адресат один, ответ нужен сразу и операция действительно синхронная.
Сложность возникает, когда адресатов несколько или они работают в разном темпе. Например, после изменения заказа CRM должна обновить карточку, склад — резерв, витрина — статус, а аналитика — витрину данных. Если источник вызывает все системы по цепочке, он зависит от их доступности, времени ответа и порядка вызовов. Повторы становятся его ответственностью, а подключение нового получателя требует менять источник.
Kafka задаёт другую границу взаимодействия: источник публикует событие в именованный поток, а получатели читают его независимо. Это не замена любой очереди и не универсальная шина для всех вызовов. Kafka полезна для событийных потоков, в которых важны независимые получатели, сохранение истории и возможность повторно обработать данные.
| Ситуация | Обычно подходит | Почему |
|---|---|---|
| Получить немедленный ответ от одного сервиса: проверить лимит, рассчитать цену, создать ресурс | HTTP/RPC | Результат нужен в рамках текущего запроса; событие не заменяет синхронный контракт |
| Один фоновый исполнитель должен выполнить одну короткую задачу | Очередь задач | Модель «взял работу — выполнил» может быть проще и дешевле в эксплуатации |
| Один бизнес-факт нужно независимо доставлять нескольким системам | Kafka | Каждая группа получателей хранит собственную позицию чтения; подключение нового получателя не требует менять отправителя |
| Нужно повторно обработать поток после исправления обработчика или собрать новое представление данных | Kafka | Записи хранятся по правилам хранения, а получатель может начать чтение с более ранней позиции |
| Нужна параллельная обработка множества независимых сущностей с порядком по каждой из них | Kafka | Разделы дают параллелизм, а стабильный ключ помогает сохранять порядок для одной сущности |
Kafka часто соседствует с HTTP и очередями задач. Создание заказа может синхронно вернуть ответ клиенту по HTTP, а затем сервис публикует order.created в Kafka. Склад, уведомления и аналитика реагируют независимо. Для этого нужны не только контейнеры: определите владельца события, имя темы, ключ, схему, срок хранения, правила повторов и порядок действий при ошибке. Вопросы повторной доставки и идемпотентности раскрыты в материале НТА об API-интеграциях, сбоях и повторах.
Лицензирование, исходный код и границы использования
Apache Kafka — проект с открытым исходным кодом. Официальный исходный код находится в публичном репозитории Apache Kafka; файл LICENSE в нём содержит лицензию Apache License 2.0.
Эта лицензия разрешает использовать Kafka внутри компании и в коммерческих системах, копировать, изменять и распространять исходный или собранный код. За использование самой Kafka не требуется лицензионный платёж. При распространении Kafka или производной сборки действуют условия полного текста Apache License 2.0: передать получателю текст лицензии, сохранить применимые уведомления об авторстве и патентах, явно отметить изменённые файлы, а при наличии файла NOTICE — сохранить относящиеся к продукту уведомления.
Лицензия также содержит патентное разрешение от участников проекта, но прекращает его для стороны, которая начинает патентный спор, утверждая нарушение со стороны Kafka или вклада в неё. Она не даёт право использовать товарные знаки Apache или Kafka для названия собственного продукта и не обещает гарантий качества, поддержки или возмещения ущерба. Подробные условия — в официальной лицензии Apache.
Для обычного развёртывания Kafka внутри организации достаточно соблюдать условия используемых дистрибутивов и образов. Если организация распространяет изменённую сборку Kafka, включает её в свой продукт, использует товарные знаки или формирует юридические уведомления для клиентов, итоговую схему проверяют с юристом. Kafbat UI, kcat и базовые компоненты контейнерных образов — отдельные проекты: их лицензии и уведомления проверяют самостоятельно, а не распространяют на них лицензию Kafka автоматически.
Словарь Kafka: термины, которые понадобятся дальше
В коде, настройках и официальной документации Kafka используются английские имена сущностей. Ниже они введены один раз; далее в тексте используются русские обозначения.
| Русское обозначение | Официальное имя | Смысл |
|---|---|---|
| Брокер | broker |
Процесс Kafka, который принимает, хранит и отдаёт сообщения клиентам |
| Тема | topic |
Именованный поток сообщений, например orders.created |
| Раздел | partition |
Независимая упорядоченная часть темы |
| Отправитель | producer |
Приложение, которое записывает сообщение в тему |
| Получатель | consumer |
Приложение, которое читает сообщения из темы |
| Группа получателей | consumer group |
Несколько экземпляров одного логического обработчика, распределяющих между собой разделы |
| Смещение | offset |
Позиция сообщения внутри одного раздела |
| Ключ сообщения | key |
Значение, по которому отправитель выбирает раздел для связанной последовательности событий |
| Правило хранения | retention |
Срок или объём, до которого Kafka сохраняет записи |
| Точка подключения | listener |
Адрес и протокол, по которым клиенты или компоненты Kafka подключаются к брокеру |
| KRaft | KRaft |
Внутренний механизм Kafka для хранения сведений о кластере; он заменил отдельный ZooKeeper |
Как устроена Kafka: журнал событий, брокер и KRaft
Kafka хранит не «сообщение, которое исчезло после получения», а последовательность записей в журнале событий. Запись обычно называют событием или сообщением. У неё есть значение, например JSON, необязательный ключ, время записи и технические метаданные. Отправитель записывает сообщение в Kafka, а получатель читает его. Базовую модель описывает введение Apache Kafka.
Брокер — запущенный процесс Kafka, принимающий записи и отдающий их клиентам. В реальном кластере брокеров обычно несколько: данные распределяются между ними, а разделы могут иметь копии. В этом руководстве брокер один. Этого достаточно, чтобы разобраться в протоколе, контрактах и диагностике, но недостаточно для устойчивости к отказу машины.
Тема — логическое имя потока сообщений. Примеры: orders.created, orders.status-changed, customers.updated. Тема не равна таблице базы данных и не должна становиться свалкой разнородных сообщений. Имя выражает предметный факт, а контракт определяет, кто его публикует, что означает запись и кто вправе её читать.
KRaft не является отдельным сервером, программой или «дополнением к Kafka». Это внутренний механизм Kafka для хранения и управления сведениями о кластере: какие есть брокеры, темы и разделы. В старых установках Kafka эту работу выполнял отдельный ZooKeeper. В KRaft-режиме Kafka выполняет её сама: узлы-контроллеры образуют кворум метаданных. Поэтому в этом руководстве не нужно ставить ZooKeeper и настраивать ещё один контейнер. Переход и различия режимов описаны в руководстве Kafka по KRaft.
В контуре ниже один контейнер исполняет две роли: брокер хранит и отдаёт сообщения, а контроллер ведёт служебные сведения о кластере. Это одноконтейнерная схема: она уменьшает число компонентов и подходит для разработки и тестирования, но связывает хранение данных и управление кластером с одной точкой отказа. В рабочей среде роли обычно проектируют отдельно на нескольких узлах.
Отправитель ── записывает событие ──► брокер Kafka ── хранит записи ──► Получатель
│
├── тема: orders.created
├── разделы: P0, P1, P2
└── метаданные: контроллер KRaft
Брокер не знает бизнес-смысл события: он хранит байты записи и обслуживает протокол. Смысл определяют название темы, схема, ключ и правила повторов; это обязанность интеграционного контракта и приложений.
Темы и разделы: порядок, ключ, смещение и хранение
Тема разбивается на разделы. Каждый раздел — упорядоченный журнал, в который новые записи добавляются в конец. Внутри раздела Kafka присваивает каждой записи смещение — позицию в журнале. Смещение 152 означает не «152-е сообщение всего кластера», а «позиция 152 в конкретном разделе конкретной темы».
Тема orders.created
P0: смещение 0 ── смещение 1 ── смещение 2 ── ...
P1: смещение 0 ── смещение 1 ── смещение 2 ── ...
P2: смещение 0 ── смещение 1 ── смещение 2 ── ...
Главное следствие: Kafka сохраняет порядок внутри одного раздела, но не между P0, P1 и P2. Если два связанных события попали в разные разделы, общего порядка их обработки обещать нельзя. Это не недостаток реализации: разделы позволяют распараллелить поток, поэтому единого порядка на всю тему нет.
Зачем нужен ключ сообщения
Отправитель может передать сообщение с ключом. При стандартной стратегии выбора раздела одинаковый стабильный ключ направляет записи в один и тот же раздел. Поэтому для событий одного заказа обычно выбирают order_id как ключ, для клиента — customer_id. Тогда order.created, order.paid и order.cancelled одного заказа сохраняют относительный порядок в одном журнале.
Ключ не нужно путать с идентификатором события. event_id обычно нужен для поиска повторов и трассировки: каждая публикация имеет свой уникальный идентификатор. order_id — ключ маршрутизации и бизнес-связи. Одна запись может содержать оба поля.
ключ = order-42 ─► P1: created → paid → shipped
ключ = order-77 ─► P0: created → cancelled
Собственная стратегия распределения может изменить выбор раздела, поэтому правило ключа фиксируют в контракте отправителя, а не оставляют неявным. Документация по устройству Kafka описывает выбор раздела и порядок сообщений.
Почему запись не исчезает после чтения
Kafka хранит записи по правилу хранения: например, ограниченное время или до заданного объёма данных. Получатель читает журнал с позиции и не удаляет запись своим чтением. Благодаря этому новую систему можно подключить позже и дать ей прочитать историю; существующего получателя можно запустить с более раннего смещения после исправления ошибки.
Правило хранения не равно бессрочному архиву. Оно влияет на стоимость диска, доступность истории и требования к данным. Срок, режим удаления и допустимость персональных данных должен определить владелец потока, а не случайное значение по умолчанию.
Отправитель, получатель и группа получателей: как делится работа
Отправитель публикует запись в тему. После этого он не вызывает конкретного получателя и не ждёт, пока тот обработает событие. Получатель запрашивает записи у Kafka и движется по смещениям. Позицию, до которой обработка завершена, получатель обычно фиксирует как подтверждённое смещение.
Группы получателей дают две модели одновременно:
- Независимая подписка. CRM, склад и аналитика используют разные имена групп. Каждая группа проходит тему со своей позиции и получает весь нужный ей поток.
- Параллельная обработка одного назначения. Несколько экземпляров одного сервиса используют общее имя группы. Kafka распределяет разделы между ними: один раздел внутри одной группы в момент времени читает только один получатель.
Тема orders.created: P0 P1 P2
Группа синхронизации склада: C1←P0, C1←P1, C2←P2
Группа загрузки аналитики: A1←P0, A1←P1, A1←P2
Группы синхронизации склада и загрузки аналитики читают одну тему независимо. Но если в первой группе запустить четыре получателя при трёх разделах, один участник останется без назначенного раздела: параллелизм группы ограничен количеством разделов. Это решение принимают до роста нагрузки, а не после появления отставания обработки.
Подтверждение обработки, повторы и границы гарантий
Если получатель обработал запись, а затем зафиксировал смещение, при следующем запуске он продолжит после неё. Если приложение выполнило действие, но упало до фиксации смещения, запись может быть прочитана повторно. Если смещение зафиксировали до необратимого действия и процесс упал после фиксации, действие может не произойти. Поэтому надёжная интеграция не строится на предположении, что Kafka сама исключит дубли.
Обычный практический подход — сделать обработчик идемпотентным: сохранять event_id, использовать уникальные ограничения или иной бизнес-механизм, чтобы повтор не создавал второй заказ, платёж или уведомление. Гарантия однократной обработки зависит от полного контура отправителя, Kafka, получателя и внешнего хранилища; этот учебный пример её не реализует.
Как выбрать число разделов и масштабировать Kafka
Теперь, когда понятна работа группы получателей, можно выбрать число разделов. Это не «настройка производительности вообще», а граница параллелизма и порядка для конкретной темы. Каждый раздел является отдельным журналом: его ведущая копия находится на одном брокере, а в одной группе получателей этот раздел в момент времени обрабатывает только один участник. Поэтому у темы с тремя разделами одна группа может полезно задействовать не более трёх одновременно работающих получателей. Четвёртый экземпляр будет подключён к группе, но раздел ему не назначат.
Значение KAFKA_NUM_PARTITIONS=3 в Compose-файле ниже — только значение по умолчанию для новых тем. Оно выбрано, чтобы показать распределение работы в учебном стенде. Это не рекомендация «три раздела для любой темы» и не влияет на темы, созданные с явным параметром --partitions.
От чего зависит выбор
Сначала ответьте на четыре вопроса, а уже затем назначайте число разделов.
| Вопрос | Как влияет на выбор | Пример решения |
|---|---|---|
| Какой порядок действительно нужен? | Если нужен единый строгий порядок всех событий темы, подходит только один раздел. Если нужен порядок по заказу, клиенту или договору, используют стабильный ключ и несколько разделов | Для статусов одного заказа ключ — order_id; разные заказы можно обрабатывать параллельно |
| Сколько экземпляров одного обработчика должны работать одновременно? | Полезный параллелизм одной группы не больше числа разделов | Для двух одновременно работающих экземпляров нужен как минимум два раздела; с запасом на рост можно начать с трёх после проверки порядка по ключу |
| Каковы объём, размер события и скорость обработки? | Нужна не догадка, а нагрузочная проверка с реальными или синтетическими данными: размером сообщений, ключами, подтверждениями и внешней операцией получателя | Если один раздел устойчиво обрабатывает 200 событий в секунду, а требуется 800, исходная гипотеза — не менее четырёх разделов; её подтверждают тестом всего контура |
| Какова стоимость хранения и эксплуатации? | Больше разделов означает больше журналов, файлов, метаданных, копий, фоновых операций и работы при восстановлении | Десять разделов «на всякий случай» для редкой темы усложнят сопровождение, но не дадут выгоды |
Практическое правило для первого расчёта: число разделов должно быть не меньше максимума из двух величин — требуемого количества одновременно работающих экземпляров одной группы и отношения целевой скорости потока к измеренной устойчивой скорости одного раздела. Результат округляют в большую сторону и проверяют нагрузочным испытанием. Это начальная оценка, а не универсальная формула: реальная граница может находиться в получателе, базе данных, сети или диске брокера.
| Выбор | Когда оправдан | Что он ограничивает |
|---|---|---|
| 1 раздел | Низкий поток и единый порядок всех записей действительно важен | В группе работает только один активный получатель; горизонтально ускорить обработку этой темы нельзя |
| 2–3 раздела | Пилотный поток с независимыми ключами и небольшой потребностью в параллельной обработке | Это удобный старт для Dev/Test, но не расчёт мощности рабочей системы |
| 6, 10 и больше | Нагрузочный тест подтвердил необходимость такого параллелизма, ключи распределяются достаточно равномерно, а команда готова сопровождать больше журналов и копий | Не компенсирует «горячий» ключ: все его события всё равно останутся в одном разделе |
Не выбирайте ключ только ради равномерности. Если события одного заказа должны идти последовательно, нельзя разбивать их между разделами случайным ключом. И наоборот: если один клиент генерирует почти весь поток, ключ customer_id создаст перегруженный раздел независимо от общего числа разделов. В таком случае пересматривают бизнес-границу порядка или способ разбиения данных; простое добавление разделов проблему не устраняет.
Что можно увеличить, а что нельзя исправить задним числом
Число разделов темы можно увеличить, например с трёх до шести. Это добавляет новые пустые разделы; уже записанные сообщения Kafka автоматически не переносит. При стандартном хешировании ключа число разделов входит в расчёт маршрута, поэтому новые сообщения с прежним ключом после расширения могут попасть в другой раздел. Значит, относительный порядок старых и новых событий одного ключа может оказаться в разных журналах. Перед таким изменением нужно проверить контракт порядка, состояние получателей и план перехода. Эта оговорка прямо указана в руководстве по операциям Apache Kafka.
Исполняемый пример расширения учебной темы приведён после первого запуска Kafka. Уменьшить число разделов обратно нельзя: для этого создают новую тему, переносят поток контролируемым способом и переключают клиентов.
За счёт чего масштабируется Kafka
Kafka масштабируется по нескольким независимым направлениям. Их нельзя заменить одним действием «добавить сервер».
| Цель | Что увеличивают | Что обязательно проверить |
|---|---|---|
| Быстрее обработать одну тему одним логическим сервисом | Число экземпляров получателя, но не выше числа разделов | Время обработки, отставание, идемпотентность и отсутствие «горячих» ключей |
| Принять и отдать больше потока | Число разделов, параметры клиентов, сетевую и дисковую ёмкость; затем измеряют весь маршрут | Распределение ключей, размер сообщений, задержку, загрузку CPU/диска/сети и отставание групп |
| Увеличить хранение, пропускную способность и пережить потерю узла | Число брокеров и число копий каждого раздела | Фактор копирования не выше числа брокеров, размещение копий по зонам отказа, минимальное число синхронных копий и восстановление |
При добавлении нового брокера существующие разделы не переезжают на него автоматически. Администратор подготавливает план распределения, запускает перераспределение разделов и наблюдает его до завершения; в процессе Kafka копирует данные на новый узел. Официальная документация описывает команды создания, запуска и проверки такого плана, а также ограничение скорости переноса, чтобы миграция не вытеснила рабочую нагрузку. Для Dev/Test из этой статьи это не нужно: здесь только один брокер, а репликация равна единице.
Как получатель выбирает позицию и подтверждает обработку
Позиция чтения и подтверждение обработки — разные решения. Сначала получатель должен понять, с какого сообщения начать, затем — в какой момент считать обработку завершённой. Если смешать эти вопросы, после перезапуска появятся либо непредсказуемые пропуски, либо непонятные повторы.
Откуда начинается чтение
Когда получатель работает в группе, Kafka хранит для неё подтверждённое смещение каждого раздела. При следующем запуске участник группы обычно продолжает после этой позиции. Имя группы задаётся свойством group.id; оно обязательно, если приложение использует управление группой или хранит позиции в Kafka. Это описано в настройках получателя Apache Kafka.
Если для группы ещё нет сохранённой позиции или нужная запись уже удалена по правилу хранения, применяется auto.offset.reset:
| Значение | Откуда начнётся чтение | Когда выбирать |
|---|---|---|
earliest |
С самого раннего доступного сообщения каждого раздела | Новый получатель должен построить состояние по доступной истории |
latest |
Только с сообщений, опубликованных после подключения | Нужны только новые события; это значение по умолчанию, поэтому его нельзя оставлять неосознанно |
none |
Приложение завершится ошибкой, если позиции нет | Пропуск истории недопустим и его нужно явно разобрать оператору |
by_duration:PT24H |
С позиции, соответствующей указанному интервалу до текущего времени | Нужен контролируемый запуск с недавнего окна; пример означает «за последние 24 часа» |
auto.offset.reset не заставляет Kafka читать с начала при каждом запуске. Он применяется только когда у группы нет пригодного сохранённого смещения. Для разового просмотра истории консольная команда из этой статьи использует --from-beginning; для прикладного получателя начальную позицию нужно закрепить в конфигурации и контракте.
Подписка группы или явное назначение разделов
| Способ чтения | Как работает | Где уместен |
|---|---|---|
| Подписка на тему в группе | Приложение указывает тему и имя группы; Kafka сама распределяет разделы между активными участниками | Обычный сервис обработки событий, который можно масштабировать несколькими экземплярами |
| Явное назначение разделов приложению | Приложение само указывает, какие разделы читать, и само управляет позициями | Диагностика, контролируемое переигрывание и специализированные служебные процессы |
| Поиск по заданному смещению | Приложение временно устанавливает позицию в конкретном разделе и читает оттуда | Разбор инцидента или повторная обработка ограниченного диапазона; не замена нормальной модели группы |
Официальный Java API Kafka содержит отдельные операции для подписки, явного назначения, перехода к заданному смещению и чтения с начала или конца раздела. Другие библиотеки называют методы иначе, но предоставляют ту же модель. Справочник Consumer API полезен как первоисточник этой механики.
Когда фиксировать смещение
В API клиентских библиотек фиксация смещения обычно называется commit. Она записывает в Kafka позицию, после которой группа продолжит чтение. Значение имеет не просто способ фиксации, а её место относительно бизнес-действия.
| Режим | Что происходит | Риск и применение |
|---|---|---|
Автоматическая фиксация: enable.auto.commit=true |
Клиент периодически фиксирует позицию в фоне; стандартный интервал — 5 секунд | Фиксация не привязана к завершению внешней операции. Подходит для простого наблюдения, но требует особой осторожности для платежей, заказов и других бизнес-действий |
| Ручная синхронная фиксация | Приложение обработало запись или пачку, затем ждёт подтверждения фиксации от Kafka | Самый понятный старт для бизнес-получателя: сначала обработка, потом фиксация. При сбое возможен повтор, поэтому нужна идемпотентность |
| Ручная асинхронная фиксация | Приложение запускает фиксацию без ожидания и получает результат через обработчик ошибки | Уменьшает задержку, но требует явной обработки ошибок и упорядочивания попыток; не стоит выбирать её только «для скорости» |
Транзакционный режим чтения: isolation.level=read_committed |
Получатель видит только подтверждённые транзакционные записи Kafka | Полезен, когда отправитель действительно использует транзакции Kafka. Не делает внешнюю базу данных или HTTP-вызов автоматически транзакционными |
В рабочей интеграции безопасная базовая последовательность обычно выглядит так: получить запись → выполнить идемпотентное бизнес-действие → сохранить его результат → зафиксировать смещение. Это даёт модель «как минимум один раз»: при отказе возможен повтор, но корректно написанный обработчик не создаст второй эффект. Фиксация перед бизнес-действием даёт другую модель — сообщение может быть пропущено. Официальная документация подтверждает, что при enable.auto.commit=true смещения фиксируются периодически в фоне; синхронные и асинхронные методы фиксации есть в Consumer API.
Что подтверждает контур разработки и тестирования, а что нет
Один Kafka-узел, который одновременно работает брокером и контроллером, отвечает на полезные вопросы: запускается ли брокер, правильны ли адреса клиентов, создаётся ли тема, проходит ли запись и может ли получатель прочитать событие. Он не отвечает на вопросы о потере машины, репликации, настройке прав или восстановлении после сбоя.
| Можно проверить | Нельзя подтверждать этим стендом |
|---|---|
| Отправитель подключается и публикует запись | Отказоустойчивость при потере брокера или диска |
| Тема и разделы создаются с ожидаемыми параметрами | Репликацию, кворум и обновление без остановки |
| Получатель читает запись с указанной позиции | Защиту данных: в примере нет TLS, SASL и правил доступа ACL |
| Схема события и ключ подходят для пилотного потока | Производительность, допустимое отставание и запас мощности рабочей системы |
| Kafbat UI и командные утилиты видят брокер из нужной сети | Резервное копирование и восстановление |
Ресурсы и первичная оценка производительности
Kafka не предъявляет одной фиксированной нормы «минимум N ядер и M ГБ»: скорость определяется размером записи, числом разделов и копий, подтверждением записи, сроком хранения, скоростью сети и тем, что делает получатель после чтения. Брокер немедленно записывает данные в файловую систему, а свободная память Linux используется как страничный кэш. Поэтому для Kafka важны не только процессор и объём памяти, но прежде всего предсказуемый SSD/NVMe-диск и запас свободного места. Руководство Apache Kafka по оборудованию и ОС отдельно отмечает зависимость от пропускной способности диска и предлагает оценивать память как минимум по объёму записи за 30 секунд.
Ниже — инженерные стартовые профили для Docker Compose из этой статьи. Это не паспорт производительности Kafka и не замена нагрузочному тесту. Значения относятся к ресурсам, доступным Docker и Kafka, а не к общему размеру виртуальной машины, на которой уже работают база данных, сборщик логов или другие сервисы.
| Сценарий | Процессор | Оперативная память | Диск для данных Kafka | Сеть | Что можно делать |
|---|---|---|---|---|---|
| Личный локальный Dev | 2 виртуальных ядра | 4 ГБ доступной памяти | 25 ГБ SSD | Локальная сеть Docker | Запустить один брокер и UI, проверить тему, запись, чтение, ключи и группы. Не считать этот профиль измерением пропускной способности или отказоустойчивости |
| Общий Dev/Test с несколькими приложениями | 4 виртуальных ядра | 8 ГБ доступной памяти | 100 ГБ SSD или NVMe | 1 Гбит/с внутри сегмента | Одновременно проверять несколько интеграционных потоков, группы и диагностические инструменты; проводить ограниченные нагрузочные испытания без бизнес-данных — только в отдельной конфигурации с TLS, SASL и ACL |
| Dev/Test с проверкой отказа узла | 3 виртуальные машины, на каждой от 4 виртуальных ядер | от 8 ГБ на машину | от 100 ГБ SSD/NVMe на машину | 1 Гбит/с между машинами | Проверять копирование данных, смену ведущей копии, восстановление и поведение при недоступности узла. Это уже отдельная многоузловая конфигурация, не Compose-файл ниже |
Для первых двух профилей один брокер остаётся единственной копией данных. Добавление памяти и диска не превращает его в отказоустойчивый кластер. Если Dev/Test должен переживать потерю узла или использовать фактор копирования больше единицы, нужны несколько брокеров; число копий темы не может быть больше числа доступных брокеров.
Профиль ресурсов и защита доступа — два независимых условия. Текущий Compose-файл намеренно подходит только для личного локального Dev: его порты привязаны к 127.0.0.1 и в нём нет TLS, SASL и ACL. Для общего Dev/Test ресурсы из второй строки таблицы применимы только после создания отдельного защищённого контура по разделу «Безопасность Kafka: от локального стенда к общему контуру».
Как оценить место и память до запуска
Сначала оцените объём, который тема будет записывать в сутки:
суточный объём = средний размер записи × число записей в секунду × 86 400
Для общего объёма диска кластера умножьте результат на срок хранения в сутках, фактор копирования и запас не менее 30 % на служебные данные, неравномерность разделов и восстановление:
диск кластера ≈ суточный объём × срок хранения × фактор копирования × 1,3
При равномерном размещении ориентир для одного брокера — общий объём, делённый на число брокеров. Это нижняя оценка: она не учитывает временный дополнительный объём при перераспределении разделов, неравномерные ключи и рост потока. Свободное место отслеживают постоянно, а не только перед первым запуском.
Для памяти полезна грубая первичная оценка из документации Apache Kafka: объём входящего потока за 30 секунд. Например, при 5 МБ/с это около 150 МБ кэша; к нему добавляют память самого процесса Kafka, операционной системы, Docker и других сервисов. Поэтому 4 ГБ подходят только для маленького локального стенда, а 8 ГБ — разумный старт для общего Dev/Test, но не гарантия для потока с долгим чтением истории, большим числом разделов или тяжёлым UI.
Нагрузочный тест выполняйте только после запуска Kafka, создания и чтения первого события. Полный воспроизводимый маршрут приведён далее, после базовой сквозной проверки.
Контур разработки и тестирования: требования и состав
Способ развёртывания в этом руководстве — Docker Compose. Он выбран как способ зафиксировать проверяемую конфигурацию Kafka в одном файле. Установка Docker здесь не рассматривается.
Перед началом нужны:
- Debian или Ubuntu, виртуальная машина либо локальная Linux-среда, не являющаяся рабочим сервером;
- установленный и запущенный Docker Engine;
- дополнение Docker Compose, доступное командой
docker compose; - права на выполнение Docker-команд в вашей среде;
- доступ к реестру образов и свободные локальные порты
29092и8080; - отдельный каталог для стенда и место для именованного Docker-тома.
Порты в примере опубликованы только на 127.0.0.1. Это значит, что Kafka и интерфейс доступны с самого хоста, но не из сети. Не заменяйте адрес на 0.0.0.0: пример передаёт данные без шифрования и аутентификации.
Создайте рабочий каталог и перейдите в него:
mkdir -p ~/kafka-dev
cd ~/kafka-dev
Безопасность Kafka: от локального стенда к общему контуру
Compose-файл из следующего раздела — локальный изолированный стенд. Он передаёт данные без шифрования и не проверяет, кто подключается к брокеру. Единственная защитная мера в нём — привязка портов к 127.0.0.1, поэтому соединиться может только программа на той же машине. Такой вариант допустим для личного обучения и синтетических данных. Его нельзя открыть в локальную сеть, VPN, интернет или сделать общим Dev/Test-контуром заменой адреса на 0.0.0.0.
Для общего Dev/Test-контура с несколькими пользователями и приложениями минимальная защищённая модель состоит из четырёх обязательных свойств.
| Свойство | Что должно быть сделано | Что это защищает |
|---|---|---|
| Шифрование канала | Для клиентских и межброкерных соединений включён TLS; сертификат брокера содержит фактические DNS-имена, а клиенты проверяют цепочку доверия и имя узла | Сообщения, пароли и служебные данные нельзя прочитать или подменить в сети |
| Проверка личности | У каждого приложения и оператора своя учётная запись. Для общего контура предпочтителен SASL/SCRAM-SHA-512 поверх TLS; механизм SASL/PLAIN допустим только внутри TLS и как осознанное упрощение | Брокер различает подключающиеся приложения и не принимает анонимные соединения |
| Права | Включён механизм авторизации и правила доступа ACL: отправитель имеет право записывать только в свои темы, получатель — читать только из своих тем и групп, оператор получает отдельные административные права | Учётная запись одного сервиса не получает доступ ко всем данным кластера |
| Граница сети и секреты | Kafka-порты доступны только нужным подсетям и приложениям; интерфейс управления не публикуется наружу. Закрытые ключи, пароли и хранилища сертификатов выдаёт корпоративное хранилище секретов или защищённое подключение файлов, а не репозиторий | Снижается риск доступа через сеть, утечки из Compose-файла и случайного использования общих учётных данных |
TLS отвечает за конфиденциальность канала и проверку сертификата узла, но сам по себе не определяет права приложения. SASL подтверждает личность клиента, а ACL ограничивают его действия. Эти три слоя дополняют друг друга; включить только один из них недостаточно для общего контура. Обзор механизмов и точек их настройки даёт раздел безопасности Apache Kafka.
Практический выбор механизма
Для нового общего Dev/Test-контура выбирайте SASL_SSL с SCRAM-SHA-512 и отдельными учётными записями приложений. У механизма SCRAM учётные данные хранятся в метаданных Kafka; их нужно создать контролируемо до запуска связи между брокерами. Это требует отдельной процедуры начальной инициализации, поэтому не копируйте пароли в Compose-переменные. Apache Kafka подробно описывает подготовку SCRAM и порядок создания учётных данных в руководстве по SASL.
SASL/PLAIN проще для короткоживущего стенда, но пароль в этом механизме передаётся клиентом, поэтому он допустим только поверх TLS. Стандартный вариант Kafka хранит пароли в файле JAAS; для постоянной среды следует получать их из внешнего хранилища через поддерживаемый обработчик проверки, а не оставлять в файле или истории команд. Это не терминологическая тонкость: соединение SASL_PLAINTEXT передаст пароль без шифрования.
Официальный Docker-образ ожидает секретные файлы в /etc/kafka/secrets: для SASL туда монтируют JAAS-конфигурацию и указывают её через KAFKA_OPTS; для TLS монтируют хранилище ключа и хранилище доверенных сертификатов, а имена файлов и пароли передают предусмотренными переменными. Список этих параметров и официальные примеры опубликованы в руководстве Docker-образа Apache Kafka. Не помещайте сами секреты в статью, docker-compose.yml, Git или журнал команд.
Минимальный контрольный список перед общим доступом
- Выпущены сертификаты для всех фактических DNS-имён брокеров; клиент с недоверенным сертификатом не подключается.
- У каждого приложения отдельная учётная запись и отдельный секрет; проверена смена пароля или сертификата без публикации секрета в Git.
- Для каждой темы и группы выданы минимально необходимые ACL; отдельно проверены разрешённая и запрещённая операции от имени прикладной учётной записи.
- Внешний порт Kafka и Kafbat UI закрыты правилами сети; интерфейс управления доступен только через согласованный защищённый маршрут и имеет собственную аутентификацию.
- Журналы, метрики неуспешных подключений и изменения прав собираются; секреты в них маскируются.
- Стенд проверен клиентом с теми же настройками TLS, SASL и ACL, которые будут у приложения. Успешный запуск незашифрованного локального Compose не подтверждает эту проверку.
Настройка сертификатов, системы идентификации и ACL зависит от корпоративной инфраструктуры и состава приложений. Поэтому здесь намеренно нет универсального «безопасного Compose на пять строк»: такой файл не может корректно выпустить сертификаты, защитить секреты и выдать права на конкретные темы. Для общего контура сначала согласуйте эту модель с ИБ и владельцами интеграций, затем соберите отдельную конфигурацию и прогоните её тем же маршрутом: подключение, создание разрешённой темы, запись, чтение, попытка запрещённого действия и просмотр журналов.
Управляемая Kafka в российских облаках
Самостоятельный Compose-контур нужен, чтобы понять Kafka, проверить контракт и управлять конфигурацией полностью. Если требуется готовый кластер в облаке, часть операционной работы — развёртывание брокеров, обновления, мониторинг и обслуживание — может взять на себя провайдер. Это не отменяет проектирование тем, ключей, прав доступа, срока хранения и обработчиков: ответственность за данные и интеграционный контракт остаётся у команды.
Ниже — не рейтинг и не рекомендация поставщика, а небольшой перечень сервисов, у которых на дату подготовки материала есть публичная документация управляемой Kafka.
| Провайдер | Сервис | Что подтверждает документация провайдера |
|---|---|---|
| Yandex Cloud | Managed Service for Apache Kafka | Развёртывание и поддержка кластеров Apache Kafka в инфраструктуре Yandex Cloud |
| Cloud.ru | Evolution Managed Kafka | Управление кластерами Apache Kafka; документация содержит создание кластеров, доступ, темы, параметры, журналирование и мониторинг |
| Timeweb Cloud | Managed Service for Kafka | Создание и управление Kafka-кластерами в облаке; заявлены интерфейс и программный интерфейс управления |
Перед выбором сравните не рекламные формулировки, а технические условия своей задачи: поддерживаемую версию Kafka, режим KRaft, число зон и копий данных, правила сетевого доступа, TLS/SASL/ACL, работу с секретами, лимиты тем и разделов, окно обслуживания, экспорт и восстановление данных, метрики, журналирование, стоимость ресурсов и уровень поддержки. Условия, доступные версии и регионы меняются; сверяйте их с актуальной документацией и договором конкретного провайдера перед запуском рабочей нагрузки.
Docker Compose: готовая конфигурация
Создайте в этом каталоге файл docker-compose.yml. Это именно локальная изолированная конфигурация без TLS, SASL и ACL, а не вариант общего защищённого контура. В нём зафиксированы официальный образ Apache Kafka 4.3.1 и Kafbat UI v1.5.0 для ручной диагностики. Конкретные теги выбраны для воспроизводимости и проверены на дату подготовки материала; перед обновлением тега повторите проверку из следующего раздела.
services:
kafka:
image: apache/kafka:4.3.1
hostname: kafka
restart: unless-stopped
ports:
- "127.0.0.1:29092:29092"
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093,PLAINTEXT_HOST://:29092
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
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_NUM_PARTITIONS: 3
KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false"
volumes:
- kafka-data:/tmp/kraft-combined-logs
healthcheck:
test: ["CMD-SHELL", "/opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list >/dev/null 2>&1"]
interval: 10s
timeout: 5s
retries: 12
start_period: 20s
kafka-ui:
image: ghcr.io/kafbat/kafka-ui:v1.5.0
restart: unless-stopped
depends_on:
kafka:
condition: service_healthy
ports:
- "127.0.0.1:8080:8080"
environment:
DYNAMIC_CONFIG_ENABLED: "true"
volumes:
kafka-data:
kafka-data — именованный Docker-том. Он сохраняет журналы Kafka при обычных docker compose stop и docker compose up; данные будут удалены, если выполнить docker compose down -v. Внутри официального одноброкерного образа каталог журналов KRaft — /tmp/kraft-combined-logs, поэтому именно он подключён к тому. Не меняйте путь без сверки документации и запуска на своей версии образа.
Пояснение Kafka-параметров в Compose
Переменные вида KAFKA_* передаются официальному образу как настройки брокера. Их полный перечень и допустимые значения зависят от версии Kafka; первичный источник — каталог настроек брокера Apache Kafka и официальное руководство по Docker-образу. Ниже разобраны все Kafka-переменные из этого файла.
| Переменная | Что задаёт | Почему здесь такое значение | Что меняется вне контура разработки и тестирования |
|---|---|---|---|
KAFKA_NODE_ID=1 |
Уникальный числовой идентификатор узла Kafka | В стенде один узел, поэтому достаточно 1 |
У каждого брокера и контроллера свой уникальный идентификатор; его планируют вместе с топологией |
KAFKA_PROCESS_ROLES=broker,controller |
Роли процесса Kafka | Один контейнер одновременно принимает клиентские записи и ведёт служебные сведения о кластере | В рабочей среде роли брокера и контроллера обычно разделяют и используют несколько узлов |
KAFKA_LISTENERS=... |
Адреса, на которых процесс принимает соединения | 9092 — сеть контейнеров, 29092 — клиент с хоста, 9093 — канал контроллера |
Имена и адреса точек подключения зависят от сети, TLS/SASL и топологии |
KAFKA_ADVERTISED_LISTENERS=... |
Адреса, которые брокер сообщает клиенту после первого подключения | Контейнер получает kafka:9092, программа на хосте — localhost:29092 |
Указывают реальные DNS-имена или балансировщики, доступные соответствующей группе клиентов |
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=... |
Протокол для каждого имени точки подключения | Все три канала работают без шифрования, что допустимо только в изолированном стенде | В рабочей среде используют согласованные SSL или SASL_SSL и управляют сертификатами и секретами |
KAFKA_INTER_BROKER_LISTENER_NAME=PLAINTEXT |
Точку подключения для обмена между брокерами | В конфигурации один брокер, но свойство нужно для согласованной внутренней конфигурации Kafka | Выделяют защищённую внутреннюю точку подключения, доступную всем брокерам |
KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER |
Точку подключения роли контроллера | Отделяет обмен метаданными от клиентского канала | Канал контроллеров доступен только кворуму контроллеров и защищён сетевыми правилами |
KAFKA_CONTROLLER_QUORUM_VOTERS=1@kafka:9093 |
Состав кворума контроллеров в формате идентификатор@хост:порт |
Единственный контроллер с идентификатором 1 доступен контейнерам под именем kafka |
Список содержит несколько заранее спроектированных узлов контроллеров; это не параметр для случайного масштабирования |
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_NUM_PARTITIONS=3 |
Число разделов по умолчанию для новых тем | Три раздела помогают увидеть распределение и модель групп | Не меняет уже созданные темы и не заменяет расчёт параллелизма; критерии выбора и последствия расширения разобраны в разделе «Как выбрать число разделов и масштабировать Kafka» |
KAFKA_AUTO_CREATE_TOPICS_ENABLE=false |
Автоматическое создание темы при обращении к несуществующему имени | Опечатка не превращается в «пустую работающую тему»; создание остаётся явным шагом | В управляемых интеграциях это обычно безопаснее, но решение принимают по процессу управления темами |
Почему нужны два клиентских адреса
Kafka-клиент сначала подключается к начальному адресу, затем получает метаданные с адресом брокера и подключается к нему. Поэтому опубликованный Docker-порт сам по себе не решает задачу: брокер должен сообщить адрес, доступный именно этому клиенту.
- Команда, запущенная внутри контейнера Kafka, использует
localhost:9092. - Kafbat UI и другой контейнер в этой Compose-сети используют
kafka:9092. - Клиент, запущенный на Debian/Ubuntu-хосте, использует
localhost:29092.
Если задать Kafbat UI localhost:29092, для интерфейса это будет его собственный контейнер, а не Kafka. Если хостовой программе выдать kafka:9092, Docker DNS-имя может быть ей неизвестно. Ошибка в параметре адресов, сообщаемых клиентам (KAFKA_ADVERTISED_LISTENERS), часто выглядит так: первое соединение есть, но клиент всё равно не читает тему.
Первый запуск и проверка брокера
Сначала проверьте синтаксис Compose и поднимите сервисы:
sudo docker compose config
sudo docker compose up -d
sudo docker compose ps
Для kafka дождитесь состояния healthy. Проверка готовности запускает kafka-topics.sh --list внутри контейнера: она подтверждает, что брокер уже отвечает по протоколу Kafka, а не только то, что процесс контейнера не завершился. Если готовность не появилась, не создавайте тему: сначала прочитайте журнал.
sudo docker compose logs --tail=100 kafka
Создание темы, запись и чтение первого события
Все Kafka-команды ниже выполняются в контейнере. Java и архив Kafka на хосте не требуются: штатные kafka-*.sh уже находятся в официальном образе.
1. Создайте тему
sudo docker compose exec kafka \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--create \
--if-not-exists \
--topic demo.orders.created \
--partitions 3 \
--replication-factor 1
Проверьте фактическую конфигурацию:
sudo docker compose exec kafka \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--describe \
--topic demo.orders.created
В выводе должны быть PartitionCount: 3 и ReplicationFactor: 1. Параметр --partitions 3 задан явно и поэтому не зависит от KAFKA_NUM_PARTITIONS. В рабочем потоке число разделов выбирают из требуемого параллелизма, ключа и роста нагрузки; увеличить число можно, но это не перераспределит старые записи и может изменить распределение ключей для новых.
2. Отправьте тестовое сообщение
Консольный отправитель читает одну строку из терминала как одно событие. Запустите его:
sudo docker compose exec kafka \
/opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server localhost:9092 \
--topic demo.orders.created
Вставьте синтетическое событие и нажмите Enter:
{"event_id":"test-001","event_type":"order.created","order_id":"DEV-42","occurred_at":"2026-08-17T12:00:00Z"}
Завершите отправитель сочетанием Ctrl+C. В этом первом тесте ключ не передаётся: цель — проверить запись и чтение. В интеграции, где порядок по заказу имеет значение, отправитель должен передавать стабильно выбранный ключ, например order_id; для консольного отправителя это задаётся отдельными параметрами сериализации ключа.
3. Прочитайте запись отдельным получателем
Откройте второй терминал и выполните:
sudo docker compose exec kafka \
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic demo.orders.created \
--group demo-inspector-01 \
--from-beginning \
--max-messages 1
На экране появится отправленный JSON, а команда завершится после одного сообщения. Это минимальное доказательство сквозного пути: брокер сохранил запись в теме и отдал её получателю. --group demo-inspector-01 создаёт имя учебной группы: оно понадобится далее для просмотра её позиции и отставания. --from-beginning важен для учебного теста: без него новая группа может ожидать только события, опубликованные после её запуска. При повторном прохождении инструкции замените окончание имени группы на новое, например demo-inspector-02: сохранённая позиция существующей группы не сбрасывается параметром --from-beginning.
Увеличение числа разделов учебной темы
Только после успешного создания и чтения записи можно безопасно повторить операцию расширения на учебной теме. Она добавит три пустых раздела; уже сохранённые сообщения не будут перемещены.
sudo docker compose exec kafka \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--alter \
--topic demo.orders.created \
--partitions 6
Проверьте итоговое число разделов командой kafka-topics.sh --describe. Не запускайте это на теме с бизнес-данными без решения о порядке: после расширения новые сообщения с прежним ключом могут попасть в другой раздел.
Нагрузочный тест после сквозной проверки
Не пишите в документации «сервер выдержит 10 000 сообщений в секунду», пока не проверены именно ваши параметры. Для первого сравнимого теста зафиксируйте размер записи, число разделов, длительность, способ подтверждения записи и состав получателей. Тест брокера без бизнес-обработки измеряет способность Kafka принимать записи; он не измеряет скорость базы данных, HTTP-вызовов и вашей прикладной логики.
| Профиль измерения | Нагрузка | Для чего нужен | Критерий пригодности |
|---|---|---|---|
| Минимальный локальный | 1 000 записей/с по 1 КиБ в течение 5 минут: всего 300 000 записей | Убедиться, что стенд не теряет запись и не создаёт устойчивого отставания на небольшой синтетической нагрузке | Нет ошибок отправителя и брокера; тема доступна; тестовый получатель догоняет поток за согласованное время |
| Рекомендуемый общий Dev/Test | 5 000 записей/с по 1 КиБ в течение 10 минут: всего 3 000 000 записей | Найти очевидное ограничение диска, сети, настроек Docker или группы получателей до подключения нескольких команд | Нет перезапуска брокера, свободен запас диска, нет устойчиво растущего LAG; фактическая скорость и задержка записаны в результатах теста |
Эти числа — стартовые сценарии измерения, а не обещание, что любой сервер с указанными ресурсами достигнет такой скорости. Для события размером 10 КиБ поток в 1 000 сообщений/с уже передаёт примерно 10 МБ/с; для получателя с внешним запросом на 50 мс предел может определяться вовсе не Kafka. Повышайте нагрузку ступенями, пока не увидите устойчивую границу, затем оставляйте запас мощности, согласованный с владельцем интеграции.
Создайте отдельную тему для теста, чтобы не смешивать синтетические записи с учебным потоком:
sudo docker compose exec kafka \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--create \
--if-not-exists \
--topic demo.performance \
--partitions 3 \
--replication-factor 1
Затем отправьте первый профиль штатной утилитой Kafka. --throughput 1000 ограничивает отправку примерно тысячей записей в секунду, --record-size 1024 задаёт размер одной синтетической записи в байтах, а acks=all требует подтверждения от всех синхронных копий. В одноброкерном стенде это единственная копия, но явное значение позволяет сравнивать будущие прогоны. compression.type=none исключает влияние сжатия на первый базовый замер.
sudo docker compose exec kafka \
/opt/kafka/bin/kafka-producer-perf-test.sh \
--bootstrap-server localhost:9092 \
--topic demo.performance \
--num-records 300000 \
--record-size 1024 \
--throughput 1000 \
--command-property acks=all compression.type=none
Утилита выводит фактическую скорость, среднюю и максимальную задержку записи. Сохраните этот вывод вместе с параметрами машины и Docker. Для второго профиля увеличьте --num-records до 3000000, а --throughput до 5000; сохраните те же acks и сжатие. Не запускайте второй профиль на машине с 25 ГБ диска, пока не рассчитан объём хранения и не подтверждён запас места. Формат и параметры штатной утилиты приведены в официальном примере Apache Kafka.
Чтобы проверить не только запись, но и догоняет ли получатель поток, прочитайте тестовую тему отдельной временной группой. -T отключает интерактивный режим Docker, а перенаправление в /dev/null не выводит в терминал сотни тысяч синтетических записей:
sudo docker compose exec -T kafka \
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic demo.performance \
--group perf-inspector-01 \
--from-beginning \
--max-messages 300000 > /dev/null
Затем подтвердите, что эта группа догнала тему:
sudo docker compose exec kafka \
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--describe \
--group perf-inspector-01
Для минимального профиля ожидается LAG: 0 по всем трём разделам. После этого проверьте журналы Kafka, свободное место и отсутствие ошибок отправителя. Перед повтором замените окончание имени группы на новое, ещё не использованное. Для второго профиля замените --max-messages на 3000000; на реальной интеграции вместо консольного получателя используют сам сервис, потому что только он покажет фактическую скорость бизнес-обработки.
Клиентские библиотеки для приложений
Команды внутри контейнера подходят для первого теста, но прикладной сервис подключается через библиотеку своего языка. Ниже — распространённые поддерживаемые варианты. Только Java-клиент входит в проект Apache Kafka; остальные библиотеки сопровождает Confluent и они совместимы с Apache Kafka. Сводный перечень и ссылки на примеры публикует документация клиентских библиотек.
| Язык или платформа | Пакет или библиотека | Что предоставляет | Руководство |
|---|---|---|---|
| Java и JVM-языки | org.apache.kafka:kafka-clients |
Базовые отправитель, получатель и административный клиент Apache Kafka | Apache Kafka Java API |
| .NET | Confluent.Kafka из NuGet |
Высокоуровневые отправитель, получатель и административный клиент; использует librdkafka |
документация .NET-клиента |
| Python | confluent-kafka из PyPI |
Отправитель, получатель и административный клиент; обёртка над librdkafka |
документация Python-клиента |
| Go | github.com/confluentinc/confluent-kafka-go/v2/kafka |
Отправитель и получатель на базе librdkafka |
документация Go-клиента |
| JavaScript и TypeScript | @confluentinc/kafka-javascript из npm |
Отправитель и получатель, включая вариант API на промисах | документация JavaScript-клиента |
| C и C++ | librdkafka |
Низкоуровневая библиотека, на которой основаны ряд клиентов выше | репозиторий librdkafka |
Для этого стенда все библиотеки получают одинаковые исходные параметры, хотя названия полей в коде могут отличаться:
| Где запущено приложение | Адрес брокера | Обязательные решения получателя |
|---|---|---|
| В другом контейнере той же Compose-сети | kafka:9092 |
Имя группы, начальная позиция, способ фиксации, имя темы и формат ключа/значения |
| Непосредственно на Debian/Ubuntu-хосте | localhost:29092 |
Те же решения; не использовать Docker-имя kafka, если оно не настроено в DNS хоста |
Не начинайте с незафиксированной версии пакета. Версию библиотеки закрепляют в зависимостях приложения, сверяют с её официальной матрицей поддержки и проверяют на таком же стенде после обновления. В прикладной конфигурации явно задайте как минимум адрес брокера, имя группы, начальную позицию и стратегию фиксации, а не полагайтесь на значения по умолчанию.
Kafbat UI и kcat: наблюдение и диагностика
Kafbat UI
Kafbat UI — независимый проект с открытым исходным кодом, не часть Apache Kafka. Он удобен, когда нужно просмотреть темы, разделы, сообщения, группы получателей и их позиции чтения без набора команд. Интерфейс не доказывает работоспособность Kafka сам по себе: сначала выполните проверку записи и чтения из предыдущего раздела.
Откройте на том же хосте http://localhost:8080. Благодаря DYNAMIC_CONFIG_ENABLED=true интерфейс покажет форму добавления кластера. Создайте подключение со значениями:
| Поле | Значение |
|---|---|
| Name | kafka-dev |
| Bootstrap servers | kafka:9092 |
После сохранения Kafbat UI перезапускает своё приложение, чтобы применить новое подключение. Подождите, пока на главной странице появится один кластер со статусом Online, затем откройте demo.orders.created, список сообщений и разделов. kafka:9092 указан потому, что интерфейс работает в Docker Compose-сети. Порт интерфейса также привязан к 127.0.0.1: не публикуйте его в сеть без отдельной аутентификации, TLS и согласованного доступа.
kcat
kcat — компактный Kafka-клиент для просмотра метаданных и чтения сообщений. Он полезен как независимая проверка, когда нужно увидеть, какие брокеры, темы и разделы обнаруживает клиент. Разовый запуск из контейнера:
sudo docker run --rm --network kafka-dev_default \
edenhill/kcat:1.7.1 \
-b kafka:9092 -L
-L выводит метаданные. Имя сети Compose обычно складывается из имени каталога и _default; если каталог или имя проекта другое, узнайте точное имя командой sudo docker network ls и подставьте его вместо kafka-dev_default.
Последние десять доступных сообщений из каждого раздела темы можно прочитать так:
sudo docker run --rm --network kafka-dev_default \
edenhill/kcat:1.7.1 \
-b kafka:9092 -C -t demo.orders.created -o -10 -e
При запуске kcat на самом хосте используйте -b localhost:29092. Образ edenhill/kcat:1.7.1 на момент проверки ориентирован на linux/amd64; на ARM Docker может предупредить об эмуляции. Для постоянной работы на ARM предпочтительнее нативный пакет kcat для вашей ОС.
Наблюдение за Dev/Test-контуром: что смотреть регулярно
Диагностика начинается не в момент аварии. Даже у учебного контура полезно иметь короткий регулярный маршрут проверки: отвечает ли брокер, существуют ли ожидаемые темы, движутся ли группы получателей и остаётся ли место для журналов. В рабочей среде эти же признаки превращают в метрики и оповещения; документация мониторинга Apache Kafka отдельно перечисляет показатели ошибок запросов, неполных копий, состояния каталогов журналов, задержки получателей и работы координатора групп.
1. Состояние контейнеров и журнал запуска
sudo docker compose ps
sudo docker compose logs --since=15m kafka
Для этого материала нормальное состояние брокера — healthy. Это означает, что встроенная проверка смогла выполнить запрос списка тем через протокол Kafka. Статус Up без healthy недостаточен: процесс может существовать, но ещё не принимать клиентов. Журнал нужен не для чтения всех строк подряд, а для поиска первой ошибки после запуска или перед сбоем: несовместимых ролей KRaft, проблем с каталогом данных, портом, адресом контроллера или правами на том.
Не ограничивайтесь строкой ERROR: обычный запуск Kafka тоже создаёт много служебных сообщений. Ищите первую запись с ошибкой и несколько десятков строк до неё. После изменения конфигурации сохраняйте саму изменённую настройку и время перезапуска — без этого невозможно отличить новый дефект от старого сообщения журнала.
2. Темы, ведущие копии и синхронность данных
sudo docker compose exec kafka \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--describe \
--topic demo.orders.created
Команда показывает каждую часть темы отдельно. На одноброкерном стенде ожидаются Leader: 1, Replicas: 1 и Isr: 1 для каждого раздела. Здесь Isr — список синхронных копий. В кластере из нескольких брокеров проблема возникает, если число Isr меньше числа Replicas: одна из копий отстала или недоступна. Если нет ведущей копии либо Isr уменьшился после отказа узла, не лечите это удалением темы — сначала зафиксируйте вывод --describe, состояние брокеров и журналы.
В этом учебном стенде одна копия данных, поэтому потеря брокера равна недоступности темы. Это ожидаемое ограничение, а не ложная тревога. В рабочем кластере отслеживают как минимум число неполностью скопированных разделов, число недоступных копий, скорость сжатия и расширения списка синхронных копий, а также свободное место в каталогах журналов.
3. Группы получателей и отставание
sudo docker compose exec kafka \
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--describe \
--group demo-inspector-01
В таблице особенно важны три значения:
| Поле | Что означает | Нормальный результат учебного теста |
|---|---|---|
CURRENT-OFFSET |
Позиция, до которой группа дошла в конкретном разделе | После одного прочитанного сообщения в его разделе значение увеличилось |
LOG-END-OFFSET |
Конец журнала в момент запроса | Не меньше текущей позиции |
LAG |
Сколько сообщений группа ещё не обработала: разница между концом журнала и её текущей позицией | 0, когда группа догнала поток |
Отставание — не всегда сбой. Оно закономерно растёт при всплеске входящих событий, остановке получателя или медленной внешней базе. Сигнал для разбора — отставание, которое растёт дольше согласованного времени, либо не уменьшается после восстановления получателя. Сначала определите, растёт ли конец журнала, движется ли текущая позиция и есть ли активные участники группы. Затем проверяйте само приложение: его журнал, время внешнего вызова, ошибки схемы данных и момент фиксации смещения. Добавление экземпляров поможет только до числа разделов темы.
Статус «нет активных участников» после завершения учебной консольной команды нормален: команда закончила работу, но Kafka сохранила позицию группы. Для постоянно работающего сервиса такой статус вместе с растущим LAG означает, что экземпляры не запущены, потеряли соединение или вышли из группы.
4. Место для данных и ресурсные признаки
df -h
sudo docker system df
Kafka хранит журналы в Docker-томе kafka-data. Следите не только за размером образов, но и за свободным местом файловой системы, на которой Docker хранит тома. Полный диск приводит не к «медленной Kafka», а к риску остановки записи, сбоя репликации и затяжного восстановления. Не удаляйте тома или не запускайте очистку Docker как первую реакцию: это может стереть единственную копию учебных данных. Сначала определите владельца тома и необходимость сохранения сообщений.
Для рабочего кластера к диску добавляют сеть, загрузку процессора и время обработки запросов. Метрики Kafka позволяют увидеть скорость входящих и исходящих данных, ошибки записи и чтения, заполнение очередей запросов и состояние каталогов журналов. Удалённый JMX по умолчанию выключен; если его включают для мониторинга, доступ к нему также защищают — иначе через него можно получить несанкционированное наблюдение или управление процессом.
Диагностика инцидента: от симптома к причине
Когда поток не работает, не начинайте с перезапуска и тем более с удаления контейнера или темы. Такой шаг уничтожает часть следов и может ухудшить состояние группы. Используйте одинаковую последовательность: зафиксировать время и симптом → проверить брокер → проверить тему → проверить отправителя → проверить группу получателей → проверить внешнюю систему. Она отделяет проблему Kafka от ошибки приложения или сети.
Сценарий 1. Брокер не запускается или не становится healthy
- Выполните
sudo docker compose psи убедитесь, что проверяете именно Kafka, а не только интерфейс управления. - Прочитайте журнал запуска командой
sudo docker compose logs --since=15m kafka. - Сверьте изменения с таблицей Kafka-параметров: роли KRaft, три точки подключения, адреса, сообщаемые клиентам, кворум контроллера и путь Docker-тома должны составлять одну согласованную конфигурацию.
- Проверьте, не занят ли локальный порт и есть ли свободное место. Не меняйте несколько переменных одновременно: после каждого исправления повторите
sudo docker compose configи запуск.
Типичные признаки: ошибка кворума или контроллера указывает на несогласованные KAFKA_CONTROLLER_* и адрес контроллера; отказ с каталогом данных — на путь тома, права или повреждённое/чужое хранилище; невозможность открыть порт — на конфликт с другим локальным сервисом. Не удаляйте kafka-data, если в нём могут быть нужные сообщения: удаление тома равносильно удалению данных этого стенда.
Сценарий 2. Отправитель не подключается или запись не появляется
Сначала выясните, где работает приложение. Из контейнера Compose оно подключается к kafka:9092, с Debian/Ubuntu-хоста — к localhost:29092. Адреса нельзя взаимозаменять: localhost внутри контейнера означает сам контейнер, а имя kafka вне Docker-сети обычно неизвестно. Если начальное соединение есть, но клиент сообщает ошибку после получения метаданных, почти всегда проверяют KAFKA_ADVERTISED_LISTENERS: именно этот параметр сообщает клиенту последующий адрес брокера.
Затем выполните kafka-topics.sh --describe для точного имени темы и сверьте имя, ключ и формат события с журналом отправителя. В защищённом контуре также проверьте срок сертификата, имя узла в сертификате, выбранный механизм SASL, учётную запись и ACL на запись. Ошибка авторизации не исправляется запуском отправителя от имени администратора: сначала определите, какое минимальное право требуется прикладной учётной записи.
Сценарий 3. Сообщение есть в теме, но бизнес-сервис его не обработал
Разделите два факта: «сообщение существует» и «конкретная группа его обработала». Первое подтверждает консольный получатель с --from-beginning, но только с новым одноразовым именем группы: сохранённая позиция существующей группы не будет сброшена этим параметром. Например, для безопасной проверки создайте отдельную диагностическую группу и не используйте имя бизнес-сервиса:
sudo docker compose exec kafka \
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic demo.orders.created \
--group diagnostic-01 \
--from-beginning \
--max-messages 1
Перед повторной проверкой замените окончание имени на новое, ещё не использованное. Такая команда создаёт лишь отдельную диагностическую позицию и не изменяет позицию рабочей группы. Kafbat UI также подходит для просмотра, но его данные подтверждайте независимой командой. Второй факт подтверждает kafka-consumer-groups.sh --describe --group <имя-группы>. Если LAG положителен и растёт, сервис не успевает или не работает. Если LAG равен нулю, но бизнес-эффекта нет, ищите проблему после Kafka: обработчик мог зафиксировать смещение слишком рано, отбросить событие как дубликат, получить ошибку схемы или не выполнить внешнюю операцию.
Не меняйте начальную позицию группы и не «сбрасывайте» её на начало во время разбора без отдельного решения владельца потока. Повторное чтение может повторить внешние действия: отправку уведомления, создание записи или изменение статуса. Сначала сохраните имя группы, смещения, версию обработчика и пример проблемного event_id; затем проверяйте идемпотентность и согласованный сценарий переигрывания.
Сценарий 4. Растёт отставание группы получателей
Отставание показывает очередь необработанных записей, но не называет причину. Смотрите динамику несколько раз: растёт ли LOG-END-OFFSET, движется ли CURRENT-OFFSET, появились ли активные участники. Если конец журнала растёт быстрее текущей позиции, получатель медленнее входящего потока. Если обе позиции остановились, проверяйте отправителя или саму тему. Если часть разделов отстаёт сильнее других, вероятен перекос ключей: один ключ создаёт «горячий» раздел.
Возможные действия выбирают только после причины: оптимизировать обработчик и внешнюю зависимость, увеличить число экземпляров до числа разделов, изменить число разделов с учётом последствий для ключа и порядка либо пересмотреть контракт разбиения. Не масштабируйте группу вслепую: четвёртый экземпляр при трёх разделах не уменьшит отставание.
Сценарий 5. В кластере из нескольких брокеров уменьшился Isr или нет ведущей копии
Сравните в --describe поля Replicas и Isr по каждому разделу. Если синхронных копий меньше, зафиксируйте затронутые темы, брокеры и время, затем проверьте сеть, диск и журналы отставшего брокера. Во время восстановления не переносите данные вручную и не отключайте защитные ограничения ради скорости: это повышает риск потери сообщений. При отсутствии ведущей копии или повторяющихся переключениях требуется процедура эксплуатации кластера, а не команды одноброкерного стенда.
Сценарий 6. Веб-интерфейс доступен, но показывает пустой или недоступный кластер
Сначала подтвердите Kafka штатной консольной командой. Затем откройте sudo docker compose logs --since=15m kafka-ui и проверьте адрес интерфейса: внутри Compose-сети он должен быть kafka:9092, не localhost:29092. После сохранения подключения Kafbat UI перезапускает своё приложение; дождитесь статуса Online. Если Kafka работает, а интерфейс — нет, это проблема диагностического инструмента, а не доказательство потери данных в Kafka.
Что приложить к обращению или сохранить перед изменениями
| Артефакт | Зачем нужен |
|---|---|
Время начала проблемы и затронутая тема, группа, event_id |
Позволяет сопоставить журналы брокера и приложения, не искать «все ошибки за неделю» |
Вывод docker compose ps и kafka-topics.sh --describe |
Показывает доступность брокера, ведущие и синхронные копии без изменения данных |
Вывод kafka-consumer-groups.sh --describe --group <имя> |
Фиксирует позиции, отставание и назначение разделов конкретной группы |
| Журналы Kafka и приложения за ограниченный интервал | Помогают увидеть первичную ошибку и исключают нерелевантный шум |
| Последнее изменение конфигурации, версии образа или сертификата | Позволяет связать сбой с изменением, а не угадывать причину |
Пароли, закрытые ключи, полные персональные данные и содержимое секретов в обращение не включают. Если нужно показать событие, маскируйте бизнес-данные и оставляйте технические поля: имя темы, раздел, смещение, время и event_id.
Что нужно спроектировать до рабочей эксплуатации
Этот Compose-файл подходит для разработки, демонстрации и проверки контракта. Переход к рабочей эксплуатации — не увеличение числа копий одной строкой. До него необходимо принять и проверить как минимум следующие решения:
- Топология и отказоустойчивость. Число брокеров, кворум контроллеров, число копий данных, минимальное число синхронных копий, диски, зоны отказа и процедура восстановления.
- Сеть и доступ. TLS, SASL или другой механизм аутентификации, правила доступа ACL, хранение секретов, сегментация сети и безопасный доступ к интерфейсу.
- Контракты данных. Владелец темы, схема и её эволюция, ключ,
event_id, правила обратной совместимости, классификация данных и правило хранения. - Обработка ошибок. Идемпотентность получателя, момент фиксации смещения, повтор, согласованный маршрут ошибок и операционная процедура повторной обработки.
- Наблюдаемость. Метрики брокеров и отставание получателей, журналы, оповещения, планирование ёмкости, контроль объёма данных и проверка восстановления.
- Эксплуатация. Обновления, резервное копирование, восстановление, проверка изменений на стенде и владельцы операционных действий.
Первая работающая тема — начало проектирования, а не доказательство готовности к критичному потоку. Начинайте с синтетических данных и некритичного сценария: так команда проверит контракт и операции, не смешивая обучение Kafka с риском для бизнеса.
Важное уведомление
Материал носит информационный характер и не является индивидуальной консультацией, проектной документацией, инструкцией по внедрению или гарантией результата. Применимость описанных подходов зависит от процессов, систем, данных, требований безопасности и иных условий конкретной организации. Перед внедрением или изменением ИТ‑систем необходимо самостоятельно оценить риски и привлечь профильных специалистов.
«Интеграция начинается с факта: заказ создан, заявка согласована, статус доставки изменился, платёж отклонён.»
- 01Kafka полезна для независимой обработки событий и повторного чтения потока; она не заменяет каждый HTTP-вызов или очередь задач.
- 02Учебный одноброкерный Compose-контур проверяет контракт и инструменты, но не заменяет отказоустойчивый защищённый кластер.
- 03Число разделов, ключи сообщений, фиксация смещений и правила повторов определяют корректность обработки не меньше, чем параметры брокера.
