【问题标题】:Spring-Cloud-Stream-Kafka-Binder functional style ignores custom De/Serializer and/or useNativeEncoding?Spring-Cloud-Stream-Kafka-Binder 功能样式忽略自定义 De/Serializer 和/或 useNativeEncoding?
【发布时间】:2020-04-10 03:17:54
【问题描述】:

我们刚刚升级到 Spring-Cloud-Stream 的 3.0.0-Release,遇到以下问题:

当使用这样的函数样式时:

public class EventProcessor {

    private final PriceValidator priceValidator;

    @Bean
    public Function<Flux<EnrichedValidationRequest>, Flux<ValidationResult>> validate() {
        return enrichedValidationRequestFlux -> enrichedValidationRequestFlux
                .map(ProcessingContext::new)
                .flatMap(priceValidator::validateAndMap);
    }
}

application.yaml 如下所示:

spring.cloud.stream:
  default-binder: kafka
  kafka:
    binder:
      brokers: ${kafka.broker.prod}
      auto-create-topics: false
  function.definition: validate

# INPUT: enrichedValidationRequests
spring.cloud.stream.bindings.validate-in-0:
  destination: ${kafka.topic.${spring.application.name}.input.enrichedValidationRequests}
  group: ${spring.application.name}.${STAGE:NOT_SET}
  consumer:
    useNativeDecoding: true


spring.cloud.stream.kafka.bindings.validate-in-0:
  consumer:
    configuration:
      key.deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value.deserializer: de.pricevalidator.deserializer.EnrichedValidationRequestDeserializer


# OUTPUT: validationResults
spring.cloud.stream.bindings.validate-out-0:
  destination: validationResultsTmp
  producer:
    useNativeEncoding: true

spring.cloud.stream.kafka.bindings.validate-out-0:
  producer:
    compression.type: lz4
    messageKeyExpression: payload.offerKey
    configuration:
      key.serializer: org.apache.kafka.common.serialization.StringSerializer
      value.serializer: de.pricevalidator.serializer.ValidationResultSerializer

似乎序列化完成了两次——当我们拦截在 kafka 主题中产生的消息时,消费者只会将它们显示为 JSON(字符串),但现在它是一个不可读的字节 []。此外,生产中的下游消费者无法再反序列化消息。奇怪的是,输入消息的反序列化似乎工作得很好,无论我们将什么放入消费者属性(在 binder 或默认 kafka 级别) 我们有一种感觉,这个bug“又回来了”,但是我们在代码中找不到确切的位置:https://github.com/spring-cloud/spring-cloud-stream/issues/1536

我们的(丑陋的)解决方法:

@Slf4j
@Configuration
public class KafkaMessageConverterConfiguration {

    @ConditionalOnProperty(value = "spring.cloud.stream.default-binder", havingValue = "kafka")
    @Bean
    public MessageConverter validationResultConverter(BinderTypeRegistry binder, ObjectMapper objectMapper) {
        return new AbstractMessageConverter(MimeType.valueOf("application/json")) {
            @Override
            protected boolean supports(final Class<?> clazz) {
                return ValidationResult.class.isAssignableFrom(clazz);
            }

            @Override
            protected Object convertToInternal(final Object payload, final MessageHeaders headers, final Object conversionHint) {
                return payload;
            }
        };
    }
}

是否有一种“正确”的方式来设置自定义序列化程序或像以前一样获取本机编码?

【问题讨论】:

  • 有趣的是,我们在反序列化时(在另一个不使用函数式样式,但使用 StreamListener 方法的反应式项目中)有一个可能相关的问题,并且 MappingJackson2MessageConverter.convertFromInternal 中的 targetClass 最终成为 Flux -它应该简单地返回有效负载,但有效负载当然不是 Flux,而是一个不同的对象。在 3.0.0 版本之前,我们总是包含 spring-cloud-stream-reactive 依赖,它提供了一个名为“MessageChannelToInputFluxParameterAdapter”的 bean,它就是这样做的。

标签: java spring-boot apache-kafka spring-cloud-stream


【解决方案1】:

所以这是在 3.0.0.RELEASE - https://github.com/spring-cloud/spring-cloud-stream/commit/74aee8102898dbff96a570d2d2624571b259e141 之后报告的问题。它已得到解决,几天后将在 3.0.1.RELEASE (Horsham.SR1) 中提供。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-12-16
    • 2021-02-17
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多