Kafka Producer — это клиент, который формирует и отправляет сообщения (records) в нужные топики Kafka. Он сам выбирает партицию, буферизует, батчит, сжимает данные и управляет подтверждениями доставки. Взаимодействие с Kafka broker идёт по бинарному протоколу: producer запрашивает метаданные кластера и отправляет запросы Produce на лидеры партиций, получая от них подтверждения и смещения сообщений.
Kafka Producer — это библиотека-клиент, встроенная в ваше приложение, которая отвечает за надёжную и эффективную доставку сообщений в Kafka cluster. Он не хранит данные сам, а только подготавливает их и отправляет на соответствующие брокеры, которые уже записывают сообщения в лог партиций.
Основные задачи Kafka Producer:
topic, опциональный ключ и значение (payload), а также заголовки.acks, retries, delivery.timeout.ms управляет тем, насколько гарантированной должна быть доставка.Схема взаимодействия Kafka Producer с Kafka broker выглядит так:
1. При старте producer обращается к так называемому bootstrap server (списку брокеров), чтобы получить metadata о кластере: какие есть топики, их партиции, какие брокеры являются лидерами этих партиций.
2. Producer строит у себя в памяти карту: «партиция → лидер-брокер», и для каждой партиции формирует очередь записей.
3. Когда срабатывает порог по размеру батча или таймауту, producer формирует сетевой запрос Produce и по TCP-соединению отправляет батч сообщений напрямую на брокер-лидер соответствующей партиции.
4. Брокер принимает батч, добавляет записи в лог партиции, реплицирует на follower-реплики (в зависимости от настроек кластера и acks), и затем возвращает ответ producer-у — со смещением (offset) каждой записи и возможными кодами ошибок.
5. В зависимости от уровня подтверждений acks:
— acks=0: broker не шлёт подтверждение, producer не ждёт ответа, максимальная скорость, но есть риск потерь сообщений.
— acks=1: подтверждает только лидер после записи на свой диск — баланс между скоростью и надёжностью.
— acks=all (или -1): лидер ждёт подтверждений от всех реплик в ISR, максимальная надёжность при меньшей скорости.
Важные дополнительные механизмы взаимодействия:
Повторные отправки и порядок сообщений. При временных сбоях сети или брокера producer может автоматически ретраить запросы. Чтобы при этом не нарушать порядок сообщений в пределах одной партиции, используются настройки max.in.flight.requests.per.connection и идемпотентный режим (enable.idempotence=true), который предотвращает дубликаты при повторной отправке.
Идемпотентный и транзакционный producer. Идемпотентный producer получает от кластера уникальный producerId и ведёт счётчики последовательности сообщений; брокер отбрасывает дубликаты по этим метаданным, обеспечивая semantics «не больше одного раза». Транзакционный producer поверх этого объединяет несколько записей (даже в разные топики/партиции) в одну атомарную транзакцию, которая либо целиком видна потребителям, либо целиком откатывается.
Отметьте свой прогресс