【问题标题】:Exactly Once Two Kafka Clusters恰好两个 Kafka 集群
【发布时间】:2021-07-21 11:50:14
【问题描述】:

我有 2 个 Kafka 集群。集群 A 和集群 B。这些集群是完全独立的。 我有一个 spring-boot 应用程序,它侦听集群 A 上的主题,转换事件,然后将其生成到集群 B。我只需要一次,因为这些是金融事件。我注意到,在我当前的应用程序中,有时会出现重复项以及错过一些事件。我试图尽我所能只实施一次。其中一篇帖子说 flink 将是比 spring-boot 更好的选择。我应该转到flink吗?请看下面的 Spring 代码。

消费者配置

@Configuration
public class KafkaConsumerConfig {

    @Value("${kafka.server.consumer}")
    String server;

    @Value("${kafka.kerberos.service.name:}")
    String kerberosServiceName;

    @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> config = new HashMap<>();

        config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, server);
        config.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        config.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);
        config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, AvroDeserializer.class);
        config.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
        
        // skip kerberos if no value is provided
        if (kerberosServiceName.length() > 0) {
          config.put("security.protocol", "SASL_PLAINTEXT");
          config.put("sasl.kerberos.service.name", kerberosServiceName);
        }

        return config;
    }

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

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, AccrualSchema> kafkaListenerContainerFactory(ConsumerErrorHandler errorHandler) {
        ConcurrentKafkaListenerContainerFactory<String, AccrualSchema> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        factory.setAutoStartup(true);
        factory.getContainerProperties().setAckMode(AckMode.RECORD);
        
        factory.setErrorHandler(errorHandler);
        
        return factory;
    }

    @Bean
    public KafkaConsumerAccrual receiver() {
        return new KafkaConsumerAccrual();
    }
}

生产者配置

    @Configuration
public class KafkaProducerConfig {

    @Value("${kafka.server.producer}")
    String server;

    @Bean
    public ProducerFactory<String, String> producerFactory() {
        Map<String, Object> config = new HashMap<>();
        config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, server);
        config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        config.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
        config.put(ProducerConfig.ACKS_CONFIG, "all");
        config.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "prod-1");
        config.put(ProducerConfig.COMPRESSION_TYPE_CONFIG,"snappy");
        config.put(ProducerConfig.LINGER_MS_CONFIG, "10");
        config.put(ProducerConfig.BATCH_SIZE_CONFIG, Integer.toString(32*1024));
        return new DefaultKafkaProducerFactory<>(config);
    }

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

}

卡夫卡生产者

    @Service
public class KafkaTopicProducer {
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;
        
    public void topicProducer(String payload, String topic) {
        kafkaTemplate.executeInTransaction(kt->kt.send(topic, payload));            
    }
}

Kafka消费者

public class KafkaConsumerAccrual {

    @Autowired
    KafkaTopicProducer kafkaTopicProducer;

    @Autowired
    Gson gson;
 
    @KafkaListener(topics = "topicname", groupId = "groupid", id = "listenerid")
    public void consume(AccrualSchema accrual,
            @Header(KafkaHeaders.RECEIVED_PARTITION_ID) Integer partition, @Header(KafkaHeaders.OFFSET) Long offset,
            @Header(KafkaHeaders.CONSUMER) Consumer<?, ?> consumer) {       
        
        
        AccrualEntity accrualEntity = convertBusinessObjectToAccrual(accrual,partition,offset);

        kafkaTopicProducer.topicProducer(gson.toJson(accrualEntity, AccrualEntity.class), accrualTopic);

    }

    public AccrualEntity convertBusinessObjectToAccrual(AccrualSchema ao, Integer partition,
            Long offset) {
        //Transform code goes here
        return ae;
    }
}

【问题讨论】:

  • 为什么要分离 Kafka 集群?
  • 两个不同的公司/系统。我们只需要整合一个主题的数据。
  • 这是我们无法改变的遗憾

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


【解决方案1】:

集群不支持 Exactly Once 语义;关键是在一个原子事务中生成记录并提交消耗的偏移量。

这在您的环境中是不可能的,因为一个事务不能跨越两个集群。

【讨论】:

  • 是否有任何解决方法来实现这一点?我得到的 1 个建议是使用 flink 和 2 阶段提交。
  • 有没有办法至少获得一次?
  • 我对flink不熟悉,所以无法评论。至少一次是 Spring for Apache Kafka 的默认行为。
  • 请您帮忙看看我发送的代码示例。我需要更改什么才能获得此设置的至少一次行为
  • 如果我们在同一个集群上,加里的另一件事是上面的代码只执行一次正确的实现?
猜你喜欢
  • 2017-01-14
  • 2020-10-04
  • 1970-01-01
  • 2020-05-12
  • 1970-01-01
  • 2017-06-09
  • 2018-11-21
  • 1970-01-01
  • 2017-11-20
相关资源
最近更新 更多