首页
学习
活动
专区
圈层
工具
发布
首页
学习
活动
专区
圈层
工具
MCP广场
社区首页 >问答首页 >Spring错误地列出了StringDeserializer而不是KafkaAvroDeserializer for value.deserializer

Spring错误地列出了StringDeserializer而不是KafkaAvroDeserializer for value.deserializer
EN

Stack Overflow用户
提问于 2021-05-12 18:19:38
回答 1查看 389关注 0票数 2

当用spring-kafa启动Spring 2.4.5时,value.deserializer值显示为StringDeserializer而不是KafkaAvroDeserializer

代码语言:javascript
运行
复制
2021-05-12 13:46:05.313  INFO 12632 --- [           main] o.a.k.clients.consumer.ConsumerConfig    : ConsumerConfig values: 

...  Elided for brevity ...

key.deserializer = class org.apache.kafka.common.serialization.StringDeserializer
value.deserializer = class org.apache.kafka.common.serialization.StringDeserializer

My配置:

代码语言:javascript
运行
复制
@Bean
public Map<String, User> consumerConfigAvro() {
    Map<String, Object> props = new HashMap<>();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
    props.put(ConsumerConfig.GROUP_ID_CONFIG, "dnk23");
    props.put( AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://localhost:18081");
    props.put( KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, true );
    props.put( KafkaAvroSerializerConfig.VALUE_SUBJECT_NAME_STRATEGY, TopicRecordNameStrategy.class.getName());

    return props;
}

@Bean
public ConsumerFactory<String, User> consumerFactoryAvro() {
    KafkaAvroDeserializer kafkaAvroDeserializer = new KafkaAvroDeserializer();
    kafkaAvroDeserializer.configure(consumerConfigAvro(), false);
    return new DefaultKafkaConsumerFactory(consumerConfigAvro(),
            new StringDeserializer(),
            kafkaAvroDeserializer);
}

@Bean
public ConcurrentKafkaListenerContainerFactory<String, User> kafkaListenerContainerFactoryAvro() {
    ConcurrentKafkaListenerContainerFactory<String, User> factory =
            new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactoryAvro());
    return factory;
}

我遇到了以下异常,但在花费了大量时间之后,我认为根本原因是value.deserializer设置不正确。

代码语言:javascript
运行
复制
    org.springframework.kafka.listener.ListenerExecutionFailedException: Listener method could not be invoked with the incoming message
Endpoint handler details:
Method [public void com.example.sandbox.SandboxApplication.listen3(com.example.sandbox.avro.User)]
Bean [com.example.sandbox.SandboxApplication$$EnhancerBySpringCGLIB$$9af3650b@5afcde28]; nested exception is org.springframework.messaging.converter.MessageConversionException: Cannot handle message; nested exception is org.springframework.messaging.converter.MessageConversionException: Cannot convert from [java.lang.String] to [com.example.sandbox.avro.User] for GenericMessage [payload=    $Henry Green Engine, headers={kafka_offset=18, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@259a8378, kafka_timestampType=CREATE_TIME, X-APP-EVENT=ApplicationCreatedEvent, kafka_receivedPartitionId=0, kafka_receivedMessageKey=some-key, kafka_receivedTopic=topic3, kafka_receivedTimestamp=1620841566436, kafka_groupId=myId3}], failedMessage=GenericMessage [payload=    $Henry Green Engine, headers={kafka_offset=18, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@259a8378, kafka_timestampType=CREATE_TIME, X-APP-EVENT=ApplicationCreatedEvent, kafka_receivedPartitionId=0, kafka_receivedMessageKey=some-key, kafka_receivedTopic=topic3, kafka_receivedTimestamp=1620841566436, kafka_groupId=myId3}]; nested exception is org.springframework.messaging.converter.MessageConversionException: Cannot handle message; nested exception is org.springframework.messaging.converter.MessageConversionException: Cannot convert from [java.lang.String] to [com.example.sandbox.avro.User] for GenericMessage [payload=    $Henry Green Engine, headers={kafka_offset=18, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@259a8378, kafka_timestampType=CREATE_TIME, X-APP-EVENT=ApplicationCreatedEvent, kafka_receivedPartitionId=0, kafka_receivedMessageKey=some-key, kafka_receivedTopic=topic3, kafka_receivedTimestamp=1620841566436, kafka_groupId=myId3}], failedMessage=GenericMessage [payload=    $Henry Green Engine, headers={kafka_offset=18, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@259a8378, kafka_timestampType=CREATE_TIME, X-APP-EVENT=ApplicationCreatedEvent, kafka_receivedPartitionId=0, kafka_receivedMessageKey=some-key, kafka_receivedTopic=topic3, kafka_receivedTimestamp=1620841566436, kafka_groupId=myId3}]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.decorateException(KafkaMessageListenerContainer.java:2114) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeErrorHandler(KafkaMessageListenerContainer.java:2102) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeRecordListener(KafkaMessageListenerContainer.java:2001) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeWithRecords(KafkaMessageListenerContainer.java:1928) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeRecordListener(KafkaMessageListenerContainer.java:1814) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeListener(KafkaMessageListenerContainer.java:1531) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1178) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1075) ~[spring-kafka-2.6.7.jar:2.6.7]
    at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) ~[na:na]
    at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) ~[na:na]
    at java.base/java.lang.Thread.run(Thread.java:834) ~[na:na]
