From 87a1b9988e69a049e57ecbb2dc4256bfb50b4fdf Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Fri, 2 Nov 2018 19:53:32 -0400 Subject: [PATCH] Log TX Exceptions - Log exceptions on commit/abort transaction - Don't attempt to abort if commit fails, otherwise we see `Abort failed` `org.apache.kafka.common.KafkaException: Cannot execute transactional method because we are in an error state` **cherry-pick to 2.1.x** * Rename variable to detect whether the commit failed * Polishing; use try/catch around commit --- .../core/DefaultKafkaProducerFactory.java | 2 ++ .../kafka/core/KafkaTemplate.java | 19 ++++++++++++++++++- 2 files changed, 20 insertions(+), 1 deletion(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java index cea814a0..8f6ade0a 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java @@ -466,6 +466,7 @@ public class DefaultKafkaProducerFactory implements ProducerFactory, this.delegate.commitTransaction(); } catch (RuntimeException e) { + logger.error("Commit failed", e); this.txFailed = true; throw e; } @@ -477,6 +478,7 @@ public class DefaultKafkaProducerFactory implements ProducerFactory, this.delegate.abortTransaction(); } catch (RuntimeException e) { + logger.error("Abort failed", e); this.txFailed = true; throw e; } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java index cdc8b72f..93530300 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java @@ -276,9 +276,17 @@ public class KafkaTemplate implements KafkaOperations { this.producers.set(producer); try { T result = callback.doInOperations(this); - producer.commitTransaction(); + try { + producer.commitTransaction(); + } + catch (Exception e) { + throw new SkipAbortException(e); + } return result; } + catch (SkipAbortException e) { + throw ((RuntimeException) e.getCause()); + } catch (Exception e) { producer.abortTransaction(); throw e; @@ -411,4 +419,13 @@ public class KafkaTemplate implements KafkaOperations { } } + @SuppressWarnings("serial") + private static final class SkipAbortException extends RuntimeException { + + SkipAbortException(Throwable cause) { + super(cause); + } + + } + }