java.util.Properties实例初始化配置未复制导致Kafka消费者异常
问题原因
- Java
java.util.Properties的构造函数new Properties(Properties defaults)并非直接复制传入对象的所有属性,而是将传入的Properties设置为新对象的默认属性集合。 - Kafka的
KafkaConsumer在加载配置时,只会读取Properties对象的主属性集合,不会访问默认属性集合中的配置。这就导致你传入的key.deserializer等配置无法被Kafka读取,最终抛出配置为空的异常。
正确的惰性初始化实现方式
如果要实现props的惰性初始化,同时保证配置能被Kafka正确识别,可采用以下写法:
Scala 正确实现
import org.apache.kafka.clients.consumer.{ConsumerConfig, KafkaConsumer} import java.util.Properties object KafkaConsumerWithAuth extends App { // lazy val确保props仅在首次被使用时才初始化 lazy val props: Properties = { val p = new Properties() p.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "redacted") p.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer") p.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer") p.setProperty("security.protocol", "SASL_SSL") p.setProperty("sasl.mechanism", "PLAIN") p.setProperty("sasl.jaas.config", "redacted") p.setProperty("ssl.endpoint.identification.algorithm", "") p.setProperty("ssl.truststore.location", "redacted.jks") p.setProperty("ssl.truststore.password", "redacted") p } lazy val consumer = new KafkaConsumer(props) println(consumer.listTopics()) }
Java 正确复制Properties的写法
如果需要在Java中复制Properties的所有属性,不能依赖带defaults的构造函数,需手动调用putAll将属性复制到新对象的主集合:
import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; import java.util.Properties; public class KafkaConsumerJava { public static void main(String[] args) { Properties props = new Properties(); props.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "redacted"); props.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); props.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); props.setProperty("security.protocol", "SASL_SSL"); props.setProperty("sasl.mechanism", "PLAIN"); props.setProperty("sasl.jaas.config", "redacted"); props.setProperty("ssl.endpoint.identification.algorithm", ""); props.setProperty("ssl.truststore.location", "redacted"); props.setProperty("ssl.truststore.password", "redacted"); // 手动复制所有属性到新对象的主集合 Properties props2 = new Properties(); props2.putAll(props); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props2); System.out.println(consumer.listTopics()); } }
内容的提问来源于stack exchange,提问作者Pavel Orekhov
相关产品推荐
相关产品推荐

