【问题标题】:Using kafka streams trying to write messages from input topic to output topic使用 kafka 流尝试将消息从输入主题写入输出主题
【发布时间】:2020-11-05 04:36:21
【问题描述】:

下面是我尝试将数据从一个主题写入另一个主题的代码逻辑。

//Computational logic
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> kStream = builder.stream(topicName_In);
kStream.foreach((k, v) -> System.out.println("Key = " + k + " Value = " + v));
//kStream.peek((k, v) -> System.out.println("Key = " + k + " Value = " + v));
kStream.to(topicName_Out);
Topology topology = builder.build();

输入主题数据格式: 简单消息-1

错误

Exception in thread "HelloStreams-564343a1-1709-4bae-8fe5-514b37eee595-StreamThread-1" org.apache.kafka.streams.errors.StreamsException: Deserialization exception handler is set to fail upon a deserialization error. If you would rather have the streaming pipeline continue after a deserialization error, please set the default.deserialization.exception.handler appropriately.
    at org.apache.kafka.streams.processor.internals.RecordDeserializer.deserialize(RecordDeserializer.java:80)
    at org.apache.kafka.streams.processor.internals.RecordQueue.updateHead(RecordQueue.java:175)
    at org.apache.kafka.streams.processor.internals.RecordQueue.addRawRecords(RecordQueue.java:112)
    at org.apache.kafka.streams.processor.internals.PartitionGroup.addRawRecords(PartitionGroup.java:162)
    at org.apache.kafka.streams.processor.internals.StreamTask.addRecords(StreamTask.java:765)
    at org.apache.kafka.streams.processor.internals.StreamThread.addRecordsToTasks(StreamThread.java:943)
    at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:764)
    at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:697)
    at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:670)
Caused by: org.apache.kafka.common.errors.SerializationException: Size of data received by IntegerDeserializer is not 4

【问题讨论】:

  • 您的实际问题是什么?您使用什么资源来创建从一个主题读取并写入另一个主题的 KafkaStream 示例?这似乎是使用 kafkaStreams 可以做的最基本的事情,并且想知道为什么它不能在网络上搜索到。如果您提出问题,也许它会澄清......
  • 我是 Kafka 新手,尝试从一个主题读取数据并尝试在另一个主题中写入
  • 也许尝试从Kafka documentation中提到的演示应用开始
  • 作为建议,当您在 Stackoverflow 上发布内容时,不要忘记提出一个真正的问题。当您只是发布错误堆栈跟踪时,有些人可能不明白您实际要解决的问题。

标签: apache-kafka kafka-consumer-api apache-kafka-streams kafka-producer-api


【解决方案1】:

鉴于这一行

KStream<String, String> kStream = builder.stream(topicName_In);

我假设你的输入数据是字符串。

错误信息说

Size of data received by IntegerDeserializer is not 4

这表明您确实配置了IntegerSerde(用于键和/或值),但您需要配置StringSerde 才能使其工作。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-22
    • 1970-01-01
    • 2019-06-04
    • 1970-01-01
    • 2017-09-23
    相关资源
    最近更新 更多