JavaRush /Курсы /Модуль 5. Spring /Лекция 190: Доставка сообщений и гарантии Kafka: "At leas...

Лекция 190: Доставка сообщений и гарантии Kafka: "At least once", "At most once", "Exactly once"

Модуль 5. Spring
19 уровень , 9 лекция
Открыта

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

  • "At most once" — "Не больше одного раза".
  • "At least once" — "Не меньше одного раза".
  • "Exactly once" — "Только один раз".

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

"At most once" (Не больше одного раза)

Итак, представьте, что вы отправляете другу сообщение, но хотите быть уверенным, что отправите его максимум один раз. Если вы случайно удалили сообщение до отправки или произошел сбой — вы просто "забьёте". Нет сообщения — нет проблем?

"At most once" означает, что сообщение либо доставлено, либо потеряно. Kafka не будет пытаться его повторно отправить. Это самый быстрый способ доставки, но он не гарантирует, что данные достигнут получателя.

Пример настройки "At most once":

properties.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true"); // Автокоммит смещения
  • Плюсы: высокая производительность, минимальная задержка.
  • Минусы: возможная потеря данных.

"At least once" (Не меньше одного раза)

А теперь представьте, что ваш друг говорит: "Лучше пришли мне сообщение дважды, чем вообще не пришли". Вот вам стратегия "At least once", когда сообщение может быть доставлено несколько раз, но уж точно хотя бы раз.

Kafka будет повторно отправлять сообщение, если не получит подтверждение от системы, что оно обработано. Если ваш консьюмер сбоит, то при следующем запуске он начнет обработку с последнего "не зафиксированного" сообщения.

Настройка:

properties.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); // Ручной коммит смещения

Пример ручного коммита:


consumer.commitSync(); // Исполняет явный коммит после успешной обработки сообщения
  • Плюсы: надежная доставка.
  • Минусы: возможны дубли сообщений.

"Exactly once" (Только один раз)

Теперь представьте, что ваш друг — перфекционист. "Никаких дубликатов, никаких потерь, я хочу, чтобы сообщение пришло ровно один раз". Это стратегия "Exactly once".

Kafka вводит поддержку "Exactly once" через механизм транзакций и идемпотентность. Это значит, что сообщение будет доставлено ровно один раз, даже если произойдет сбой на продюсере или консьюмере.

Настройка:

properties.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true"); // Включаем идемпотентность

Использование транзакций:


producer.initTransactions(); // Инициализация транзакций
producer.beginTransaction();
producer.send(record);
producer.commitTransaction(); // Коммит транзакции
  • Плюсы: высочайшая надежность данных.
  • Минусы: более высокая задержка и сложность конфигурации.

Как настроить правильную гарантию доставки

Давайте разберем, как эти гарантии работают на практике. Рассмотрим каждую стратегию более детально.

Настройка продюсера

Kafka-продюсер играет ключевую роль в достижении нужного уровня доставки. Вот несколько параметров, которые вам нужно знать:

1. acks (Acknowledgments): Этот параметр определяет, сколько подтверждений требуется от брокеров.

  • acks=0 : Продюсер не ждет подтверждения от брокера (подходит для "At most once").
  • acks=1 : Запись считается успешной после подтверждения лидера.
  • acks=all : Все реплики должны подтвердить запись (для "At least once" и "Exactly once").

Пример настройки:


properties.put(ProducerConfig.ACKS_CONFIG, "all"); // Все реплики подтверждают запись

2. idempotence (Идемпотентность): повторная отправка того же сообщения не создаст дубликатов.

  • Для "Exactly once" сделайте это обязательным:
    properties.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
    

3. retries (Повторы): определяет, сколько раз продюсер будет пытаться отправить сообщение, если произошла временная ошибка.

properties.put(ProducerConfig.RETRIES_CONFIG, Integer.MAX_VALUE);

Настройка консьюмера

Консьюмеры тоже должны быть правильно настроены для достижения необходимых гарантий.

  • 1. Auto Commit (Автоматический коммит):
    • Для "At most once" включите:
      
      properties.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true");
      
  • 2. Ручное управление смещением (Offsets):
    • Для "At least once" и "Exactly once" фиксируйте смещение после успешной обработки сообщения:
      
      try {
          consumer.commitSync();
      } catch (CommitFailedException e) {
          // Логируем ошибку
      }
      

Практический пример: "Exactly once"

Давайте напишем простое приложение, которое будет отправлять и получать сообщения с гарантией "Exactly once".

Продюсер с транзакциями:


Properties producerProps = new Properties();
producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
producerProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true"); // Идемпотентность

KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps);

producer.initTransactions(); // Инициализация транзакции

try {
    producer.beginTransaction(); // Начало транзакции
    producer.send(new ProducerRecord<>("topic1", "key", "value"));
    producer.commitTransaction(); // Завершение транзакции
} catch (ProducerFencedException e) {
    producer.abortTransaction(); // Откат транзакции в случае ошибки
} finally {
    producer.close();
}

Консьюмер с ручным коммитом:


Properties consumerProps = new Properties();
consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "group1");
consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); // Ручной коммит

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps);
consumer.subscribe(Arrays.asList("topic1"));

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, String> record : records) {
        // Обработка сообщения
        System.out.printf("Consumed record: key = %s, value = %s%n", record.key(), record.value());
    }
    consumer.commitSync(); // Подтверждение обработки
}

Ошибки и их предотвращение

При настройке гарантий доставки важно учитывать:

  • Дублирование сообщений: даже с "Exactly once" вы можете получить дубли, если транзакции не настроены корректно.
  • Потеря данных: "At most once" чревата потерями при любых сбоях.
  • Производительность: "Exactly once" требует больше ресурсов, поэтому важно учитывать этот фактор при проектировании системы.

Легендарная Kafka даёт вам инструменты для каждой ситуации. Главное — знать, какие гарантии нужны вашему приложению.

Комментарии
ЧТОБЫ ПОСМОТРЕТЬ ВСЕ КОММЕНТАРИИ ИЛИ ОСТАВИТЬ КОММЕНТАРИЙ,
ПЕРЕЙДИТЕ В ПОЛНУЮ ВЕРСИЮ