如何使用KafkaClusterTestKit在kraft模式下配置自定义端口?
在KafkaClusterTestKit的KRaft模式下配置自定义端口
要在KRaft模式下通过KafkaClusterTestKit指定自定义端口,核心是通过KafkaClusterTestKit.Builder注入自定义配置,覆盖默认的随机端口逻辑,具体实现方式如下:
构建自定义Broker配置
创建Properties对象,明确指定KRaft模式下Broker的端口及必要配置:listeners:设置Broker监听的地址与端口(需包含控制器监听端口,如PLAINTEXT://localhost:9092,CONTROLLER://localhost:9093)advertised.listeners:对外暴露的地址与端口,通常与listeners中的PLAINTEXT地址一致- 补充KRaft基础配置:
process.roles、node.id、controller.quorum.voters等
注入配置并启动TestKit
使用KafkaClusterTestKit.Builder的withBrokerProperties方法传入自定义配置,替代默认的随机端口生成逻辑。示例代码(Scala)
import kafka.test.KafkaClusterTestKit import java.util.Properties val customBrokerProps = new Properties() // 指定自定义端口 customBrokerProps.put("listeners", "PLAINTEXT://localhost:9092,CONTROLLER://localhost:9093") customBrokerProps.put("advertised.listeners", "PLAINTEXT://localhost:9092") // KRaft必填配置 customBrokerProps.put("process.roles", "broker,controller") customBrokerProps.put("node.id", "1") customBrokerProps.put("controller.quorum.voters", "1@localhost:9093") customBrokerProps.put("controller.listener.names", "CONTROLLER") val testKit = KafkaClusterTestKit.Builder .withNumBrokers(1) .withKraft() .withBrokerProperties(customBrokerProps) .build() testKit.start()注意事项
- 多Broker场景下需为每个Broker配置独立端口,避免冲突
controller.quorum.voters中的控制器地址必须与listeners里的CONTROLLER端口匹配- 该方式与ZK模式下通过
TestUtils指定端口的逻辑一致,都是通过覆盖默认配置实现自定义端口
内容的提问来源于stack exchange,提问作者sobychacko
相关产品推荐
相关产品推荐

