【问题标题】:Spring Kafka multiple serializers and consumer/container factoriesSpring Kafka 多个序列化器和消费者/容器工厂
【发布时间】:2017-11-23 21:10:57
【问题描述】:

在我的 Spring Boot 应用程序中,我使用以下 Sender/Receiver 配置了 Kafka:

@Configuration
public class KafkaSenderConfig {

    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Bean
    public Map<String, Object> producerConfigs() {

        Map<String, Object> props = new HashMap<>();

        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);

        return props;
    }

    @Bean
    public ProducerFactory<String, ImportDecisionMessage> producerFactory() {
        return new DefaultKafkaProducerFactory<>(producerConfigs());
    }

    @Bean
    public KafkaTemplate<String, ImportDecisionMessage> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }

}

@Configuration
public class KafkaReceiverConfig {

    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Value("${spring.kafka.consumer.group-id}")
    private String consumerGroupId;

    @Bean
    public Map<String, Object> consumerConfigs() {

        Map<String, Object> props = new HashMap<>();

        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroupId);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);

        return props;
    }

    @Bean
    public ConsumerFactory<String, ImportDecisionMessage> consumerFactory() {
        return new DefaultKafkaConsumerFactory<>(consumerConfigs(), new StringDeserializer(), new JsonDeserializer<>(ImportDecisionMessage.class));
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, ImportDecisionMessage> kafkaListenerContainerFactory() {

        ConcurrentKafkaListenerContainerFactory<String, ImportDecisionMessage> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());

        return factory;
    }

}

基于此配置,我现在只能使用我的 POJO - ImportDecisionMessage 并使用 JsonDeserializer 序列化程序发送/接收消息。

我还必须有可能发送另一个 POJO 作为消息,例如,ProductCarCategory

另外,我还想用另一个org.apache.kafka.common.serialization.BytesSerializer 代替Car

如何正确扩展我的配置以支持这些类型和序列化程序?

【问题讨论】:

    标签: spring spring-boot apache-kafka spring-kafka


    【解决方案1】:

    这是不可能的,因为:

    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
    

    完全是 Apache Kafka 属性。而那个只支持简单的普通反序列化策略。您应该考虑已经从 Kafka 消费者返回的原始 byte[] 下游自定义逻辑。

    【讨论】:

    • 感谢您的回答。这是否意味着可以为所有应用程序(例如 String、JSON 或 Bytes.. 等)配置唯一一对 Serializer/Deserializer?您能否还展示如何扩展我的配置以支持其他消息类型,例如 Product、Car、Category(不仅是 ImportDecisionMessage)?
    • 不,这是每个消费者的。您可以使用不同的选项配置多个使用者。但实际上他们至少都应该研究不同的分区。 OTOH ContainerFactory 也可以配置MessageConverterdocs.spring.io/spring-kafka/docs/2.0.1.RELEASE/reference/html/…
    猜你喜欢
    • 2018-11-08
    • 2017-03-06
    • 2019-08-30
    • 2019-04-07
    • 2018-08-04
    • 1970-01-01
    • 2020-02-24
    • 2019-08-25
    • 1970-01-01
    相关资源
    最近更新 更多