如何使用Reactor Kafka KafkaSender API向不同集群的两个主题发送消息?
解决方案
因为你要发送消息到两个不同的Kafka集群,单个KafkaSender实例只能绑定一个集群的配置,所以必须为每个集群单独创建KafkaSender,分别执行发送操作。
步骤说明
为每个集群配置独立的Producer参数
每个Kafka集群需要单独的ProducerConfig配置,核心是指定各自的bootstrap.servers:// 第一个集群的配置 Map<String, Object> firstClusterConfig = new HashMap<>(); firstClusterConfig.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-first-broker.com:9092"); firstClusterConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); firstClusterConfig.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); // 可添加其他必要配置(如acks、retries等) // 第二个集群的配置 Map<String, Object> anotherClusterConfig = new HashMap<>(); anotherClusterConfig.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-another-broker-not-the-same-as-above.com:9093"); anotherClusterConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); anotherClusterConfig.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); // 可添加其他必要配置创建对应集群的KafkaSender实例
每个配置对应一个独立的KafkaSender:KafkaSender<String, String> firstClusterSender = KafkaSender.create(SenderOptions.create(firstClusterConfig)); KafkaSender<String, String> anotherClusterSender = KafkaSender.create(SenderOptions.create(anotherClusterConfig));构造SenderRecord并分别发送
针对每个主题构造SenderRecord,通过对应的Sender发送。如果需要等待两个发送操作都完成,可以用Mono.zip合并结果:String myMessage = "your-message-content"; // 发送到第一个集群的first_topic Mono<SenderResult<String>> sendToFirst = firstClusterSender.send(Mono.just( SenderRecord.create(new ProducerRecord<>("first_topic", null, myMessage), myMessage) )).next(); // 发送到第二个集群的another_topic Mono<SenderResult<String>> sendToAnother = anotherClusterSender.send(Mono.just( SenderRecord.create(new ProducerRecord<>("another_topic", null, myMessage), myMessage) )).next(); // 合并两个发送操作,等待全部完成 Mono.zip(sendToFirst, sendToAnother) .doOnSuccess(tuple -> { // 处理成功结果,tuple.getT1()对应第一个发送的结果,tuple.getT2()对应第二个 System.out.println("两个主题消息发送完成"); }) .doOnError(error -> { // 处理发送失败的情况 error.printStackTrace(); }) .block(); // 同步场景用block等待结果;异步场景直接subscribe即可
为什么之前的方法无效?
你尝试的逗号分隔主题("first_topic,another_topic")仅适用于同一个Kafka集群下的多主题发送,本质是让Producer把同一条消息发往同一个集群内的多个主题。但你的场景是两个完全独立的集群,单个Producer无法同时连接两个不同的集群,所以这种方式完全不适用。
内容的提问来源于stack exchange,提问作者PatPanda
相关产品推荐
相关产品推荐

