【发布时间】:2022-01-12 01:10:48
【问题描述】:
生产者属性
spring.kafka.producer.bootstrap-servers=127.0.0.1:9092
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer
消费者属性
spring.kafka.consumer.bootstrap-servers=127.0.0.1:9092
spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer
spring.kafka.consumer.properties.spring.json.trusted.packages=*
spring.kafka.consumer.group-id=user-group
server.port=8085
消费者服务
@Service
public class UserConsumerService {
@KafkaListener(topics = { "user-topic" })
public void consumerUserData(User user) {
System.out.println("Users Age Is: " + user.getAge() + " Fav Genre " + user.getFavGenre());
}
}
生产者服务
@Service
public class UserProducerService {
@Autowired
private KafkaTemplate<String, User> kafkaTemplate;
public void sendUserData(User user) {
kafkaTemplate.send("user-topic", user.getName(), user);
}
}
创建主题的生产者配置
@Configuration public class KafkaConfig {
@Bean
public NewTopic topicOrder() {
return TopicBuilder.name("user-topic").partitions(2).replicas(1).build();
}
}
生产者运行良好,但消费者给出了类似的错误
2021-12-06 21:45:50.299 ERROR 4936 --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer : Consumer exception java.lang.IllegalStateException: This error handler cannot process 'SerializationException's directly; please consider configuring an'ErrorHandlingDeserializer' 在值和/或键反序列化器中 org.springframework.kafka.listener.DefaultErrorHandler.handleOtherException(DefaultErrorHandler.java:149) ~[spring-kafka-2.8.0.jar:2.8.0] DefaultErrorHandler.java:149 at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.handleConsumerException(KafkaMessageListenerContainer.java:1760) ~[spring-kafka-2.8.0.jar:2.8.0] KafkaMessageListenerContainer.java:1760 在 org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1283) ~[spring-kafka-2.8.0.jar:2.8.0] KafkaMessageListenerContainer.java:1283 在 java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) ~[na:na] Executors.java:539 在 java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) ~[na:na] FutureTask.java:264 at java.base/java.lang.Thread.run(Thread.java:833) ~[na:na] Thread.java:833 原因: org.apache.kafka.common.errors.RecordDeserializationException:错误 在偏移量 1 处反序列化分区 user-topic-0 的键/值。如果 有需要,请追查过去的记录继续消费。在 org.apache.kafka.clients.consumer.internals.Fetcher.parseRecord(Fetcher.java:1429) ~[kafka-clients-3.0.0.jar:na] Fetcher.java:1429 at org.apache.kafka.clients.consumer.internals.Fetcher.access$3400(Fetcher.java:134) ~[kafka-clients-3.0.0.jar:na] Fetcher.java:134 at org.apache.kafka.clients.consumer.internals.Fetcher$CompletedFetch.fetchRecords(Fetcher.java:1652) ~[kafka-clients-3.0.0.jar:na] Fetcher.java:1652 at org.apache.kafka.clients.consumer.internals.Fetcher$CompletedFetch.access$1800(Fetcher.java:1488) ~[kafka-clients-3.0.0.jar:na] Fetcher.java:1488 at org.apache.kafka.clients.consumer.internals.Fetcher.fetchRecords(Fetcher.java:721) ~[kafka-clients-3.0.0.jar:na] Fetcher.java:721 at org.apache.kafka.clients.consumer.internals.Fetcher.fetchedRecords(Fetcher.java:672) ~[kafka-clients-3.0.0.jar:na] Fetcher.java:672 at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1277) ~[kafka-clients-3.0.0.jar:na] KafkaConsumer.java:1277 at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1238) ~[kafka-clients-3.0.0.jar:na] KafkaConsumer.java:1238 在 org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1211) ~[kafka-clients-3.0.0.jar:na] KafkaConsumer.java:1211 at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollConsumer(KafkaMessageListenerContainer.java:1507) ~[spring-kafka-2.8.0.jar:2.8.0] KafkaMessageListenerContainer.java:1507 在 org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doPoll(KafkaMessageListenerContainer.java:1497) ~[spring-kafka-2.8.0.jar:2.8.0] KafkaMessageListenerContainer.java:1497 在 org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1325) ~[spring-kafka-2.8.0.jar:2.8.0] KafkaMessage
如果您能提供帮助,我会很高兴,因为我是 kafka 的新手并试图找出为什么会出现此错误
【问题讨论】:
-
您确定您的消费者属性正确吗?请向我们展示一个真正的
consumer部分,并确保您在那里使用Deserializer。 -
抱歉,Artem,我的错误,刚刚更新了消费者属性
-
谢谢。现在一切都很好。我们可以在堆栈跟踪中看到更多关于该错误的信息吗?我相信它应该报告这样一个
SerializationException的cause... -
尽可能添加更多内容,当我在 ide (vscode) 中运行消费者服务时,ide 冻结,甚至不允许我复制整个日志,所以我需要重新启动
标签: spring-boot apache-kafka spring-kafka