【问题标题】:spring-cloud-stream message conversion exceptionspring-cloud-stream 消息转换异常
【发布时间】:2018-03-21 14:33:52
【问题描述】:

在将我们的一项服务升级到 spring-cloud-stream 2.0.0.RC3 时,我们在尝试使用由使用旧版本 spring-cloud-stream 的服务生成的消息时遇到异常 - Ditmars.RELEASE:

错误 31241 --- [container-4-C-1] o.s.integration.handler.LoggingHandler : org.springframework.messaging.converter.MessageConversionException: 无法从 [[B] 转换为 [com.watercorp.messaging.types .incoming.UsersDeletedMessage] for GenericMessage [payload=byte[371], headers={kafka_offset=1, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@62029d0d, kafka_timestampType=CREATE_TIME, message_id=1645508761, id=f4e947de- 22e6-b629-229b-4fa961c73f2d,type=USERS_DELETED,kafka_receivedPartitionId=4,contentType=text/plain,kafka_receivedTopic=user,kafka_receivedTimestamp=1521641760698,timestamp=1521641772477}],failedMessage=GenericMessages={payload= kafka_offset=1, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@62029d0d, kafka_timestampType=CREATE_TIME, message_id=1645508761, id=f4e947de-22e6-b629-229b-4fa961c73f2d, type=USERS_DELETEDIdreive=4,d content=文本/纯文本,kafka_receivedTopic=用户,kafka_receivedTimesta mp=1521641760698,时间戳=1521641772477}] 在 org.springframework.messaging.handler.annotation.support.PayloadArgumentResolver.resolveArgument(PayloadArgumentResolver.java:144) 在 org.springframework.messaging.handler.invocation.HandlerMethodArgumentResolverComposite.resolveArgument(HandlerMethodArgumentResolverComposite.java:116) 在 org.springframework.messaging.handler.invocation.InvocableHandlerMethod.getMethodArgumentValues(InvocableHandlerMethod.java:137) 在 org.springframework.messaging.handler.invocation.InvocableHandlerMethod.invoke(InvocableHandlerMethod.java:109) 在 org.springframework.cloud.stream.binding.StreamListenerMessageHandler.handleRequestMessage(StreamListenerMessageHandler.java:55) 在 org.springframework.integration.handler.AbstractReplyProducingMessageHandler.handleMessageInternal(AbstractReplyProducingMessageHandler.java:109) 在 org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:164) 在 org.springframework.cloud.stream.binding.DispatchingStreamListenerMessageHandler.handleRequestMessage(DispatchingStreamListenerMessageHandler.java:87) 在 org.springframework.integration.handler.AbstractReplyProducingMessageHandler.handleMessageInternal(AbstractReplyProducingMessageHandler.java:109) 在 org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:157) 在 org.springframework.integration.dispatcher.AbstractDispatcher.tryOptimizedDispatch(AbstractDispatcher.java:116) 在 org.springframework.integration.dispatcher.UnicastingDispatcher.doDispatch(UnicastingDispatcher.java:132) 在 org.springframework.integration.dispatcher.UnicastingDispatcher.dispatch(UnicastingDispatcher.java:105) 在 org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:73) 在 org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:463) 在 org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:407) 在 org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:181) 在 org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:160) 在 org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:47) 在 org.springframework.messaging.core.AbstractMessageSendingTemplate.send(AbstractMessageSendingTemplate.java:108) 在 org.springframework.integration.endpoint.MessageProducerSupport.sendMessage(MessageProducerSupport.java:203) 在 org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter.access 300 美元(KafkaMessageDrivenChannelAdapter.java:70) 在 org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter$IntegrationRecordMessageListener.onMessage(KafkaMessageDrivenChannelAdapter.java:387) 在 org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter$IntegrationRecordMessageListener.onMessage(KafkaMessageDrivenChannelAdapter.java:364) 在 org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeRecordListener(KafkaMessageListenerContainer.java:1001) 在 org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeWithRecords(KafkaMessageListenerContainer.java:981) 在 org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeRecordListener(KafkaMessageListenerContainer.java:932) 在 org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeListener(KafkaMessageListenerContainer.java:801) 在 org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:689) 在 java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) 在 java.util.concurrent.FutureTask.run(FutureTask.java:266) 在 java.lang.Thread.run(Thread.java:745)

看起来原因是与消息一起发送的contentType 标头是text/plain,尽管它应该是application/json
生产者配置:

