【问题标题】:Spring Cloud Stream Kafka > consuming Avro messages from Confluent REST ProxySpring Cloud Stream Kafka > 使用来自 Confluent REST 代理的 Avro 消息
【发布时间】:2017-02-17 21:14:03
【问题描述】:

我有以下场景:

我的应用程序如下所示:

@SpringBootApplication
@EnableBinding(Sink.class)
public class MyApplication {
  private static Logger log = LoggerFactory.getLogger(MyApplication.class);

  public static void main(String[] args) {
    SpringApplication.run(MyApplication.class, args);
  }

  @StreamListener(Sink.INPUT)
  public void myMessageSink(MyMessage message) {
    log.info("Received new message: {}", message);
  }
}

而 MyMessage 是 Avro 从 Avro 架构创建的类。

我的 application.properties 如下所示:

spring.cloud.stream.bindings.input.destination=myTopic
spring.cloud.stream.bindings.input.group=${spring.application.name}
spring.cloud.stream.bindings.input.contentType=application/*+avro

我现在的问题是每次收到新消息,都会抛出如下异常:

org.springframework.messaging.MessagingException: Exception thrown while invoking MyApplication#myMessageSink[1 args]; nested exception is org.apache.avro.AvroRuntimeException: Malformed data. Length is negative: -27
    at org.springframework.cloud.stream.binding.StreamListenerAnnotationBeanPostProcessor$StreamListenerMessageHandler.handleRequestMessage(StreamListenerAnnotationBeanPostProcessor.java:316) ~[spring-cloud-stream-1.1.0.RELEASE.jar:1.1.0.RELEASE]
    at org.springframework.integration.handler.AbstractReplyProducingMessageHandler.handleMessageInternal(AbstractReplyProducingMessageHandler.java:109) ~[spring-integration-core-4.3.2.RELEASE.jar:4.3.2.RELEASE]
    at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:127) ~[spring-integration-core-4.3.2.RELEASE.jar:4.3.2.RELEASE]
    at org.springframework.integration.dispatcher.AbstractDispatcher.tryOptimizedDispatch(AbstractDispatcher.java:116) ~[spring-integration-core-4.3.2.RELEASE.jar:4.3.2.RELEASE]
    at org.springframework.integration.dispatcher.UnicastingDispatcher.doDispatch(UnicastingDispatcher.java:148) ~[spring-integration-core-4.3.2.RELEASE.jar:4.3.2.RELEASE]
    ...
Caused by: org.apache.avro.AvroRuntimeException: Malformed data. Length is negative: -27
    at org.apache.avro.io.BinaryDecoder.doReadBytes(BinaryDecoder.java:336) ~[avro-1.8.1.jar:1.8.1]
    at org.apache.avro.io.BinaryDecoder.readString(BinaryDecoder.java:263) ~[avro-1.8.1.jar:1.8.1]
    at org.apache.avro.io.ResolvingDecoder.readString(ResolvingDecoder.java:201) ~[avro-1.8.1.jar:1.8.1]
    at org.apache.avro.generic.GenericDatumReader.readString(GenericDatumReader.java:430) ~[avro-1.8.1.jar:1.8.1]
    at org.apache.avro.generic.GenericDatumReader.readString(GenericDatumReader.java:422) ~[avro-1.8.1.jar:1.8.1]
    at org.apache.avro.generic.GenericDatumReader.readMapKey(GenericDatumReader.java:335) ~[avro-1.8.1.jar:1.8.1]
    at org.apache.avro.generic.GenericDatumReader.readMap(GenericDatumReader.java:321) ~[avro-1.8.1.jar:1.8.1]
    at org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:177) ~[avro-1.8.1.jar:1.8.1]
    at org.apache.avro.specific.SpecificDatumReader.readField(SpecificDatumReader.java:116) ~[avro-1.8.1.jar:1.8.1]
    at org.apache.avro.generic.GenericDatumReader.readRecord(GenericDatumReader.java:230) ~[avro-1.8.1.jar:1.8.1]
    at org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:174) ~[avro-1.8.1.jar:1.8.1]
    at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:152) ~[avro-1.8.1.jar:1.8.1]
    at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:144) ~[avro-1.8.1.jar:1.8.1]
    at org.springframework.cloud.stream.schema.avro.AbstractAvroMessageConverter.convertFromInternal(AbstractAvroMessageConverter.java:91) ~[spring-cloud-stream-schema-1.1.0.RELEASE.jar:1.1.0.RELEASE]
    at org.springframework.messaging.converter.AbstractMessageConverter.fromMessage(AbstractMessageConverter.java:175) ~[spring-messaging-4.3.3.RELEASE.jar:4.3.3.RELEASE]
    at org.springframework.messaging.converter.CompositeMessageConverter.fromMessage(CompositeMessageConverter.java:67) ~[spring-messaging-4.3.3.RELEASE.jar:4.3.3.RELEASE]
    at org.springframework.messaging.handler.annotation.support.PayloadArgumentResolver.resolveArgument(PayloadArgumentResolver.java:117) ~[spring-messaging-4.3.3.RELEASE.jar:4.3.3.RELEASE]
    at org.springframework.messaging.handler.invocation.HandlerMethodArgumentResolverComposite.resolveArgument(HandlerMethodArgumentResolverComposite.java:112) ~[spring-messaging-4.3.3.RELEASE.jar:4.3.3.RELEASE]
    at org.springframework.messaging.handler.invocation.InvocableHandlerMethod.getMethodArgumentValues(InvocableHandlerMethod.java:138) ~[spring-messaging-4.3.3.RELEASE.jar:4.3.3.RELEASE]
    at org.springframework.messaging.handler.invocation.InvocableHandlerMethod.invoke(InvocableHandlerMethod.java:107) ~[spring-messaging-4.3.3.RELEASE.jar:4.3.3.RELEASE]
    at org.springframework.cloud.stream.binding.StreamListenerAnnotationBeanPostProcessor$StreamListenerMessageHandler.handleRequestMessage(StreamListenerAnnotationBeanPostProcessor.java:307) ~[spring-cloud-stream-1.1.0.RELEASE.jar:1.1.0.RELEASE]
    ... 35 common frames omitted

据我了解,问题在于 Confluent 堆栈包含消息架构的 ID 作为消息负载的一部分,并且客户端应该在架构 ID 之后开始读取实际的 Avro 消息。 看来我需要配置 Kafka 绑定以使用 Confluent 的 KafkaAvroDeserializer,但我不知道如何实现这一点。

(我可以使用 Confluent 的 avro 控制台消费者完美地检索消息,因此 Avro 编码似乎不是问题)

我还尝试使用 @EnableSchemaRegistry 注释并配置 ConfluentSchemaRegistryClient bean,但在我看来,这仅控制模式的存储/检索位置,而不是实际的反序列化。

这甚至应该以某种方式工作吗?

【问题讨论】:

    标签: java spring apache-kafka spring-cloud-stream confluent-platform


    【解决方案1】:

    per-binding 属性spring.cloud.stream.kafka.bindings.input.consumer.configuration.value.deserializer 设置为Confluent's KafkaAvroDeserializer class name 时是否有效?

    【讨论】:

    • 看起来KafkaMessageChannelBinder 的消费者端点总是将ByteArrayDeserializer 用于键/值序列化程序。这可能是KafkaMessageChannelBinder 的生产者默认使用ByteArraySerializer 的结果。对于非基于 Spring Cloud Stream 的生产者,消费者应该能够使用上述属性覆盖所需的反序列化器。
    • 很遗憾没有,我已经试过了。从我在source code 中看到的内容来看,解串器被硬编码为默认的 ByteArrayDeserializer。
    • 是的,请在此处创建问题:github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues。我们可以从那里追踪它。谢谢!
    【解决方案2】:

    有点回答我自己的问题。我现在所做的是实现一个 MessageConverter,它只删除任何消息的前 4 个字节,然后再将它们传递给 Avro 解码器。代码大部分取自spring-cloud-stream的AbstractAvroMessageConverter:

    public class ConfluentAvroSchemaMessageConverter extends AvroSchemaMessageConverter {
    
    public ConfluentAvroSchemaMessageConverter() {
        super(new MimeType("application", "avro+confluent"));
    }
    
    @Override
    protected Object convertFromInternal(Message<?> message, Class<?> targetClass, Object conversionHint) {
        Object result = null;
        try {
            byte[] payload = (byte[]) message.getPayload();
    
            // byte array to contain the message without the confluent header (first 4 bytes)
            byte[] payloadWithoutConfluentHeader = new byte[payload.length - 4];
            ByteBuffer buf = ByteBuffer.wrap(payload);
            MimeType mimeType = getContentTypeResolver().resolve(message.getHeaders());
            if (mimeType == null) {
                if (conversionHint instanceof MimeType) {
                    mimeType = (MimeType) conversionHint;
                }
                else {
                    return null;
                }
            }
    
            // read first 4 bytes and copy the rest to the new byte array
            // see https://groups.google.com/forum/#!topic/confluent-platform/rI1WNPp8DJU
            buf.getInt();
            buf.get(payloadWithoutConfluentHeader);
    
            Schema writerSchema = resolveWriterSchemaForDeserialization(mimeType);
            Schema readerSchema = resolveReaderSchemaForDeserialization(targetClass);
            DatumReader<Object> reader = getDatumReader((Class<Object>) targetClass, readerSchema, writerSchema);
            Decoder decoder = DecoderFactory.get().binaryDecoder(payloadWithoutConfluentHeader, null);
            result = reader.read(null, decoder);
        }
        catch (IOException e) {
                throw new MessageConversionException(message, "Failed to read payload", e);
        }
        return result;
    
    }
    

    然后我通过 application.properties 将传入 Kafka 主题的内容类型设置为 application/avro+confluent。

    这至少让我可以检索使用 Confluent 堆栈编码的消息,但它当然不会以任何方式与模式注册表交互。

    【讨论】:

      猜你喜欢
      • 2020-03-20
      • 2018-03-09
      • 1970-01-01
      • 2016-11-10
      • 2018-07-18
      • 2018-01-31
      • 2018-05-17
      • 1970-01-01
      • 2021-06-12
      相关资源
      最近更新 更多