You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 10:37:54