【问题标题】:Kafka Stream with Avro in JAVA , schema.registry.url" which has no default valueKafka Stream with Avro in JAVA , schema.registry.url" 没有默认值
【发布时间】:2018-10-25 04:51:54
【问题描述】:

我的 Kafka Stream 应用程序有以下配置

    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(AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG,"http://localhost:8081");

config.put("schema.registry.url","http://localhost:8081");

但同样的错误

在准备我的 Stream 应用程序时,我已按照 this url 的说明进行操作。

有什么建议吗?

【问题讨论】:

  • 从链接的例子中,你的Collections.singletonMap在哪里?

标签: java apache-kafka avro apache-kafka-streams confluent-schema-registry


【解决方案1】:

如果您有 Avro 格式的键和值,那么以下几行应该可以为您解决问题,

config.put("key.converter.schema.registry.url", "http://localhost:8081");  
config.put("value.converter.schema.registry.url", "http://localhost:8081");

如果这似乎不起作用,您可以明确覆盖 Serdes。例如,如果您有 Avro 密钥:

final Map<String, String> serdeConfig = Collections.singletonMap("schema.registry.url",
                                                                 "http://localhost:8081");
final Serde<GenericRecord> keyGenericAvroSerde = new GenericAvroSerde();
keyGenericAvroSerde.configure(serdeConfig, true); // true for record keys
final Serde<GenericRecord> valueGenericAvroSerde = new GenericAvroSerde();
valueGenericAvroSerde.configure(serdeConfig, false); // false for record values

StreamsBuilder builder = new StreamsBuilder();
KStream<GenericRecord, GenericRecord> textLines =
builder.stream(keyGenericAvroSerde, valueGenericAvroSerde, "my-avro-topic");
// Do whatever you like

【讨论】:

  • 非常感谢 Giorgos。结果我添加了行“final Map serdeConfig = Collections.singletonMap("schema.registry.url", "localhost:8081");" 一切正常。我以为这行只是一个可选设置要运行,但结果似乎这是必须的~
  • 跟进:来自StreamConfig 的配置仅传递给内部创建的所有 Serdes KafkaStreams。对于您的情况,您在代码中创建 Serdes,因此您有责任提供正确的配置。参照。 docs.confluent.io/current/streams/developer-guide/…
  • 有谁知道如何在 Spring Boot 应用程序中进行配置?也许直接通过application.yml
  • @Naveen 在下面查看我的答案stackoverflow.com/a/68574923/1128216
  • 为什么要为 Kafka Streams 使用 Kafka Connect 转换器属性?
【解决方案2】:

如果您使用的是application.yml,您可以将属性设置为:

spring:
  kafka:
    properties:
      schema.registry.url: your-schema-registy-url
    consumer:
      auto-offset-reset: latest
      group-id: simple-consumer

我在confluent blog的教程中找到了它

【讨论】:

  • 它记录了一个警告The configuration 'schema.registry.url' was supplied but isn't a known config. 有一个answer 说通过更改属性的位置(将它放在producer 属性之后)该警告消失了该建议不适用于我。
  • 警告是因为 AdminClient 配置不知道该属性。您可以将其移至consumer.properties 下以使其静音
【解决方案3】:

显然答案已经过时了。
使用此代码 sn-p 对我有用(这里的示例只是将一个主题传递给另一个主题,没有任何更改)。

    final Map<String, String> serdeConfig = Collections.singletonMap(AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG,"http://localhost:8081");
    final Serde<GenericRecord> keyGenericAvroSerde = new GenericAvroSerde();
    keyGenericAvroSerde.configure(serdeConfig, true); // true for record keys
    final Serde<GenericRecord> valueGenericAvroSerde = new GenericAvroSerde();
    valueGenericAvroSerde.configure(serdeConfig, false); // false for record values

    StreamsBuilder builder = new StreamsBuilder();

    builder.stream("my-avro-topic",Consumed.with(keyGenericAvroSerde, valueGenericAvroSerde))
                .peek((k, v) -> log.info("Processed a new record")).to(outputTopic,Produced.with(keyGenericAvroSerde, valueGenericAvroSerde));

【讨论】:

  • 这似乎重复了已接受答案的下半部分。如果您说需要 Consumed / Produced 对象,最好编辑/评论帖子
猜你喜欢
  • 1970-01-01
  • 2014-05-21
  • 2018-01-20
  • 2020-11-11
  • 2019-02-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多