Sobes.tech
Junior — Middle

Как се осигурява възможността за доставка на съобщения поне веднъж в Kafka?

sobes.tech AI

Отговор от AI

В Apache Kafka доставката на съобщения с гаранция " поне веднъж " се осигурява чрез механизма за потвърждения (acknowledgments) и управлението на офсетите (offsets).

Процесът изглежда така:

  1. Производителят изпраща съобщение към Kafka и чака потвърждение от брокера (например, ack=all — потвърждение от всички реплики).
  2. Съобщението се запазва в журнала (логовете) на брокера.
  3. Консуматорът чете съобщението и след успешно обработване изпраща потвърждение (commit) на офсета (offset) към Kafka.

Ако консуматорът се срине преди потвърждението, при рестартиране ще прочете съобщението отново, което гарантира доставката "поне веднъж".

Пример на Go с използване на библиотеката sarama:

consumer, err := sarama.NewConsumerGroup(brokers, groupID, config)
// ...
// В обработчика на съобщения:
func (h *handler) ConsumeClaim(sess sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
    for message := range claim.Messages() {
        process(message) // обработка на съобщението
        sess.MarkMessage(message, "") // потвърждение за обработка
    }
    return nil
}

Така, ако потвърждението не бъде изпратено, съобщението ще бъде доставено отново.