Kafka Streams配置SpecificAvroSerDe时缺失schema.registry.url报错求助
我为Kafka Streams应用配置了如下参数:
Properties config = new Properties(); config.put(StreamsConfig.APPLICATION_ID_CONFIG,this.applicaionId); config.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG,svrConfig.getBootstrapServers()); config.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); // we disable the cache to demonstrate all the "steps" involved in the transformation - not recommended in prod config.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, svrConfig.getCacheMaxBytesBufferingConfig()); // Exactly once processing!! config.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE); config.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG,SpecificAvroSerDe.class); config.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG,SpecificAvroSerDe.class); config.put(AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG,"http://localhost:8081");
运行时出现如下错误:
Exception in thread "main" io.confluent.common.config.ConfigException: Missing required configuration "schema.registry.url" which has no default value. at io.confluent.common.config.ConfigDef.parse(ConfigDef.java:243) at io.confluent.common.config.AbstractConfig.<init>(AbstractConfig.java:78) at io.confluent.kafka.serializers.AbstractKafkaAvroSerDeConfig.<init>(AbstractKafkaAvroSerDeConfig.java:100) at io.confluent.kafka.serializers.KafkaAvroSerializerConfig.<init>(KafkaAvroSerializerConfig.java:32) at io.confluent.kafka.serializers.KafkaAvroSerializer.configure(KafkaAvroSerializer.java:48) at io.confluent.kafka.streams.serdes.avro.SpecificAvroSerializer.configure(SpecificAvroSerializer.java:58) at io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde.configure(SpecificAvroSerde.java:107)
我尝试将配置行替换为:
config.put("schema.registry.url","http://localhost:8081");
但仍出现相同错误。我是按照官方文档配置的应用,请问有解决建议吗?
解决建议
别着急,这个问题通常是因为SerDe实例没有正确获取到配置导致的,这里有几个可行的排查和解决方向:
显式为SerDe实例单独配置Schema Registry URL
有时候依赖全局默认SerDe配置可能会因为初始化顺序、配置传递失效等问题出问题,你可以手动创建SpecificAvroSerde实例并单独传入配置,确保它能直接拿到需要的参数:SpecificAvroSerde<YourKeyType> keySerde = new SpecificAvroSerde<>(); SpecificAvroSerde<YourValueType> valueSerde = new SpecificAvroSerde<>(); Map<String, String> serdeConfig = new HashMap<>(); serdeConfig.put(AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://localhost:8081"); keySerde.configure(serdeConfig, true); // 第二个参数true表示这是用于Key的SerDe valueSerde.configure(serdeConfig, false); // false表示这是用于Value的SerDe // 在构建流的时候明确指定这些SerDe KStream<String, YourValueType> stream = builder.stream("input-topic", Consumed.with(keySerde, valueSerde));检查配置是否被意外覆盖
可以在初始化完config对象后,打印所有配置项确认schema.registry.url是否存在:config.forEach((k, v) -> System.out.println(k + ": " + v));如果输出里没有这个配置项,说明你的代码逻辑中存在覆盖或清除该配置的操作,需要排查相关代码。
确认依赖版本兼容性
确保你的Confluent相关依赖(kafka-streams、kafka-avro-serializer、schema-registry-client等)版本完全一致。版本不匹配可能会导致配置键的识别出现偏差,比如不同版本中AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG对应的字符串可能不一致,导致配置无法被正确读取。验证Schema Registry服务可用性
虽然当前错误是配置缺失,但可以顺便确认http://localhost:8081地址的Schema Registry是否正常运行,比如用curl测试:curl http://localhost:8081/subjects服务不可达会导致后续的其他错误,但优先解决配置传递的问题。
内容的提问来源于stack exchange,提问作者Benny Chan

