多Kafka连接配置咨询:多Broker监听与单Broker写入实现方案
多Kafka连接配置方案
一、生产者配置(单集群写入)
基于你提供的单个生产者配置,优化为从配置文件读取参数,保持单集群写入逻辑:
配置文件(application.yml)
kafka: producer: bootstrap-servers: "192.168.1.100:9092" # 目标写入的Kafka Broker地址 key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer
生产者配置类
@Configuration public class KafkaProducerConfig { @Value("${kafka.producer.bootstrap-servers}") private String bootstrapServers; @Value("${kafka.producer.key-serializer}") private Class<? extends Serializer<String>> keySerializer; @Value("${kafka.producer.value-serializer}") private Class<? extends Serializer<String>> valueSerializer; @Bean public Map<String, Object> producerConfigs() { Map<String, Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, keySerializer); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, valueSerializer); // 可按需添加acks、retries等其他生产者配置 return props; } @Bean public ProducerFactory<String, String> producerFactory() { return new DefaultKafkaProducerFactory<>(producerConfigs()); } @Bean public KafkaTemplate<String, String> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); } }
二、多消费者配置(多集群监听)
要监听多个不同IP的Kafka集群,需为每个集群单独创建一套消费者配置Bean,通过不同的Bean名称区分配置组:
配置文件新增多集群消费者配置(application.yml)
kafka: consumer1: bootstrap-servers: "192.168.1.101:9092" # 第一个监听集群的Broker地址 group-id: "consumer-group-1" key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer auto-offset-reset: earliest consumer2: bootstrap-servers: "192.168.1.102:9092" # 第二个监听集群的Broker地址 group-id: "consumer-group-2" key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer auto-offset-reset: latest
多消费者配置类
@Configuration public class MultiKafkaConsumerConfig { // ------------------- 第一个Kafka集群消费者配置 ------------------- @Value("${kafka.consumer1.bootstrap-servers}") private String consumer1BootstrapServers; @Value("${kafka.consumer1.group-id}") private String consumer1GroupId; @Value("${kafka.consumer1.key-deserializer}") private Class<? extends Deserializer<String>> consumer1KeyDeserializer; @Value("${kafka.consumer1.value-deserializer}") private Class<? extends Deserializer<String>> consumer1ValueDeserializer; @Value("${kafka.consumer1.auto-offset-reset}") private String consumer1AutoOffsetReset; @Bean("consumer1Configs") public Map<String, Object> consumer1Configs() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, consumer1BootstrapServers); props.put(ConsumerConfig.GROUP_ID_CONFIG, consumer1GroupId); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, consumer1KeyDeserializer); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, consumer1ValueDeserializer); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, consumer1AutoOffsetReset); return props; } @Bean("consumer1Factory") public ConsumerFactory<String, String> consumer1Factory() { return new DefaultKafkaConsumerFactory<>(consumer1Configs()); } @Bean("kafkaListenerContainerFactory1") public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory1() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumer1Factory()); // 可按需设置concurrency等容器工厂参数 return factory; } // ------------------- 第二个Kafka集群消费者配置 ------------------- @Value("${kafka.consumer2.bootstrap-servers}") private String consumer2BootstrapServers; @Value("${kafka.consumer2.group-id}") private String consumer2GroupId; @Value("${kafka.consumer2.key-deserializer}") private Class<? extends Deserializer<String>> consumer2KeyDeserializer; @Value("${kafka.consumer2.value-deserializer}") private Class<? extends Deserializer<String>> consumer2ValueDeserializer; @Value("${kafka.consumer2.auto-offset-reset}") private String consumer2AutoOffsetReset; @Bean("consumer2Configs") public Map<String, Object> consumer2Configs() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, consumer2BootstrapServers); props.put(ConsumerConfig.GROUP_ID_CONFIG, consumer2GroupId); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, consumer2KeyDeserializer); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, consumer2ValueDeserializer); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, consumer2AutoOffsetReset); return props; } @Bean("consumer2Factory") public ConsumerFactory<String, String> consumer2Factory() { return new DefaultKafkaConsumerFactory<>(consumer2Configs()); } @Bean("kafkaListenerContainerFactory2") public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory2() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumer2Factory()); return factory; } }
三、使用多集群监听
编写消费者监听方法时,通过containerFactory参数指定对应容器工厂,即可实现对不同Kafka集群的监听:
@Component public class MultiKafkaConsumer { // 监听第一个Kafka集群的topic @KafkaListener(topics = "topic-from-cluster1", containerFactory = "kafkaListenerContainerFactory1") public void listenCluster1(String message) { System.out.println("Received from cluster1: " + message); // 业务逻辑处理 } // 监听第二个Kafka集群的topic @KafkaListener(topics = "topic-from-cluster2", containerFactory = "kafkaListenerContainerFactory2") public void listenCluster2(String message) { System.out.println("Received from cluster2: " + message); // 业务逻辑处理 } }
四、扩展说明
- 若需监听更多Kafka集群,只需复制上述消费者配置模块,新增对应配置参数和Bean即可。
- 所有配置参数建议放入配置文件,避免硬编码,便于后续维护修改。
- 生产者和消费者的其他配置(如重试机制、批量处理等)可根据业务需求添加到对应props中。
内容的提问来源于stack exchange,提问作者IgorPiven
相关产品推荐
相关产品推荐

