Ответ
Kafka позволяет настраивать баланс между надежностью и производительностью, предлагая три основные семантики доставки:
-
At-most-once (максимум один раз): Сообщения могут быть потеряны, но никогда не будут доставлены дублями. Достигается отправкой с
acks=0(producer не ждет подтверждения от брокера) илиretries=0. Подходит для метрик или логов, где допустима некоторая потеря данных. -
At-least-once (минимум один раз): Сообщение гарантированно будет доставлено, но в случае ошибок сети или ребалансировки потребителя возможны дубли. Стандартная настройка при
acks=all(producer ждет подтверждения от всех in-sync реплик) иretries > 0. Потребитель должен быть идемпотентным или уметь обрабатывать дубли. -
Exactly-once (ровно один раз): Гарантирует, что сообщение будет обработано ровно один раз, без потерь и дублей. Требует комбинации:
- Идемпотентный producer (
enable.idempotence=true): предотвращает дубли из-за повторных отправок от producer. - Транзакции между producer и consumer (
transactional.id): обеспечивают атомарную запись в несколько партиций/топиков. - Изоляция
read_committedу consumer: потребитель читает только закоммиченные транзакции.
- Идемпотентный producer (
// Пример конфигурации 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();
}