Ответ
В Kafka параллелизм обработки определяется количеством партиций в топике. Два потребителя в одной consumer group могут работать бесперебойно, если топик имеет как минимум две партиции, и каждый потребитель закрепится за своей партицией.
Способ 1: Автоматическое распределение (рекомендуется)
Используйте subscribe(). Kafka автоматически сбалансирует партиции между потребителями в группе.
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "my-app"); // Одинаковый group.id для всех потребителей группы
props.put("key.deserializer", StringDeserializer.class);
props.put("value.deserializer", StringDeserializer.class);
KafkaConsumer<String, String> consumer1 = new KafkaConsumer<>(props);
KafkaConsumer<String, String> consumer2 = new KafkaConsumer<>(props);
// Оба потребителя подписываются на один топик
consumer1.subscribe(List.of("my-topic"));
consumer2.subscribe(List.of("my-topic"));
// Kafka назначит: Consumer1 -> Partition 0, Consumer2 -> Partition 1
Способ 2: Ручное назначение партиций
Используйте assign() для точного контроля. При этом отключается автоматический rebalance.
// Потребитель 1 закрепляется за партицией 0
TopicPartition partition0 = new TopicPartition("my-topic", 0);
consumer1.assign(List.of(partition0));
// Потребитель 2 закрепляется за партицией 1
TopicPartition partition1 = new TopicPartition("my-topic", 1);
consumer2.assign(List.of(partition1));
Ключевые правила:
- Максимальный параллелизм = количество партиций. Если потребителей больше, чем партиций, "лишние" потребители будут бездействовать.
- Гарантия порядка сохраняется только в пределах одной партиции.
- Для топика с одной партицией два потребителя в одной группе не смогут работать параллельно — один из них останется неактивным.