【问题标题】:Spring Kafka and exactly once delivery guaranteeSpring Kafka 和一次性交付保证
【发布时间】:2019-03-05 07:46:50
【问题描述】:

我使用 Spring Kafka 和 Spring Boot,只是想知道如何配置我的消费者,例如:

@KafkaListener(topics = "${kafka.topic.post.send}", containerFactory = "postKafkaListenerContainerFactory")
public void sendPost(ConsumerRecord<String, Post> consumerRecord, Acknowledgment ack) {

    // do some logic

    ack.acknowledge();
}

要使用一次性交货保证吗?

我应该只在sendPost 方法上添加org.springframework.transaction.annotation.Transactional 注释,仅此而已,还是我需要执行一些额外的步骤才能实现这一点?

更新

这是我当前的配置

@Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(KafkaProperties kafkaProperties, KafkaTransactionManager<Object, Object> transactionManager) {

        kafkaProperties.getProperties().put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, kafkaConsumerMaxPollIntervalMs);
        kafkaProperties.getProperties().put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, kafkaConsumerMaxPollRecords);

        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        //factory.getContainerProperties().setAckMode(AckMode.MANUAL_IMMEDIATE);
        factory.getContainerProperties().setTransactionManager(transactionManager);
        factory.setConsumerFactory(consumerFactory(kafkaProperties));

        return factory;
    }


    @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);
        props.put(ProducerConfig.MAX_REQUEST_SIZE_CONFIG, 15000000);

        return props;
    }

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

    @Bean
    public KafkaTemplate<String, Post> postKafkaTemplate() {
        return new KafkaTemplate<>(postProducerFactory());
    }

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

    @Bean
    public KafkaTemplate<String, Update> updateKafkaTemplate() {
        return new KafkaTemplate<>(updateProducerFactory());
    }

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

    @Bean
    public KafkaTemplate<String, Message> messageKafkaTemplate() {
        return new KafkaTemplate<>(messageProducerFactory());
    }

但它失败并出现以下错误:

***************************
APPLICATION FAILED TO START
***************************

Description:

Parameter 0 of method kafkaTransactionManager in org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration required a single bean, but 3 were found:
    - postProducerFactory: defined by method 'postProducerFactory' in class path resource [com/example/domain/configuration/messaging/KafkaProducerConfig.class]
    - updateProducerFactory: defined by method 'updateProducerFactory' in class path resource [com/example/domain/configuration/messaging/KafkaProducerConfig.class]
    - messageProducerFactory: defined by method 'messageProducerFactory' in class path resource [com/example/domain/configuration/messaging/KafkaProducerConfig.class]

我做错了什么?

【问题讨论】:

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


    【解决方案1】:

    您不应使用手动确认。相反,将KafkaTransactionManager 注入到侦听器容器中,当侦听器方法正常退出(否则回滚)时,容器将向事务发送偏移量。

    您不应该只通过消费者执行一次确认。

    编辑

    application.yml

    spring:
      kafka:
        consumer:
          auto-offset-reset: earliest
          enable-auto-commit: false
          properties:
            isolation:
              level: read_committed
        producer:
          transaction-id-prefix: myTrans.
    

    应用程序

    @SpringBootApplication
    public class So52570118Application {
    
        public static void main(String[] args) {
            SpringApplication.run(So52570118Application.class, args);
        }
    
        @Bean // override boot's auto-config to add txm
        public ConcurrentKafkaListenerContainerFactory<?, ?> kafkaListenerContainerFactory(
                ConcurrentKafkaListenerContainerFactoryConfigurer configurer,
                ConsumerFactory<Object, Object> kafkaConsumerFactory,
                KafkaTransactionManager<Object, Object> transactionManager) {
            ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
            configurer.configure(factory, kafkaConsumerFactory);
            factory.getContainerProperties().setTransactionManager(transactionManager);
            return factory;
        }
    
        @Autowired
        private KafkaTemplate<String, String> template;
    
        @KafkaListener(id = "so52570118", topics = "so52570118")
        public void listen(String in) throws Exception {
            System.out.println(in);
            Thread.sleep(5_000);
            this.template.send("so52570118out", in.toUpperCase());
            System.out.println("sent");
        }
    
        @KafkaListener(id = "so52570118out", topics = "so52570118out")
        public void listenOut(String in) {
            System.out.println(in);
        }
    
        @Bean
        public ApplicationRunner runner() {
            return args -> this.template.executeInTransaction(t -> t.send("so52570118", "foo"));
        }
    
        @Bean
        public NewTopic topic1() {
            return new NewTopic("so52570118", 1, (short) 1);
        }
    
        @Bean
        public NewTopic topic2() {
            return new NewTopic("so52570118out", 1, (short) 1);
        }
    
    }
    

    【讨论】:

    • 感谢您的回答。您能否举例说明如何在 Spring Boot 中使用 KafkaTransactionManager
    • 实际上,您不需要事务管理器 bean,启动时自动配置一个(已编辑)。
    • 谢谢!我更新了我的配置,但它因错误而失败。我已经更新了我的问题。我在那里做错了什么?
    • 由于您要定义自己的生产者工厂(并且不止一个),因此需要将其中一个(您希望为其运行事务的那个)标记为@Primary - 或者您可以明确定义您自己的事务管理器@Bean,而不是使用 Boots 自动配置的。
    • 鉴于您为每个使用相同的 producerConfigs,为什么不直接使用 &lt;String, Object&gt; 创建一个?如果您希望每个模板的发送参与同一个事务,则必须这样做,因为事务的范围是生产者。
    猜你喜欢
    • 2019-10-01
    • 2017-06-14
    • 2020-07-29
    • 1970-01-01
    • 2018-12-09
    • 2018-12-02
    • 2015-07-03
    • 1970-01-01
    • 2019-09-16
    相关资源
    最近更新 更多