This is an automated email from the ASF dual-hosted git repository.
oscerd pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new 0e12c3d0ac26 CAMEL-24782: camel-kafka - recover the transactional
producer instead of wedging the route (#27186)
0e12c3d0ac26 is described below
commit 0e12c3d0ac26586add52e6afb6a2b8c7d5a4b531
Author: Andrea Cosentino <[email protected]>
AuthorDate: Thu Oct 1 12:35:27 2026 +0200
CAMEL-24782: camel-kafka - recover the transactional producer instead of
wedging the route (#27186)
KafkaTransactionSynchronization.onDone handled a failed transaction in two
ways
that left the shared producer unusable and wedged the route:
- On a KafkaException it closed the shared producer with no recreation. The
KafkaProducer field kept pointing at the closed instance, so every later
exchange failed to begin or send a transaction.
- A failed commit (catch KafkaException) only recorded the exception,
leaving the
transaction open, so the next beginTransaction() failed.
The close heuristic was also too broad: any KafkaException closed the
producer,
even a downstream failure it could have recovered from with an abort.
Fix, following Kafka's documented transactional pattern:
- Classify fatal errors (ProducerFencedException,
OutOfOrderSequenceException,
AuthorizationException) - after which the producer cannot be reused - from
abortable ones.
- Fatal: close the producer and mark it for recreation; KafkaProducer
rebuilds it
(lazily, thread-safe, re-initialising transactions and re-setting the
transactional id) before the next transaction, so the route recovers.
- Abortable: abort the transaction so the shared producer stays usable; if
the
abort itself fails, close and recreate.
- A failed commit now goes through the same recovery (abort, or
close+recreate when
fatal) instead of leaving the transaction open.
Related: CAMEL-24780 (transaction begin), CAMEL-24783 (async batch
dispatch).
Co-authored-by: Claude Opus 4.8 <[email protected]>
Signed-off-by: Andrea Cosentino <[email protected]>
---
.../camel/component/kafka/KafkaProducer.java | 37 +++++-
.../kafka/KafkaTransactionSynchronization.java | 85 +++++++++---
.../camel/component/kafka/KafkaProducerTest.java | 33 +++++
.../kafka/KafkaTransactionSynchronizationTest.java | 143 +++++++++++++++++++++
4 files changed, 280 insertions(+), 18 deletions(-)
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java
index 43ca34c2afea..30c7fc744dfb 100755
---
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java
+++
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java
@@ -76,6 +76,8 @@ public class KafkaProducer extends DefaultAsyncProducer
implements RouteIdAware
private ExecutorService workerPool;
private boolean shutdownWorkerPool;
private volatile boolean closeKafkaProducer;
+ // set when a fatal transactional error closed the shared producer, so the
next transaction rebuilds it (CAMEL-24782)
+ private volatile boolean producerClosedForRecreation;
private final String endpointTopic;
private final Integer configPartitionKey;
private final String configKey;
@@ -523,19 +525,52 @@ public class KafkaProducer extends DefaultAsyncProducer
implements RouteIdAware
UnitOfWork uow = exchange.getUnitOfWork();
if (!uow.isTransactedBy(transactionId)) {
+ // A prior transaction may have hit a fatal error that closed the
shared producer; rebuild it before
+ // starting a new transaction so the route recovers instead of
staying wedged on a dead producer
+ // (CAMEL-24782).
+ recreateProducerIfClosed();
LOG.debug("Starting kafka transaction {} with exchange {}",
transactionId, exchange.getExchangeId());
// Begin the broker transaction first, then mark the unit of work
and register the
// synchronization. This way a failure in beginTransaction() does
not leave the unit of work
// flagged as transacted without a synchronization to commit or
roll it back (CAMEL-24780).
kafkaProducer.beginTransaction();
uow.beginTransactedBy(transactionId);
- uow.addSynchronization(new
KafkaTransactionSynchronization(transactionId, kafkaProducer));
+ uow.addSynchronization(new
KafkaTransactionSynchronization(transactionId, kafkaProducer, this));
} else {
LOG.debug("Using existing kafka transaction {} with exchange {}.",
transactionId, exchange.getExchangeId());
}
}
+ /**
+ * Rebuilds the shared producer when a previous transaction hit a fatal
error and closed it. Lazy (done here on the
+ * next transactional send rather than from the Kafka callback thread) and
guarded so concurrent callers rebuild it
+ * once. Re-initialises transactions so the new producer is ready to use
(CAMEL-24782).
+ */
+ private synchronized void recreateProducerIfClosed() {
+ if (producerClosedForRecreation) {
+ LOG.warn("Recreating kafka producer for transaction {} after a
fatal error closed the previous one",
+ transactionId);
+ Properties props = getProps();
+ // the transactional id may have been generated from the
endpoint/route id in doStart and is not in the
+ // configuration, so set it explicitly so the rebuilt producer
matches the one that was closed
+ if (transactionId != null) {
+ props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG,
transactionId);
+ }
+ createProducer(props);
+ kafkaProducer.initTransactions();
+ producerClosedForRecreation = false;
+ }
+ }
+
+ /**
+ * Signals that the shared producer was closed because a transaction hit a
fatal error, so it must be rebuilt before
+ * the next transactional send. Called by {@link
KafkaTransactionSynchronization} (CAMEL-24782).
+ */
+ void markProducerClosedForRecreation() {
+ producerClosedForRecreation = true;
+ }
+
@Override
public String getRouteId() {
return routeId;
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaTransactionSynchronization.java
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaTransactionSynchronization.java
index 2525211f7878..9758d0af5499 100644
---
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaTransactionSynchronization.java
+++
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaTransactionSynchronization.java
@@ -19,7 +19,9 @@ package org.apache.camel.component.kafka;
import org.apache.camel.Exchange;
import org.apache.camel.support.SynchronizationAdapter;
import org.apache.kafka.clients.producer.Producer;
-import org.apache.kafka.common.KafkaException;
+import org.apache.kafka.common.errors.AuthorizationException;
+import org.apache.kafka.common.errors.OutOfOrderSequenceException;
+import org.apache.kafka.common.errors.ProducerFencedException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -27,37 +29,86 @@ class KafkaTransactionSynchronization extends
SynchronizationAdapter {
private static final Logger LOG =
LoggerFactory.getLogger(KafkaTransactionSynchronization.class);
private final String transactionId;
private final Producer kafkaProducer;
+ private final KafkaProducer owner;
- public KafkaTransactionSynchronization(String transactionId, Producer
kafkaProducer) {
+ public KafkaTransactionSynchronization(String transactionId, Producer
kafkaProducer, KafkaProducer owner) {
this.transactionId = transactionId;
this.kafkaProducer = kafkaProducer;
+ this.owner = owner;
}
@Override
public void onDone(Exchange exchange) {
try {
if (exchange.getException() != null || exchange.isRollbackOnly()) {
- if (exchange.getException() instanceof KafkaException) {
- LOG.warn("Catch {} and will close kafka producer with
transaction {} ", exchange.getException(),
- transactionId);
- kafkaProducer.close();
- } else {
- LOG.warn("Abort kafka transaction {} with exchange {}",
transactionId, exchange.getExchangeId());
- kafkaProducer.abortTransaction();
- }
+ rollback(exchange);
} else {
- LOG.debug("Commit kafka transaction {} with exchange {}",
transactionId, exchange.getExchangeId());
- kafkaProducer.commitTransaction();
+ commit(exchange);
}
- } catch (KafkaException e) {
- exchange.setException(e);
+ } finally {
+ exchange.getUnitOfWork().endTransactedBy(transactionId);
+ }
+ }
+
+ private void commit(Exchange exchange) {
+ try {
+ LOG.debug("Commit kafka transaction {} with exchange {}",
transactionId, exchange.getExchangeId());
+ kafkaProducer.commitTransaction();
} catch (Exception e) {
+ // The commit failed. Record it and return the producer to a
usable state - abort the still-open
+ // transaction, or close and rebuild it when the error is fatal.
The previous code only recorded the
+ // exception, leaving the transaction open so the next
beginTransaction() failed and wedged the route
+ // (CAMEL-24782).
exchange.setException(e);
- LOG.warn("Abort kafka transaction {} with exchange {} due to {} ",
transactionId, exchange.getExchangeId(),
- e.getMessage(), e);
+ recover(e, "commit");
+ }
+ }
+
+ private void rollback(Exchange exchange) {
+ // The routing already failed; abort the transaction so the shared
producer stays usable, or - when the
+ // failure is a fatal transactional error that the producer cannot
recover from - close and rebuild it.
+ recover(exchange.getException(), "rollback");
+ }
+
+ /**
+ * Returns the shared producer to a usable state after a failed or
rolled-back transaction. A fatal error (the
+ * producer can no longer be used - fenced, out-of-order sequence,
authorization) closes it and asks the owner to
+ * rebuild it before the next transaction; anything else aborts the
current transaction. If the abort itself fails
+ * the producer is also closed and rebuilt, so the route is never left
wedged on a dead producer (CAMEL-24782).
+ */
+ private void recover(Throwable cause, String phase) {
+ if (isFatal(cause)) {
+ LOG.warn("Closing kafka producer for transaction {} after a fatal
error during {}: {}",
+ transactionId, phase, cause.toString());
+ closeForRecreation();
+ return;
+ }
+ try {
+ LOG.warn("Abort kafka transaction {} during {}", transactionId,
phase);
kafkaProducer.abortTransaction();
+ } catch (Exception abortFailure) {
+ LOG.warn("Closing kafka producer for transaction {} after a failed
abort during {}: {}",
+ transactionId, phase, abortFailure.toString(),
abortFailure);
+ closeForRecreation();
+ }
+ }
+
+ private void closeForRecreation() {
+ try {
+ kafkaProducer.close();
} finally {
- exchange.getUnitOfWork().endTransactedBy(transactionId);
+ // even if close() throws, the producer is unusable - make sure
the next transaction rebuilds it
+ owner.markProducerClosedForRecreation();
}
}
+
+ /**
+ * A fatal transactional error leaves the producer permanently unusable,
so it must be closed rather than aborted
+ * (an abort would fail too). These are the errors the Kafka transactional
producer documents as fatal.
+ */
+ private static boolean isFatal(Throwable cause) {
+ return cause instanceof ProducerFencedException
+ || cause instanceof OutOfOrderSequenceException
+ || cause instanceof AuthorizationException;
+ }
}
diff --git
a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaProducerTest.java
b/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaProducerTest.java
index b195eae3235a..ffa317f46e1d 100755
---
a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaProducerTest.java
+++
b/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaProducerTest.java
@@ -59,6 +59,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
@@ -221,6 +222,38 @@ public class KafkaProducerTest {
Mockito.verify(uow, Mockito.never()).addSynchronization(any());
}
+ @Test
+ void transactionalProducerIsRecreatedAfterAFatalError() throws Exception {
+ // CAMEL-24782: a fatal transactional error closes the shared
producer; the next transaction must rebuild it
+ // through the client factory instead of leaving the route wedged on a
dead (closed) producer.
+ endpoint.getConfiguration().setTransactionalId("test-tx");
+ endpoint.getConfiguration().setTopic("sometopic");
+ producer.doStart();
+
+ Producer recreated = Mockito.mock(Producer.class);
+ KafkaClientFactory factory = Mockito.mock(KafkaClientFactory.class);
+
Mockito.when(factory.getProducer(any(Properties.class))).thenReturn(recreated);
+ endpoint.setKafkaClientFactory(factory);
+
+ // a previous transaction hit a fatal error and closed the shared
producer
+ producer.markProducerClosedForRecreation();
+
+ UnitOfWork uow = Mockito.mock(UnitOfWork.class);
+ Mockito.when(uow.isTransactedBy(any())).thenReturn(false);
+ Mockito.when(exchange.getUnitOfWork()).thenReturn(uow);
+ Mockito.when(exchange.getIn()).thenReturn(in);
+ Mockito.when(exchange.getMessage()).thenReturn(in);
+ in.setHeader(KafkaConstants.PARTITION_KEY, 4);
+
+ producer.process(exchange, callback);
+
+ // the shared producer was rebuilt via the factory, re-initialised for
transactions, and used for this exchange
+ Mockito.verify(factory).getProducer(any(Properties.class));
+ Mockito.verify(recreated).initTransactions();
+ Mockito.verify(recreated).beginTransaction();
+ assertSame(recreated, producer.getKafkaProducer());
+ }
+
@Test
public void processSendsMessageWithTopicHeaderAndNoTopicInEndPoint()
throws Exception {
endpoint.getConfiguration().setTopic(null);
diff --git
a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaTransactionSynchronizationTest.java
b/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaTransactionSynchronizationTest.java
new file mode 100644
index 000000000000..ef42daf30a29
--- /dev/null
+++
b/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaTransactionSynchronizationTest.java
@@ -0,0 +1,143 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.kafka;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.spi.UnitOfWork;
+import org.apache.kafka.clients.producer.Producer;
+import org.apache.kafka.common.KafkaException;
+import org.apache.kafka.common.errors.OutOfOrderSequenceException;
+import org.apache.kafka.common.errors.ProducerFencedException;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Unit tests for {@link KafkaTransactionSynchronization#onDone}: a fatal
transactional error must close the shared
+ * producer and have the owner rebuild it (otherwise the route wedges on a
dead producer), an abortable error must abort
+ * the transaction (so the producer stays usable), and a failed commit must
not be left silently open (CAMEL-24782).
+ */
+class KafkaTransactionSynchronizationTest {
+
+ private static final String TX_ID = "tx-1";
+
+ @SuppressWarnings("rawtypes")
+ private Producer kafkaProducer;
+ private KafkaProducer owner;
+ private Exchange exchange;
+ private KafkaTransactionSynchronization sync;
+
+ @BeforeEach
+ void setUp() {
+ kafkaProducer = Mockito.mock(Producer.class);
+ owner = Mockito.mock(KafkaProducer.class);
+ exchange = Mockito.mock(Exchange.class);
+
when(exchange.getUnitOfWork()).thenReturn(Mockito.mock(UnitOfWork.class));
+ sync = new KafkaTransactionSynchronization(TX_ID, kafkaProducer,
owner);
+ }
+
+ @Test
+ void fatalExceptionClosesAndMarksForRecreation() {
+ when(exchange.getException()).thenReturn(new
ProducerFencedException("fenced"));
+
+ sync.onDone(exchange);
+
+ verify(kafkaProducer).close();
+ verify(owner).markProducerClosedForRecreation();
+ verify(kafkaProducer, never()).abortTransaction();
+ }
+
+ @Test
+ void nonFatalExceptionAbortsTransaction() {
+ when(exchange.getException()).thenReturn(new
RuntimeException("downstream failed"));
+
+ sync.onDone(exchange);
+
+ verify(kafkaProducer).abortTransaction();
+ verify(kafkaProducer, never()).close();
+ verify(owner, never()).markProducerClosedForRecreation();
+ }
+
+ @Test
+ void rollbackOnlyAbortsTransaction() {
+ when(exchange.getException()).thenReturn(null);
+ when(exchange.isRollbackOnly()).thenReturn(true);
+
+ sync.onDone(exchange);
+
+ verify(kafkaProducer).abortTransaction();
+ verify(kafkaProducer, never()).commitTransaction();
+ verify(kafkaProducer, never()).close();
+ }
+
+ @Test
+ void commitSuccessCommitsTransaction() {
+ when(exchange.getException()).thenReturn(null);
+ when(exchange.isRollbackOnly()).thenReturn(false);
+
+ sync.onDone(exchange);
+
+ verify(kafkaProducer).commitTransaction();
+ verify(kafkaProducer, never()).abortTransaction();
+ verify(kafkaProducer, never()).close();
+ }
+
+ @Test
+ void commitFatalErrorClosesAndMarks() {
+ when(exchange.getException()).thenReturn(null);
+ when(exchange.isRollbackOnly()).thenReturn(false);
+ Mockito.doThrow(new
OutOfOrderSequenceException("gap")).when(kafkaProducer).commitTransaction();
+
+ sync.onDone(exchange);
+
+ verify(kafkaProducer).close();
+ verify(owner).markProducerClosedForRecreation();
+ verify(kafkaProducer, never()).abortTransaction();
+
verify(exchange).setException(Mockito.isA(OutOfOrderSequenceException.class));
+ }
+
+ @Test
+ void commitAbortableErrorAbortsTransaction() {
+ when(exchange.getException()).thenReturn(null);
+ when(exchange.isRollbackOnly()).thenReturn(false);
+ Mockito.doThrow(new KafkaException("commit failed,
abortable")).when(kafkaProducer).commitTransaction();
+
+ sync.onDone(exchange);
+
+ // the still-open transaction is aborted so the next
beginTransaction() can succeed, and the producer is kept
+ verify(kafkaProducer).abortTransaction();
+ verify(kafkaProducer, never()).close();
+ verify(owner, never()).markProducerClosedForRecreation();
+ verify(exchange).setException(Mockito.isA(KafkaException.class));
+ }
+
+ @Test
+ void failedAbortDuringRollbackClosesAndMarks() {
+ when(exchange.getException()).thenReturn(new
RuntimeException("downstream failed"));
+ Mockito.doThrow(new KafkaException("abort
failed")).when(kafkaProducer).abortTransaction();
+
+ sync.onDone(exchange);
+
+ // abort failed, so the producer is unusable and must be rebuilt
rather than left wedged
+ verify(kafkaProducer).close();
+ verify(owner).markProducerClosedForRecreation();
+ }
+}