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();
+    }
+}

Reply via email to