【问题标题】:Kafka Streams StreamsException when I work with multiple streams当我使用多个流时,Kafka Streams StreamsException
【发布时间】:2020-05-21 09:38:15
【问题描述】:

我正在与 Kafka Streams 和 Kotlin 合作开发一项服务,该服务具有三个主题的流。第一个具有 Avro 值,另外两个具有 String 值。

在我的properties 文件中,我将SpecificAvroSerde 作为默认值Serde,然后我使用Consumed.with(Serdes.String(), Serdes.String()) 来使用字符串值。

    val topicOneStream = streamsBuilder.stream<String, AvroObject>(topicOne)
        .peek { k, _ -> logger.info("Received message with key: $k") }
        .flatMapValues { v -> listOf(v) }.groupByKey().reduce { v1, _ -> v1 }

    val topicTwoStream = streamsBuilder
        .stream<String, String>(topicTwo, Consumed.with(Serdes.String(), Serdes.String()))
        .peek { k, _ -> logger.info("Received message with key: $k") }
        .flatMapValues { v -> listOf(v) }.groupByKey().reduce { v1, _ -> v1 }

    val topicThreeStream = streamsBuilder.stream<String, String>(topicThree, Consumed.with(Serdes.String(), Serdes.String()))
        .peek { k, _ -> logger.info("Received message with key: $k") }
        .mapValues { v -> objectMapper.readValue(v, AdviceCreated::class.java) }
        .flatMapValues { v -> listOf(v) }.groupByKey().reduce { v1, _ -> v1 }

当我将以下流的配置作为值的默认值时,我看到 Avro 流(第一个流)工作正常并使用了我在该主题上发布的内容。但是当我使用相同的配置发布到字符串值流时出现异常。

default.value.serde: io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde

这是发布到 topicTwo 和 topicThree 的异常:

org.apache.kafka.streams.errors.StreamsException: A serializer (io.confluent.kafka.streams.serdes.avro.SpecificAvroSerializer) is not compatible to the actual value type (value type: java.lang.String). Change the default Serdes in StreamConfig or provide correct Serdes via method parameters.

PS。它必须是同一服务中的三个流,因为稍后会有一个连接。

【问题讨论】:

  • 您的代码看起来正确。您可以使用TopologyTestDriver 重现该问题吗?
  • 我已经使用 TopologyTestDriver 为它编写了一个测试,它会显示同样的错误。如果SpecificAvroSerde 是默认的字符串流将不会被消费,但如果StringSerde 是默认的那么avro 流将不会被消费!

标签: kotlin apache-kafka apache-kafka-streams spring-kafka


【解决方案1】:

感谢一位朋友 (Mario Boikov),当 Kafka 进行分组以生成新的 KTable 时,问题就出现了。它不知道为分组采用哪个序列化程序,因此它采用默认的 Serde 作为值,在我的情况下是 SpecificAvroSerde

通过为 groupByKey 提供分组所需的序列化程序解决了这个问题:

val topicTwoStream = streamsBuilder
    .stream<String, String>(topicTwo, Consumed.with(Serdes.String(), Serdes.String()))
    .peek { k, _ -> logger.info("Received message with key: $k") }
    .flatMapValues { v -> listOf(v) }.groupByKey(Grouped.with(Serdes.String(), Serdes.String())).reduce { v1, _ -> v1 }

val topicThreeStream = streamsBuilder.stream<String, String>(topicThree, Consumed.with(Serdes.String(), Serdes.String()))
    .peek { k, _ -> logger.info("Received message with key: $k") }
    .flatMapValues { v -> listOf(v) }.groupByKey(Grouped.with(Serdes.String(), Serdes.String())).reduce { v1, _ -> v1 }
    .mapValues { v -> objectMapper.readValue(v, AdviceCreated::class.java) }

干杯?

【讨论】:

    猜你喜欢
    • 2019-01-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-06-20
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多