JMeter通过JSR223发消息至CloudKarafka时SASL认证失败求助
问题:CloudKarafka Kafka发送消息报错Failed to configure SaslClientAuthenticator
Failed to configure SaslClientAuthenticator
Error connecting to node dory.srvs.cloudkafka.com:9094 (id: -1 rack: null)
用户提供的原始代码
JSR223脚本
System.setProperty("java.security.auth.login.config" , "C:/!work/apache-jmeter-5.4.1/bin/kafka-jaas.conf"); //System.setProperty("java.security.auth.login.config" , "kafka-jaas.conf"); import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.*; import java.util.Properties; Properties properties = new Properties(); properties.put("bootstrap.servers", "dory.srvs.cloudkafka.com:9094"); properties.put("security.protocol", "SASL_PLAINTEXT"); properties.put("sasl.mechanism", "PLAINTEXT"); properties.put("acks", "1"); properties.put("retries", 1); properties.put("batch.size", 16384); properties.put("linger.ms", 0); properties.put("buffer.memory", 33554432); properties.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); properties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); properties.put("compression.type", "none"); properties.put("send.buffer.bytes", 131072); properties.put("receive.buffer.bytes", 32768); properties.put("sasl.kerberos.service.name", "kafka"); properties.put("ssl.keystore.type", "JKS"); properties.put("ssl.truststore.type", "JKS"); Producer<String, String> producer = new KafkaProducer<>(properties); producer.send(new ProducerRecord<String, String>("lzejnhew-qwe", "{\n\t\"messageId\":{{SEQUENCE(\"messageId\", 1, 1)}},\n\t\"messageBody\":\"{{RANDOM_ALPHA_NUMERIC(\"abcedefghijklmnopqrwxyzABCDEFGHIJKLMNOPQRWXYZ\", 100)}}\",\n\t\"messageCategory\":\"{{RANDOM_STRING(\"Finance\", \"Insurance\", \"Healthcare\", \"Shares\")}}\",\n\t\"messageStatus\":\"{{RANDOM_STRING(\"Accepted\",\"Pending\",\"Processing\",\"Rejected\")}}\",\n\t\"messageTime\":{{TIMESTAMP()}}\n}")); producer.close();
kafka-jaas.conf配置
KafkaClient { org.apache.kafka.common.security.plain.PlainLoginModule required serviceName="lzejnhew-qwe" username="lzejnhew" password="***"; };
错误排查与修复方案
- Sasl机制配置错误:将
sasl.mechanism值从PLAINTEXT改为PLAIN(CloudKarafka采用SASL PLAIN认证机制,PLAINTEXT是协议类型而非认证机制) - 移除多余的SSL配置:当前使用
SASL_PLAINTEXT协议,无需配置ssl.keystore.type和ssl.truststore.type,直接删除这两项 - 移除Kerberos相关配置:
sasl.kerberos.service.name是Kerberos认证专用参数,SASL PLAIN认证不需要,删除该配置 - JAAS配置清理:CloudKarafka的JAAS配置不需要
serviceName参数,删除该行
修正后的代码与配置
修正后的JSR223脚本
System.setProperty("java.security.auth.login.config" , "C:/!work/apache-jmeter-5.4.1/bin/kafka-jaas.conf"); import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.*; import java.util.Properties; Properties properties = new Properties(); properties.put("bootstrap.servers", "dory.srvs.cloudkafka.com:9094"); properties.put("security.protocol", "SASL_PLAINTEXT"); properties.put("sasl.mechanism", "PLAIN"); // 修正认证机制 properties.put("acks", "1"); properties.put("retries", 1); properties.put("batch.size", 16384); properties.put("linger.ms", 0); properties.put("buffer.memory", 33554432); properties.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); properties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); properties.put("compression.type", "none"); properties.put("send.buffer.bytes", 131072); properties.put("receive.buffer.bytes", 32768); Producer<String, String> producer = new KafkaProducer<>(properties); producer.send(new ProducerRecord<String, String>("lzejnhew-qwe", "{\n\t\"messageId\":{{SEQUENCE(\"messageId\", 1, 1)}},\n\t\"messageBody\":\"{{RANDOM_ALPHA_NUMERIC(\"abcedefghijklmnopqrwxyzABCDEFGHIJKLMNOPQRWXYZ\", 100)}}\",\n\t\"messageCategory\":\"{{RANDOM_STRING(\"Finance\", \"Insurance\", \"Healthcare\", \"Shares\")}}\",\n\t\"messageStatus\":\"{{RANDOM_STRING(\"Accepted\",\"Pending\",\"Processing\",\"Rejected\")}}\",\n\t\"messageTime\":{{TIMESTAMP()}}\n}")); producer.close();
修正后的kafka-jaas.conf
KafkaClient { org.apache.kafka.common.security.plain.PlainLoginModule required username="lzejnhew" password="***"; };
内容的提问来源于stack exchange,提问作者yndingo
相关产品推荐
相关产品推荐

