Проверяет понимание настройки Kafka для конкурентного потребления сообщений несколькими потоками и инстансами приложения.
Kafka хранит сообщения в партициях топика. Чтобы несколько инстансов приложения обрабатывали сообщения параллельно, они должны быть в одной группе потребителей (consumer group). Kafka гарантирует, что каждая партиция назначается только одному потребителю из группы, поэтому сообщения из одной партиции обрабатываются последовательно, но разные партиции могут обрабатываться разными потребителями одновременно.
Важно помнить, что объект KafkaConsumer не является потокобезопасным. Если несколько потоков используют один и тот же экземпляр, возникают блокировки и ошибки. Правильный подход — создавать отдельный KafkaConsumer для каждого потока. Например, если у вас 4 потока в одном инстансе, создайте 4 потребителя, каждый из которых подписан на свой набор партиций.
// Каждый поток создает свой KafkaConsumer
for (int i = 0; i < numThreads; i++) {
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList(topic));
executor.submit(() -> {
while (running) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
process(record);
}
}
});
}Здесь каждый поток имеет собственный consumer, что исключает блокировки. Количество потоков лучше выбирать равным числу партиций топика, чтобы максимально использовать параллелизм.
Если у вас несколько инстансов приложения, просто укажите одинаковый group.id в конфигурации. Kafka автоматически распределит партиции между всеми потребителями группы. Например, топик с 12 партициями и 3 инстансами по 4 потока каждый даст 12 потребителей, каждый получит по одной партиции.
Используйте группы потребителей и отдельные KafkaConsumer на поток, чтобы избежать блокировок и обеспечить масштабируемость. Это стандартный паттерн для высоконагруженных систем на Kafka.