【发布时间】: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