Kafka Producer抛出“key.serializer”配置缺失异常求助
排查Kafka Producer "key.serializer"配置缺失异常
尝试过构造函数(部分)设置属性、setProperties、put方法、ProducerConfig及文本值配置等多种方式,始终遇到Kafka Producer抛出的ConfigException: Missing required configuration "key.serializer" which has no default value异常,以下是相关代码及异常信息:
原代码
public class KafkaProducer { private <T> void produce(T data, String topic) { Gson gson = new GsonBuilder() .setPrettyPrinting() .registerTypeAdapter(LocalDate.class, new LocalDateAdapter()) .create(); String jsonString = gson.toJson(data); Properties kafkaProperties = new Properties(); try(Producer<String, String> producer = new KafkaProducer<>(kafkaProperties)) { kafkaProperties.setProperty(CLIENT_ID_CONFIG, MainProperties.get().kafkaProducerProperties.getClientId()); kafkaProperties.setProperty(BOOTSTRAP_SERVERS_CONFIG, MainProperties.get().kafkaProducerProperties.getUrl()); kafkaProperties.setProperty(KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); kafkaProperties.setProperty(VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); createTopic(topic, kafkaProperties); producer.send(new ProducerRecord<>(topic, jsonString)); } } public void produceDataType1(KafkaType1Message kafkaType1Values) { produce(kafkaType1Values, MainProperties.get().kafkaProducerProperties.getType1Topic()); } public void produceDataType2(KafkaType2Message kafkaType2Values) { produce(kafkaType2Values, MainProperties.get().kafkaProducerProperties.getType2ValuesTopic()); } public KafkaProducerProperties(Source source) { super(source); this.url = value("url",""); this.clientId = value("clientId", "TestProducer"); this.type1ValuesTopic = value("type1", "type1_topic"); this.type2ValuesTopic = value("type2", "type1_topic"); } public static Factory<KafkaProducerProperties> factory() { return KafkaProducerProperties::new; } }
异常信息
org.apache.kafka.common.config.ConfigException: Missing required configuration "key.serializer" which has no default value. at org.apache.kafka.common.config.ConfigDef.parseValue(ConfigDef.java:493) at org.apache.kafka.common.config.ConfigDef.parse(ConfigDef.java:483) at org.apache.kafka.common.config.AbstractConfig.<init>(AbstractConfig.java:113) at org.apache.kafka.common.config.AbstractConfig.<init>(AbstractConfig.java:133) at org.apache.kafka.clients.producer.ProducerConfig.<init>(ProducerConfig.java:553) at org.apache.kafka.clients.producer.KafkaProducer.<init>(KafkaProducer.java:289) at org.apache.kafka.clients.producer.KafkaProducer.<init>(KafkaProducer.java:316) at org.apache.kafka.clients.producer.KafkaProducer.<init>(KafkaProducer.java:301)
问题原因
核心错误是配置属性的设置时机错误:代码中先创建了KafkaProducer实例,之后才给kafkaProperties设置key.serializer等必要配置。KafkaProducer在构造时会立即读取传入的Properties配置,后续对Properties的修改不会同步到已创建的Producer实例中,导致构造时配置缺失,触发异常。
解决方法
将所有Properties的配置代码移到创建KafkaProducer实例之前,确保初始化Producer时所有必要配置已就绪。修改后的produce方法代码如下:
private <T> void produce(T data, String topic) { Gson gson = new GsonBuilder() .setPrettyPrinting() .registerTypeAdapter(LocalDate.class, new LocalDateAdapter()) .create(); String jsonString = gson.toJson(data); Properties kafkaProperties = new Properties(); // 先设置所有必要配置 kafkaProperties.setProperty(CLIENT_ID_CONFIG, MainProperties.get().kafkaProducerProperties.getClientId()); kafkaProperties.setProperty(BOOTSTRAP_SERVERS_CONFIG, MainProperties.get().kafkaProducerProperties.getUrl()); kafkaProperties.setProperty(KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); kafkaProperties.setProperty(VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); // 再用配置好的Properties创建Producer实例 try(Producer<String, String> producer = new KafkaProducer<>(kafkaProperties)) { createTopic(topic, kafkaProperties); producer.send(new ProducerRecord<>(topic, jsonString)); } }
内容的提问来源于stack exchange,提问作者One String Samurai
相关产品推荐
相关产品推荐

