Producer и Consumer
Producer — отправляет сообщения в Kafka. Consumer — читает их оттуда. Они не знают друг о друге напрямую — оба общаются только с Kafka, что и обеспечивает развязку (см. тему «Зачем очереди»).
// Producer — отправка
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("orders", "order-123", "{...данные заказа...}"));
// Consumer — получение
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(List.of("orders"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
process(record.value());
}
}
Копнуть глубже
Consumer Group — несколько consumer’ов, работающих вместе над одним топиком, делят нагрузку между собой. Каждое сообщение из топика обрабатывается только одним consumer’ом из группы — это даёт горизонтальное масштабирование обработки:
props.put("group.id", "email-service-group"); // все consumer'ы с этим group.id делят нагрузку
Если в группе 3 consumer’а и топик разбит на 3 партиции (см. тему «Offset, партиции») — каждый consumer обрабатывает свою партицию параллельно, обработка идёт в 3 раза быстрее, чем одним consumer’ом.
poll() — модель вытягивания (pull), а не проталкивания (push). Consumer сам решает, когда забрать новую порцию сообщений, а не Kafka “толкает” их насильно — это даёт consumer’у контроль над собственной скоростью обработки, не позволяя перегрузить себя потоком сообщений быстрее, чем он успевает их обрабатывать.
• что такое Consumer Group и как она помогает масштабировать обработку (если дошёл до 2-го слоя).