Java环境Kafka明文模式创建生产者的Properties配置及SSL代码处理咨询
Kafka明文模式(未启用SSL)Java生产者配置说明
一、必填配置字段
你当前场景是明文传输,结合代码里已经配置了SASL账号密码认证的逻辑,需要保留的有效配置如下:
bootstrap.servers:Kafka集群节点地址,多个节点用英文逗号分隔,格式为ip1:port1,ip2:port2key.serializer、value.serializer:消息键、值的序列化实现类,你当前使用的StringSerializer符合字符串类型消息的需求security.protocol:明文+SASL认证场景下固定配置为SASL_PLAINTEXT;如果连SASL认证也没开启,配置为PLAINTEXT即可sasl.jaas.config:仅开启SASL认证时需要配置,内容就是你当前拼接的PlainLoginModule账号密码配置,逻辑正确
二、现有代码调整建议
- 冗余SSL配置直接删除
你代码里的ssl.cipher.suites、ssl.enabled.protocols都是SSL加密场景才需要的配置,未启用SSL的场景下完全无效,建议直接删除,避免冗余配置引发不必要的连接异常。 - 修复bootstrap.servers配置逻辑
你现有代码在循环里逐个给bootstrap.servers赋值、每次循环都新建销毁生产者,存在两个问题:
- 每次只配置单个节点,若该节点故障会直接导致连接失败
- Kafka生产者是重量级对象,频繁创建销毁会严重损耗性能
建议提前把所有节点地址拼接为逗号分隔的字符串,单次赋值后复用同一个生产者实例。
三、修正后参考代码
Properties pro = new Properties(); // 拼接所有节点地址 String bootstrapServers = String.join(",", ipAndPorts); pro.setProperty("bootstrap.servers", bootstrapServers); // 仅保留需要的配置,删除SSL相关项 String sasl = String.format("org.apache.kafka.common.security.plain.PlainLoginModule required username=\"%s\" password=\"%s\";", userName, password); pro.setProperty("sasl.jaas.config", sasl); // 明文+SASL认证固定为SASL_PLAINTEXT,无SASL则改为PLAINTEXT pro.setProperty("security.protocol", "SASL_PLAINTEXT"); pro.setProperty("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); pro.setProperty("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); // 复用生产者实例,不要循环创建 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); }
内容的提问来源于stack exchange,提问作者pi_ka_qiu_1999
相关产品推荐
相关产品推荐

