Какие гарантии доставки сообщений предоставляет Apache Kafka?

«Какие гарантии доставки сообщений предоставляет Apache Kafka?» — вопрос из категории Брокеры сообщений, который задают на 23% собеседований Devops Инженер. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

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

  1. At-most-once (максимум один раз): Сообщения могут быть потеряны, но никогда не будут доставлены дублями. Достигается отправкой с acks=0 (producer не ждет подтверждения от брокера) или retries=0. Подходит для метрик или логов, где допустима некоторая потеря данных.

  2. At-least-once (минимум один раз): Сообщение гарантированно будет доставлено, но в случае ошибок сети или ребалансировки потребителя возможны дубли. Стандартная настройка при acks=all (producer ждет подтверждения от всех in-sync реплик) и retries > 0. Потребитель должен быть идемпотентным или уметь обрабатывать дубли.

  3. Exactly-once (ровно один раз): Гарантирует, что сообщение будет обработано ровно один раз, без потерь и дублей. Требует комбинации:

    • Идемпотентный producer (enable.idempotence=true): предотвращает дубли из-за повторных отправок от producer.
    • Транзакции между producer и consumer (transactional.id): обеспечивают атомарную запись в несколько партиций/топиков.
    • Изоляция read_committed у consumer: потребитель читает только закоммиченные транзакции.
// Пример конфигурации producer для exactly-once семантики
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

// Ключевые настройки
props.put("enable.idempotence", "true"); // Включает идемпотентность
props.put("acks", "all"); // Ждем подтверждения от всех реплик
props.put("transactional.id", "my-app-transactional-id"); // Уникальный ID для транзакций

Producer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions(); // Инициализация транзакций

try {
    producer.beginTransaction();
    producer.send(new ProducerRecord<>("orders", "key1", "value1"));
    producer.send(new ProducerRecord<>("payments", "key2", "value2"));
    producer.commitTransaction(); // Атомарный коммит в оба топика
} catch (ProducerFencedException e) {
    producer.close();
}