如何通过SslBundle以编程方式配置Kafka Admin Client的SSL属性及发布消息触发Admin Client的问题解析
我来帮你拆解这个问题,先搞清楚为什么会触发Admin Client,再给你用SslBundle编程配置它的具体方案:
一、为什么发布消息会触发Admin Client?
你看到Admin Client被调用,核心原因是Spring Kafka默认会在生产者发送消息前检查目标主题是否存在(这个行为由KafkaTemplate的checkTopicExists属性控制,默认值为true)。哪怕你确认主题已经存在,这个检查逻辑还是会自动触发,进而创建Admin Client实例。
而你的问题在于:这个自动创建的Admin Client没有继承生产者的SSL配置,导致它无法连接SSL保护的Kafka集群,所以才会出现Node -1 disconnected、Cancelled in-flight API_VERSIONS request这类连接失败的日志。
二、用SslBundle编程配置Admin Client的两种方案
方案1:自定义AdminClient Bean并关联到生产者工厂
这种方案适合你后续可能还要用AdminClient做其他操作(比如管理主题、查看集群元数据)的场景,直接给AdminClient配置好SslBundle,让生产者工厂复用它:
第一步:创建自定义AdminClient Bean
@Configuration public class KafkaAdminConfiguration { @Bean public AdminClient kafkaAdminClient(SslBundles sslBundles) { Map<String, Object> adminProps = new HashMap<>(); // 基础集群配置 adminProps.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "mykafka.com:9092"); adminProps.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SSL"); // 配置SslBundle相关的SSL引擎工厂 adminProps.put(SslConfigs.SSL_ENGINE_FACTORY_CLASS_CONFIG, SslBundleSslEngineFactory.class.getName()); adminProps.put(SslBundle.class.getName(), sslBundles.getBundle("my-bundle")); return AdminClient.create(adminProps); } }
第二步:让生产者工厂使用这个自定义AdminClient
修改你现有的KafkaProducerConfiguration,给DefaultKafkaProducerFactory设置adminClientSupplier,指定用我们刚才创建的AdminClient来做主题检查:
@Configuration public class KafkaProducerConfiguration { @Bean public KafkaTemplate<String, MyPojo> kafkaTemplate(final MeterRegistry meterRegistry, final SslBundles sslBundles, final AdminClient kafkaAdminClient) { final Map<String, Object> properties = new ConcurrentHashMap<>(); properties.put("security.protocol", "SSL"); properties.put(SslConfigs.SSL_ENGINE_FACTORY_CLASS_CONFIG, SslBundleSslEngineFactory.class.getName()); properties.put(SslBundle.class.getName(), sslBundles.getBundle("my-bundle")); properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "mykafka.com:9092"); properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); final DefaultKafkaProducerFactory<String, MyPojo> defaultKafkaProducerFactory = new DefaultKafkaProducerFactory<>(properties); // 关联自定义AdminClient,用于主题检查 defaultKafkaProducerFactory.setAdminClientSupplier(() -> kafkaAdminClient); defaultKafkaProducerFactory.addListener(new MicrometerProducerListener<>(meterRegistry)); KafkaTemplate<String, MyPojo> kafkaTemplate = new KafkaTemplate<>(defaultKafkaProducerFactory); kafkaTemplate.setObservationEnabled(true); kafkaTemplate.setObservationConvention(new KafkaTemplateObservation.DefaultKafkaTemplateObservationConvention()); return kafkaTemplate; } }
方案2:直接关闭主题检查(快速绕过)
如果你100%确认目标主题一定存在,也可以直接关闭Spring Kafka的主题检查行为,这样就不会触发AdminClient了,省去配置它的麻烦:
// 在创建KafkaTemplate后添加这行代码 kafkaTemplate.setCheckTopicExists(false);
或者直接在生产者配置里添加:
properties.put(ProducerConfig.CHECK_TOPIC_EXISTS_CONFIG, false);
⚠️ 注意:这种方法只是绕过问题,如果你后续需要使用AdminClient做其他集群操作,还是得回到方案1配置它的SSL。
三、补充:为什么你的YAML配置没生效?
顺便提一句,你之前尝试的YAML配置没生效,是因为Spring Kafka的AdminClient默认读取spring.kafka.admin.*前缀的配置,而不是spring.kafka.*。如果想用YAML配置AdminClient的SSL,应该写成:
spring: kafka: admin: properties: security.protocol: SSL ssl: key-password: abc keystore-location: 'file:///Users/me/keystore.p12' keystore-password: abc truststore-location: 'file:///Users/me/truststore.p12' truststore-password: abc
不过既然你想用SslBundle的编程方式,方案1会更贴合你的需求。
内容来源于stack exchange