春天: 云: 溪流: 卡夫卡: 粘合剂: 经纪人:卡夫卡 默认代理端口:9092 zkNodes:动物园管理员 默认ZkPort:2181 最小分区计数:2 复制因子:1 自动创建主题:真 自动添加分区:真 标头:类型,message_id 必填项:1 配置: “[security.protocol]”:PLAINTEXT #TODO:这是一种解决方法。应该是 security.protocol 绑定: 用户输出: 生产商: 同步:真 配置: 重试次数:10000 默认: 粘合剂:卡夫卡 内容类型:应用程序/json 组:用户服务 消费者: 最大尝试次数:1 生产商: partitionKeyExtractorClass:com.watercorp.user_service.messaging.PartitionKeyExtractor 绑定: 用户输出: 目的地:用户 生产商: 分区数:5 消费者配置: 春天: 云: 溪流: 卡夫卡: 粘合剂: 经纪人:卡夫卡 默认代理端口:9092 最小分区计数:2 复制因子:1 自动创建主题:真 自动添加分区:真 标头:类型,message_id 必填项:1 配置: “[security.protocol]”:PLAINTEXT #TODO:这是一种解决方法。应该是 security.protocol 绑定: 用户输入: 消费者: 自动重新平衡启用:真 自动提交错误:真 enableDlq: 真 默认: 粘合剂:卡夫卡 内容类型:应用程序/json 组:注册服务 消费者: 最大尝试次数:1 标头模式:嵌入标头 生产商: partitionKeyExtractorClass: com.watercorp.messaging.PartitionKeyExtractor 标头模式:嵌入标头 绑定: 用户输入: 目的地:用户 消费者: 并发:5 分区:真

消费者@StreamListener:

@StreamListener(target = UserInput.INPUT, 条件 = "headers['type']=='" + USERS_DELETED + "'") 公共无效句柄UsersDeletedMessage(@Valid UsersDeletedMessage usersDeletedMessage,@Header(值=“kafka_receivedPartitionId”, required = false) String partitionId, @Header(value = KAFKA_TOPIC_HEADER_NAME, required = false) String topic, @Header(MESSAGE_ID_HEADER_NAME) String messageId) throws Throwable { logger.info(String.format("收到的用户删除消息message, message id: %s topic: %s partition: %s", messageId, topic, partitionId)); handleMessageWithRetry(_usersDeletedMessageHandler, usersDeletedMessage, messageId, topic); }

【问题讨论】:

    标签: spring-cloud-stream


    【解决方案1】:

    这是 RC3 中的一个错误; recently fixed on master;它将在月底到期的 GA 版本中发布。同时,您可以尝试使用 2.0.0.BUILD-SNAPSHOT 吗?

    我能够重现该问题并使用快照为我修复了它...

        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-stream</artifactId>
            <version>2.0.0.BUILD-SNAPSHOT</version>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-stream-binder-kafka</artifactId>
            <version>2.0.0.BUILD-SNAPSHOT</version>
            <exclusions>
                <exclusion>
                    <groupId>org.springframework.cloud</groupId>
                    <artifactId>spring-cloud-stream-binder-kafka-core</artifactId>
                </exclusion>
            </exclusions>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-stream-binder-kafka-core</artifactId>
            <version>2.0.0.BUILD-SNAPSHOT</version>
        </dependency>
    

    编辑

    为了完整性:

    Ditmars 制作人

    @SpringBootApplication
    @EnableBinding(Source.class)
    public class So49409104Application {
    
        public static void main(String[] args) {
            SpringApplication.run(So49409104Application.class, args);
        }
    
        @Bean
        public ApplicationRunner runner(MessageChannel output) {
            return args -> {
                Foo foo = new Foo();
                foo.setBar("bar");
                output.send(new GenericMessage<>(foo));
            };
        }
    
    
        public static class Foo {
    
            private String bar;
    
            public String getBar() {
                return this.bar;
            }
    
            public void setBar(String bar) {
                this.bar = bar;
            }
    
            @Override
            public String toString() {
                return "Foo [bar=" + this.bar + "]";
            }
    
        }
    
    }
    

    spring:
      cloud:
        stream:
          bindings:
            output:
              destination: so49409104a
              content-type: application/json
              producer:
                header-mode: embeddedHeaders
    

    艾姆赫斯特消费者:

    @SpringBootApplication
    @EnableBinding(Sink.class)
    public class So494091041Application {
    
        public static void main(String[] args) {
            SpringApplication.run(So494091041Application.class, args);
        }
    
        @StreamListener(Sink.INPUT)
        public void listen(Foo foo) {
            System.out.println(foo);
        }
    
        public static class Foo {
    
            private String bar;
    
            public String getBar() {
                return this.bar;
            }
    
            public void setBar(String bar) {
                this.bar = bar;
            }
    
            @Override
            public String toString() {
                return "Foo [bar=" + this.bar + "]";
            }
    
        }
    
    }
    

    spring:
      cloud:
        stream:
          bindings:
            input:
              group: so49409104
              destination: so49409104a
              consumer:
                header-mode: embeddedHeaders
              content-type: application/json
    

    结果:

    Foo [bar=bar]
    

    header-mode 是必需的,因为 2.0 中的默认值是 native,现在 Kafka 原生支持标头。

    【讨论】:

    • 使用 Elmhurst.BUILD-SNAPSHOT 修复它。再次感谢!
    猜你喜欢
    • 2016-05-31
    • 2019-02-15
    • 2016-09-09
    • 2021-02-20
    • 2020-09-28
    • 2021-06-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多