【发布时间】: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