Caused by: org.springframework.messaging.converter.MessageConversionException: Cannot handle message; nested exception is org.springframework.messaging.converter.MessageConversionException: Cannot convert from [java.lang.String] to [com.example.sandbox.avro.User] for GenericMessage [payload=    $Henry Green Engine, headers={kafka_offset=18, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@259a8378, kafka_timestampType=CREATE_TIME, X-APP-EVENT=ApplicationCreatedEvent, kafka_receivedPartitionId=0, kafka_receivedMessageKey=some-key, kafka_receivedTopic=topic3, kafka_receivedTimestamp=1620841566436, kafka_groupId=myId3}], failedMessage=GenericMessage [payload=    $Henry Green Engine, headers={kafka_offset=18, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@259a8378, kafka_timestampType=CREATE_TIME, X-APP-EVENT=ApplicationCreatedEvent, kafka_receivedPartitionId=0, kafka_receivedMessageKey=some-key, kafka_receivedTopic=topic3, kafka_receivedTimestamp=1620841566436, kafka_groupId=myId3}]
    at org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter.invokeHandler(MessagingMessageListenerAdapter.java:341) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.kafka.listener.adapter.RecordMessagingMessageListenerAdapter.onMessage(RecordMessagingMessageListenerAdapter.java:86) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.kafka.listener.adapter.RecordMessagingMessageListenerAdapter.onMessage(RecordMessagingMessageListenerAdapter.java:51) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeOnMessage(KafkaMessageListenerContainer.java:2069) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeOnMessage(KafkaMessageListenerContainer.java:2051) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeRecordListener(KafkaMessageListenerContainer.java:1988) ~[spring-kafka-2.6.7.jar:2.6.7]
    ... 8 common frames omitted
Caused by: org.springframework.messaging.converter.MessageConversionException: Cannot convert from [java.lang.String] to [com.example.sandbox.avro.User] for GenericMessage [payload=    $Henry Green Engine, headers={kafka_offset=18, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@259a8378, kafka_timestampType=CREATE_TIME, X-APP-EVENT=ApplicationCreatedEvent, kafka_receivedPartitionId=0, kafka_receivedMessageKey=some-key, kafka_receivedTopic=topic3, kafka_receivedTimestamp=1620841566436, kafka_groupId=myId3}]
    at org.springframework.messaging.handler.annotation.support.PayloadMethodArgumentResolver.resolveArgument(PayloadMethodArgumentResolver.java:145) ~[spring-messaging-5.3.6.jar:5.3.6]
    at org.springframework.kafka.annotation.KafkaListenerAnnotationBeanPostProcessor$KafkaNullAwarePayloadArgumentResolver.resolveArgument(KafkaListenerAnnotationBeanPostProcessor.java:926) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.messaging.handler.invocation.HandlerMethodArgumentResolverComposite.resolveArgument(HandlerMethodArgumentResolverComposite.java:117) ~[spring-messaging-5.3.6.jar:5.3.6]
    at org.springframework.messaging.handler.invocation.InvocableHandlerMethod.getMethodArgumentValues(InvocableHandlerMethod.java:148) ~[spring-messaging-5.3.6.jar:5.3.6]
    at org.springframework.messaging.handler.invocation.InvocableHandlerMethod.invoke(InvocableHandlerMethod.java:116) ~[spring-messaging-5.3.6.jar:5.3.6]
    at org.springframework.kafka.listener.adapter.HandlerAdapter.invoke(HandlerAdapter.java:48) ~[spring-kafka-2.6.7.jar:2.6.7]
    at org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter.invokeHandler(MessagingMessageListenerAdapter.java:330) ~[spring-kafka-2.6.7.jar:2.6.7]
    ... 13 common frames omitted

侦听器

代码语言:javascript
运行
复制
@KafkaListener(id = "myId3", topics = "topic3")
public void listen3(@Payload User rec) {
    System.out.println(rec.getName());
}

生产者

代码语言:javascript
运行
复制
                User user = User.newBuilder()
                        .setName("Henry Green Engine")
                        .setNumber(count)
                        .build();

                Message<User> message3 = MessageBuilder
                        .withPayload(user)
                        .setHeader(KafkaHeaders.TOPIC, "topic3")
                        .setHeader(KafkaHeaders.MESSAGE_KEY, "some-key")
                        .setHeader(KafkaHeaders.PARTITION_ID, 0)
                        .setHeader("X-APP-EVENT", "ApplicationCreatedEvent")
                        .build();

                kafkaTemplate3.send(message3);

Avro模式

代码语言:javascript
运行
复制
{
  "namespace": "com.example.sandbox.avro",
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "name", "type": "string"},
    {"name": "number", "type": "int"}
  ]
}
EN

Stack Overflow用户

回答已采纳

发布于 2021-05-12 18:39:38

当使用非标准容器工厂bean名称(默认为kafkaListenerContainerFactory)时,必须在@KafkaListener上指定工厂bean名称。

代码语言:javascript
运行
复制
public @interface KafkaListener {

...

    /**
     * The bean name of the {@link org.springframework.kafka.config.KafkaListenerContainerFactory}
     * to use to create the message listener container responsible to serve this endpoint.
     * <p>
     * If not specified, the default container factory is used, if any. If a SpEL
     * expression is provided ({@code #{...}}), the expression can either evaluate to a
     * container factory instance or a bean name.
     * @return the container factory bean name.
     */
    String containerFactory() default "";

所以..., containerFactory="kafkaListenerContainerFactoryAvro", ...)

票数 2
EN
查看全部 1 条回答
页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/67509073

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档