【问题标题】:Store sum in KTable<String, Integer>将总和存储在 KTable<String, Integer>
【发布时间】:2019-06-26 23:11:46
【问题描述】:

我正在尝试计算某个主题的信封数量。交易采用 avro 格式。我使用this example 作为参考。

final StreamsBuilder streamsBuilder = new StreamsBuilder();
final KStream<String, Transaction> transactionKStream = streamsBuilder.stream(INPUT_TOPIC);

final KStream<String, Integer> envelopes = transactionKStream.filter((k, v) -> v.getProduct().toString()
    .matches("C4|C5"))
    .map((k, v) -> KeyValue.pair("1", v.getAmount()));

final KTable<String, Integer> amount = envelopes
    .groupByKey()
    .reduce((v1, v2) -> v1 + v2);

我想将总和存储在 KTable 但是当我将数据发送到输入主题时,消费者会崩溃

A serializer (key: org.apache.kafka.common.serialization.StringSerializer / value: io.confluent.kafka.streams.serdes.avro.GenericAvroSerializer) is not compatible to the actual key or value type (key type: java.lang.String / value type: java.lang.Integer). Change the default Serdes in StreamConfig or provide correct Serdes via method parameters.

当 KTable 被注释掉时,它运行良好。但不要合计金额。

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    groupByKey() 使用默认的序列化器:

    groupByKey()

    按当前键将记录分组到一个 KGroupedStream 同时保留原始值和默认值 序列化器和反序列化器。

    您必须使用groupByKey(Serialized&lt;K,V&gt; serialized)groupByKey(Grouped&lt;K,V&gt; grouped)

    以下应该可以解决问题:

    final KTable<String, Integer> amount = envelopes
        .groupByKey(Serialized.with(Serdes.String(), Serdes.Integer()))
        .reduce((v1, v2) -> v1 + v2);
    

    【讨论】:

      猜你喜欢
      • 2017-03-21
      • 2021-02-21
      • 1970-01-01
      • 2019-12-06
      • 1970-01-01
      • 2019-05-05
      • 1970-01-01
      • 1970-01-01
      • 2015-02-22
      相关资源
      最近更新 更多