对接第三方Kafka需配置哪些信息?SSL与PLAIN模式如何适配?
对接第三方Kafka配置说明
一、核心必填配置项
通用基础配置(所有模式都需要)
- 集群接入地址:配置项为
bootstrap.servers,填写第三方提供的Kafka节点IP+端口,多个节点用逗号分隔 - 序列化规则:生产者需配置
key.serializer和value.serializer,字符串类型消息一般用org.apache.kafka.common.serialization.StringSerializer,消费者对应配置同类型的反序列化器 - 安全协议:配置项为
security.protocol,按照第三方要求选择协议类型,你当前用到的是SASL_PLAINTEXT(你提到的PLAIN模式)和SASL_SSL两种
SASL认证配置(你当前用的PLAIN认证模式需要)
配置项为sasl.jaas.config,填写第三方分配的用户名和密码,你现有代码的写法是正确的,额外建议显式配置sasl.mechanism为PLAIN,避免默认值和服务端不匹配。
二、keystore与truststore证书配置规则
两类证书不存在二选一的情况,使用场景差异如下:
- truststore:存储你信任的服务端证书,作用是客户端校验服务端身份合法性,所有使用SSL/TLS加密的场景都必须配置,不管是单向还是双向认证
- keystore:存储客户端自身的身份证书,作用是让服务端校验客户端的合法性,只有双向SSL认证场景需要配置
简单总结规则: - 用
SASL_PLAINTEXT等非SSL协议:两类证书都不需要配置 - 用
SASL_SSL协议且是单向认证(绝大多数第三方Kafka的默认模式):只需要配置truststore - 用
SASL_SSL协议且是双向认证:需要同时配置truststore和keystore
证书的格式、路径、密码都需要和第三方Kafka运维方确认后填写。
三、你现有双模式适配代码的优化建议
你当前代码没有做协议分支判断,不管用哪种模式都配置了SSL相关参数,虽然SASL_PLAIN模式下这些参数会被忽略,但还是建议做分支判断减少冗余配置,参考逻辑如下:
Properties pro = new Properties(); // 公共配置 pro.setProperty("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); pro.setProperty("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); pro.setProperty("bootstrap.servers", String.join(",", ipAndPorts)); // 建议一次性拼接所有节点,不要循环覆盖 // SASL公共配置 String sasl = String.format("org.apache.kafka.common.security.plain.PlainLoginModule required username=\"%s\" password=\"%s\";", userName, password); pro.setProperty("sasl.jaas.config", sasl); pro.setProperty("sasl.mechanism", "PLAIN"); pro.setProperty("security.protocol", protocols); // 按协议分支配置 if ("SASL_SSL".equals(protocols)) { pro.setProperty("ssl.cipher.suites", suites); pro.setProperty("ssl.enabled.protocols", authorizedEncryptionMode); // 配置truststore pro.setProperty("ssl.truststore.location", "你的truststore文件路径"); pro.setProperty("ssl.truststore.password", "truststore密码"); // 如果是双向认证,额外加keystore配置 // pro.setProperty("ssl.keystore.location", "你的keystore文件路径"); // pro.setProperty("ssl.keystore.password", "keystore密码"); // pro.setProperty("ssl.key.password", "私钥密码"); } // 消息发送逻辑统一写在分支外 try(Producer<String, String> producer = new KafkaProducer<>(pro)) { producer.send(new ProducerRecord<>(topic, null, System.currentTimeMillis(), alarmModel.getString("alarmId"), new Record(alarmModel).toString())); } catch (Exception e) { log.info("pmRuleDeliverBySouthKafka failed....", e); }
另外你原有代码里循环遍历ipAndPorts重复覆盖bootstrap.servers的写法是错误的,bootstrap.servers支持填写多个节点地址,一次性用逗号拼接即可,不需要循环创建生产者。
内容的提问来源于stack exchange,提问作者pi_ka_qiu_1999
相关产品推荐
相关产品推荐

