【问题标题】:How to use transactions in Kafka and how to use abortTransaction?如何在 Kafka 中使用事务以及如何使用 abortTransaction?
【发布时间】:2020-09-14 15:58:31
【问题描述】:

我是 kafka 新手,我使用 Kafka Producer Java api。 面对卡夫卡的这个问题, Kafka: Invalid transition attempted from state COMMITTING_TRANSACTION to state ABORTING_TRANSACTION.

人们写道,producer.abortTransaction() 只有在没有交易进行时才应该被调用.... 知道如何检查飞行中是否有交易吗?以及如何清除/阻止它们?

这是我的代码:

try { 
  producer.send(record, new Callback() { 
    @Override
    public void onCompletion(RecordMetadata recordMetadata, Exception e) { 
      if ( e != null){ 
        logger.info("Record was not sent due to kafka issue");
        throw new KafkaException("Record was not sent due to kafka issue");
      }
    }
  });
} catch (KafkaException e){
  producer.abortTransaction(); 
}

【问题讨论】:

  • try { producer.send(record, new Callback() { @Override public void onCompletion(RecordMetadata recordMetadata, Exception e) { if ( e != null){ logger.info("Record was not sent due to kafka issue"); throw new KafkaException("Record was not sent due to kafka issue"); } } } );} catch (KafkaException e){ producer.abortTransaction(); }
  • 我需要实现的是检测 kafka 何时停止,在这种情况下清除所有缓冲区,以便当 kafka 再次启动时,这些缓冲区中的记录不会出现在消费者端。我 initTransaction() 并提交它。我在正常执行发生时提交事务,并在出现问题时中止它(这里的问题是 kafka 将被停止)。

标签: apache-kafka kafka-producer-api


【解决方案1】:

我需要实现的是检测 kafka 何时停止,在这种情况下清除所有缓冲区,这样当 kafka 再次启动时,这些缓冲区中的记录就不会出现在消费者端。

在这种情况下,您通常会做的是应用KafkaProducer 的 Java 文档中描述的事务:

 Properties props = new Properties();
 props.put("bootstrap.servers", "localhost:9092");
 props.put("transactional.id", "my-transactional-id");
 Producer<String, String> producer = new KafkaProducer<>(props, new StringSerializer(), new StringSerializer());

 producer.initTransactions();

 try {
     producer.beginTransaction();
     for (int i = 0; i < 100; i++)
         producer.send(new ProducerRecord<>("my-topic", Integer.toString(i), Integer.toString(i)));
     producer.commitTransaction();
 } catch (ProducerFencedException | OutOfOrderSequenceException | AuthorizationException e) {
     // We can't recover from these exceptions, so our only option is to close the producer and exit.
     producer.close();
 } catch (KafkaException e) {
     // For all other exceptions, just abort the transaction and try again.
     producer.abortTransaction();
 }
 producer.close();

这样,如果 isolation.level 设置为 read_committed,则 100 条记录要么全部可见,要么都不可见。

您正在关闭 producer.close() 以处理不可恢复的异常,例如

  • ProducerFencedException:这个致命异常表明另一个具有相同transactional.id 的生产者已经启动。在任何给定时间只能有一个具有transactional.id 的生产者实例,并且要启动的最新实例会“隔离”先前的实例,以便它们无法再发出事务请求。遇到此异常时,必须关闭生产者实例。

  • OutOfOrderSequenceException:该异常表示broker从生产者那里收到了一个意外的序列号,这意味着数据可能已经丢失。如果生产者仅配置为幂等性(即,如果设置了enable.idempotence,但未配置transactional.id),则可以使用相同的生产者实例继续发送,但这样做有重新排序发送记录的风险。对于事务性生产者,这是一个致命错误,您应该关闭生产者。

  • AuthorizationException:[不言自明]

【讨论】:

  • 感谢您的帮助,您能否解释一下为什么将producer.close() 用于ProducerFencedExceptionproducer.abortTransaction() 用于KafkaException。另外,我需要添加回调到producersend() API吗?
猜你喜欢
  • 2020-12-25
  • 1970-01-01
  • 1970-01-01
  • 2022-01-27
  • 2012-02-20
  • 2019-11-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多