Как в Kafka обеспечить параллельную обработку сообщений топика двумя потребителями, назначив каждому свою партицию?

«Как в Kafka обеспечить параллельную обработку сообщений топика двумя потребителями, назначив каждому свою партицию?» — вопрос из категории Брокеры сообщений, который задают на 10% собеседований Java Разработчик. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

В 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));

Ключевые правила:

  • Максимальный параллелизм = количество партиций. Если потребителей больше, чем партиций, "лишние" потребители будут бездействовать.
  • Гарантия порядка сохраняется только в пределах одной партиции.
  • Для топика с одной партицией два потребителя в одной группе не смогут работать параллельно — один из них останется неактивным.