This is an automated email from the ASF dual-hosted git repository.

chia7712 pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git


The following commit(s) were added to refs/heads/trunk by this push:
     new 8411977aaf9 MINOR: Move transaction coordinator leaf classes from 
Scala to Java (#22576)
8411977aaf9 is described below

commit 8411977aaf96c3949aebc0f18a095ab6d8560cae
Author: Sanskar Jhajharia <[email protected]>
AuthorDate: Tue Jun 23 09:22:41 2026 +0530

    MINOR: Move transaction coordinator leaf classes from Scala to Java (#22576)
    
    ### Description
    
    As part of the ongoing effort to migrate `kafka.coordinator.transaction`
    from Scala (core) to Java (the transaction-coordinator module), this PR
    moves the self-contained "leaf" classes - pure data holders and a
    stateful helper with no ReplicaManager/KafkaConfig coupling. These can
    move independently while the remaining Scala classes continue to compile
    against them, since core depends on transaction-coordinator.
    
    ### Changes
    - The three data classes become Java records. PendingCompleteTxn keeps a
    custom toString to preserve the original log output format.
    - TxnMarkerQueue becomes a plain Java class: kafka.utils.Logging →
    SLF4J, Scala getOrElseUpdate → computeIfAbsent, and Option → Optional on
    the public API.
    - Scala call sites updated
    - Import repointed in KafkaApis.scala and the affected test files.
    
    Reviewers: Chia-Ping Tsai <[email protected]>
---
 .../transaction/TransactionCoordinator.scala       |  10 +--
 .../TransactionMarkerChannelManager.scala          |  83 +++--------------
 ...TransactionMarkerRequestCompletionHandler.scala |   1 +
 core/src/main/scala/kafka/server/KafkaApis.scala   |   3 +-
 .../TransactionCoordinatorConcurrencyTest.scala    |   2 +-
 .../transaction/TransactionCoordinatorTest.scala   |  70 +++++++--------
 ...sactionMarkerRequestCompletionHandlerTest.scala |   6 +-
 .../scala/unit/kafka/server/KafkaApisTest.scala    |   8 +-
 .../transaction/InitProducerIdResult.java          |  22 +++++
 .../transaction/PendingCompleteTxn.java            |  25 ++++++
 .../PendingCompleteTxnAndMarkerEntry.java          |  25 ++++++
 .../coordinator/transaction/TxnMarkerQueue.java    | 100 +++++++++++++++++++++
 12 files changed, 232 insertions(+), 123 deletions(-)

diff --git 
a/core/src/main/scala/kafka/coordinator/transaction/TransactionCoordinator.scala
 
b/core/src/main/scala/kafka/coordinator/transaction/TransactionCoordinator.scala
index 771c6727072..36a2fa034cf 100644
--- 
a/core/src/main/scala/kafka/coordinator/transaction/TransactionCoordinator.scala
+++ 
b/core/src/main/scala/kafka/coordinator/transaction/TransactionCoordinator.scala
@@ -29,7 +29,7 @@ import org.apache.kafka.common.record.internal.RecordBatch
 import org.apache.kafka.common.requests.{AddPartitionsToTxnResponse, 
TransactionResult}
 import org.apache.kafka.common.utils.{ProducerIdAndEpoch, Time}
 import org.apache.kafka.common.utils.internals.LogContext
-import org.apache.kafka.coordinator.transaction.{ProducerIdManager, 
TransactionConfig, TransactionLogConfig, TransactionMetadata, TransactionState, 
TransactionStateManagerConfig, TransactionalIdAndProducerIdEpoch, 
TxnTransitMetadata}
+import org.apache.kafka.coordinator.transaction.{InitProducerIdResult, 
ProducerIdManager, TransactionConfig, TransactionLogConfig, 
TransactionMetadata, TransactionState, TransactionStateManagerConfig, 
TransactionalIdAndProducerIdEpoch, TxnTransitMetadata}
 import org.apache.kafka.metadata.MetadataCache
 import org.apache.kafka.server.common.{RequestLocal, TransactionVersion}
 import org.apache.kafka.server.record.BrokerCompressionType
@@ -76,11 +76,11 @@ object TransactionCoordinator {
   }
 
   private def initTransactionError(error: Errors): InitProducerIdResult = {
-    InitProducerIdResult(RecordBatch.NO_PRODUCER_ID, 
RecordBatch.NO_PRODUCER_EPOCH, error)
+    new InitProducerIdResult(RecordBatch.NO_PRODUCER_ID, 
RecordBatch.NO_PRODUCER_EPOCH, error)
   }
 
   private def initTransactionMetadata(txnMetadata: TxnTransitMetadata): 
InitProducerIdResult = {
-    InitProducerIdResult(txnMetadata.producerId, txnMetadata.producerEpoch, 
Errors.NONE)
+    new InitProducerIdResult(txnMetadata.producerId, 
txnMetadata.producerEpoch, Errors.NONE)
   }
 }
 
@@ -126,7 +126,7 @@ class TransactionCoordinator(txnConfig: TransactionConfig,
       // if the transactional id is null, then always blindly accept the 
request
       // and return a new producerId from the producerId manager
       try {
-        
responseCallback(InitProducerIdResult(producerIdManager.generateProducerId(), 
producerEpoch = 0, Errors.NONE))
+        responseCallback(new 
InitProducerIdResult(producerIdManager.generateProducerId(), 0, Errors.NONE))
       } catch {
         case e: Exception => 
responseCallback(initTransactionError(Errors.forException(e)))
       }
@@ -1113,5 +1113,3 @@ class TransactionCoordinator(txnConfig: TransactionConfig,
     info("Shutdown complete.")
   }
 }
-
-case class InitProducerIdResult(producerId: Long, producerEpoch: Short, error: 
Errors)
diff --git 
a/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerChannelManager.scala
 
b/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerChannelManager.scala
index c519c648d75..d270da31b84 100644
--- 
a/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerChannelManager.scala
+++ 
b/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerChannelManager.scala
@@ -20,7 +20,7 @@ package kafka.coordinator.transaction
 import 
kafka.coordinator.transaction.TransactionMarkerChannelManager.{LogAppendRetryQueueSizeMetricName,
 MetricNames, UnknownDestinationQueueSizeMetricName}
 
 import java.util
-import java.util.concurrent.{BlockingQueue, ConcurrentHashMap, 
LinkedBlockingQueue}
+import java.util.concurrent.{ConcurrentHashMap, LinkedBlockingQueue}
 import kafka.server.KafkaConfig
 import kafka.utils.Logging
 import org.apache.kafka.clients._
@@ -33,7 +33,7 @@ import org.apache.kafka.common.security.JaasContext
 import org.apache.kafka.common.utils.Time
 import org.apache.kafka.common.utils.internals.LogContext
 import org.apache.kafka.common.{Node, Reconfigurable, TopicPartition}
-import org.apache.kafka.coordinator.transaction.{TransactionMetadata, 
TxnTransitMetadata}
+import org.apache.kafka.coordinator.transaction.{PendingCompleteTxn, 
PendingCompleteTxnAndMarkerEntry, TransactionMetadata, TxnMarkerQueue, 
TxnTransitMetadata}
 import org.apache.kafka.metadata.MetadataCache
 import org.apache.kafka.server.common.RequestLocal
 import org.apache.kafka.server.metrics.KafkaMetricsGroup
@@ -113,53 +113,6 @@ object TransactionMarkerChannelManager {
 
 }
 
-class TxnMarkerQueue(@volatile var destination: Node) extends Logging {
-
-  // keep track of the requests per txn topic partition so we can easily clear 
the queue
-  // during partition emigration
-  private val markersPerTxnTopicPartition = new ConcurrentHashMap[Int, 
BlockingQueue[PendingCompleteTxnAndMarkerEntry]]().asScala
-
-  def removeMarkersForTxnTopicPartition(partition: Int): 
Option[BlockingQueue[PendingCompleteTxnAndMarkerEntry]] = {
-    markersPerTxnTopicPartition.remove(partition)
-  }
-
-  def addMarkers(txnTopicPartition: Int, pendingCompleteTxnAndMarker: 
PendingCompleteTxnAndMarkerEntry): Unit = {
-    val queue = markersPerTxnTopicPartition.getOrElseUpdate(txnTopicPartition, 
{
-      // Note that this may get called more than once if threads have a close 
race while adding new queue.
-      info(s"Creating new marker queue for txn partition $txnTopicPartition to 
destination broker ${destination.id}")
-      new LinkedBlockingQueue[PendingCompleteTxnAndMarkerEntry]()
-    })
-    queue.add(pendingCompleteTxnAndMarker)
-
-    if (markersPerTxnTopicPartition.get(txnTopicPartition).orNull != queue) {
-      // This could happen if the queue got removed concurrently.
-      // Note that it could create an unexpected state when the queue is 
removed from
-      // removeMarkersForTxnTopicPartition, we could have:
-      //
-      // 1. [addMarkers] Retrieve queue.
-      // 2. [removeMarkersForTxnTopicPartition] Remove queue.
-      // 3. [removeMarkersForTxnTopicPartition] Iterate over queue, but not 
removeMarkersForTxn because queue is empty.
-      // 4. [addMarkers] Add markers to the queue.
-      //
-      // Now we've effectively removed the markers while 
transactionsWithPendingMarkers has an entry.
-      //
-      // While this could lead to an orphan entry in 
transactionsWithPendingMarkers, sending new markers
-      // will fix the state, so it shouldn't impact the state machine 
operation.
-      warn(s"Added $pendingCompleteTxnAndMarker to dead queue for txn 
partition $txnTopicPartition to destination broker ${destination.id}")
-    }
-  }
-
-  def forEachTxnTopicPartition[B](f:(Int, 
BlockingQueue[PendingCompleteTxnAndMarkerEntry]) => B): Unit =
-    markersPerTxnTopicPartition.foreachEntry { (partition, queue) =>
-      if (!queue.isEmpty) f(partition, queue)
-    }
-
-  def totalNumMarkers: Int = markersPerTxnTopicPartition.values.foldLeft(0) { 
_ + _.size }
-
-  // visible for testing
-  def totalNumMarkers(txnTopicPartition: Int): Int = 
markersPerTxnTopicPartition.get(txnTopicPartition).fold(0)(_.size)
-}
-
 class TransactionMarkerChannelManager(
   config: KafkaConfig,
   metadataCache: MetadataCache,
@@ -219,7 +172,7 @@ class TransactionMarkerChannelManager(
       info(s"Creating new marker queue map to destination broker $brokerId")
       new TxnMarkerQueue(broker)
     })
-    brokerRequestQueue.destination = broker
+    brokerRequestQueue.setDestination(broker)
     brokerRequestQueue.addMarkers(txnTopicPartition, 
pendingCompleteTxnAndMarker)
 
     trace(s"Added marker ${pendingCompleteTxnAndMarker.txnMarkerEntry} for 
transactional id" +
@@ -238,7 +191,7 @@ class TransactionMarkerChannelManager(
   override def generateRequests(): 
util.Collection[RequestAndCompletionHandler] = {
     retryLogAppends()
     val pendingCompleteTxnAndMarkerEntries = new 
util.ArrayList[PendingCompleteTxnAndMarkerEntry]()
-    markersQueueForUnknownBroker.forEachTxnTopicPartition { case (_, queue) =>
+    markersQueueForUnknownBroker.forEachTxnTopicPartition { (_, queue) =>
       queue.drainTo(pendingCompleteTxnAndMarkerEntries)
     }
 
@@ -255,7 +208,7 @@ class TransactionMarkerChannelManager(
     val currentTimeMs = time.milliseconds()
     markersQueuePerBroker.values.map { brokerRequestQueue =>
       val pendingCompleteTxnAndMarkerEntries = new 
util.ArrayList[PendingCompleteTxnAndMarkerEntry]()
-      brokerRequestQueue.forEachTxnTopicPartition { case (_, queue) =>
+      brokerRequestQueue.forEachTxnTopicPartition { (_, queue) =>
         queue.drainTo(pendingCompleteTxnAndMarkerEntries)
       }
       (brokerRequestQueue.destination, pendingCompleteTxnAndMarkerEntries)
@@ -299,7 +252,7 @@ class TransactionMarkerChannelManager(
         if (epochAndMetadata.coordinatorEpoch == coordinatorEpoch) {
           debug(s"Sending $transactionalId's transaction markers for 
$txnMetadata with " +
             s"coordinator epoch $coordinatorEpoch succeeded, trying to append 
complete transaction log now")
-          tryAppendToLog(PendingCompleteTxn(transactionalId, coordinatorEpoch, 
txnMetadata, newMetadata))
+          tryAppendToLog(new PendingCompleteTxn(transactionalId, 
coordinatorEpoch, txnMetadata, newMetadata))
         } else {
           info(s"The cached metadata $txnMetadata has changed to 
$epochAndMetadata after " +
             s"completed sending the markers with coordinator epoch 
$coordinatorEpoch; abort " +
@@ -319,7 +272,7 @@ class TransactionMarkerChannelManager(
                           txnMetadata: TransactionMetadata,
                           newMetadata: TxnTransitMetadata): Unit = {
     val transactionalId = txnMetadata.transactionalId
-    val pendingCompleteTxn = PendingCompleteTxn(
+    val pendingCompleteTxn = new PendingCompleteTxn(
       transactionalId,
       coordinatorEpoch,
       txnMetadata,
@@ -408,7 +361,7 @@ class TransactionMarkerChannelManager(
       broker match {
         case Some(brokerNode) =>
           val marker = new TxnMarkerEntry(producerId, producerEpoch, 
coordinatorEpoch, result, topicPartitions.toList.asJava, transactionVersion)
-          val pendingCompleteTxnAndMarker = 
PendingCompleteTxnAndMarkerEntry(pendingCompleteTxn, marker)
+          val pendingCompleteTxnAndMarker = new 
PendingCompleteTxnAndMarkerEntry(pendingCompleteTxn, marker)
 
           if (brokerNode == Node.noNode) {
             // if the leader of the partition is known but node not available, 
put it into an unknown broker queue
@@ -458,7 +411,7 @@ class TransactionMarkerChannelManager(
   }
 
   def removeMarkersForTxnTopicPartition(txnTopicPartitionId: Int): Unit = {
-    
markersQueueForUnknownBroker.removeMarkersForTxnTopicPartition(txnTopicPartitionId).foreach
 { queue =>
+    
markersQueueForUnknownBroker.removeMarkersForTxnTopicPartition(txnTopicPartitionId).ifPresent
 { queue =>
       for (entry <- queue.asScala) {
         info(s"Removing $entry for txn partition $txnTopicPartitionId to 
destination broker -1")
         removeMarkersForTxn(entry.pendingCompleteTxn)
@@ -466,7 +419,7 @@ class TransactionMarkerChannelManager(
     }
 
     markersQueuePerBroker.foreach { case(brokerId, brokerQueue) =>
-      
brokerQueue.removeMarkersForTxnTopicPartition(txnTopicPartitionId).foreach { 
queue =>
+      
brokerQueue.removeMarkersForTxnTopicPartition(txnTopicPartitionId).ifPresent { 
queue =>
         for (entry <- queue.asScala) {
           info(s"Removing $entry for txn partition $txnTopicPartitionId to 
destination broker $brokerId")
           removeMarkersForTxn(entry.pendingCompleteTxn)
@@ -486,19 +439,3 @@ class TransactionMarkerChannelManager(
     }
   }
 }
-
-case class PendingCompleteTxn(transactionalId: String,
-                              coordinatorEpoch: Int,
-                              txnMetadata: TransactionMetadata,
-                              newMetadata: TxnTransitMetadata) {
-
-  override def toString: String = {
-    "PendingCompleteTxn(" +
-      s"transactionalId=$transactionalId, " +
-      s"coordinatorEpoch=$coordinatorEpoch, " +
-      s"txnMetadata=$txnMetadata, " +
-      s"newMetadata=$newMetadata)"
-  }
-}
-
-case class PendingCompleteTxnAndMarkerEntry(pendingCompleteTxn: 
PendingCompleteTxn, txnMarkerEntry: TxnMarkerEntry)
diff --git 
a/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerRequestCompletionHandler.scala
 
b/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerRequestCompletionHandler.scala
index 63990fda985..fe48251819d 100644
--- 
a/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerRequestCompletionHandler.scala
+++ 
b/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerRequestCompletionHandler.scala
@@ -22,6 +22,7 @@ import org.apache.kafka.clients.{ClientResponse, 
RequestCompletionHandler}
 import org.apache.kafka.common.TopicPartition
 import org.apache.kafka.common.protocol.Errors
 import org.apache.kafka.common.requests.WriteTxnMarkersResponse
+import 
org.apache.kafka.coordinator.transaction.PendingCompleteTxnAndMarkerEntry
 
 import scala.collection.mutable
 import scala.jdk.CollectionConverters._
diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala 
b/core/src/main/scala/kafka/server/KafkaApis.scala
index 8e4a695baa4..3bf0861d847 100644
--- a/core/src/main/scala/kafka/server/KafkaApis.scala
+++ b/core/src/main/scala/kafka/server/KafkaApis.scala
@@ -17,7 +17,7 @@
 
 package kafka.server
 
-import kafka.coordinator.transaction.{InitProducerIdResult, 
TransactionCoordinator}
+import kafka.coordinator.transaction.TransactionCoordinator
 import kafka.network.RequestChannel
 import kafka.server.QuotaFactory.{QuotaManagers, UNBOUNDED_QUOTA}
 import kafka.server.handlers.DescribeTopicPartitionsRequestHandler
@@ -60,6 +60,7 @@ import org.apache.kafka.common.{Node, TopicIdPartition, 
TopicPartition, Uuid}
 import org.apache.kafka.coordinator.group.modern.share.ShareGroupConfigProvider
 import org.apache.kafka.coordinator.group.{Group, GroupConfig, 
GroupConfigManager, GroupCoordinator}
 import org.apache.kafka.coordinator.share.ShareCoordinator
+import org.apache.kafka.coordinator.transaction.InitProducerIdResult
 import org.apache.kafka.metadata.{ConfigRepository, MetadataCache}
 import org.apache.kafka.network.Request
 import org.apache.kafka.security.DelegationTokenManager
diff --git 
a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionCoordinatorConcurrencyTest.scala
 
b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionCoordinatorConcurrencyTest.scala
index 7cbf7ddfaee..78044bf9c97 100644
--- 
a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionCoordinatorConcurrencyTest.scala
+++ 
b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionCoordinatorConcurrencyTest.scala
@@ -38,7 +38,7 @@ import org.apache.kafka.common.requests._
 import org.apache.kafka.common.utils.{MockTime, ProducerIdAndEpoch}
 import org.apache.kafka.common.utils.internals.LogContext
 import org.apache.kafka.common.{Node, TopicPartition, Uuid}
-import org.apache.kafka.coordinator.transaction.{ProducerIdManager, 
TransactionConfig, TransactionLog, TransactionMetadata, TransactionState}
+import org.apache.kafka.coordinator.transaction.{InitProducerIdResult, 
ProducerIdManager, TransactionConfig, TransactionLog, TransactionMetadata, 
TransactionState}
 import org.apache.kafka.metadata.MetadataCache
 import org.apache.kafka.server.common.{FinalizedFeatures, MetadataVersion, 
RequestLocal, TransactionVersion}
 import org.apache.kafka.server.storage.log.FetchIsolation
diff --git 
a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionCoordinatorTest.scala
 
b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionCoordinatorTest.scala
index c8094146b75..91e5a4f7ce5 100644
--- 
a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionCoordinatorTest.scala
+++ 
b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionCoordinatorTest.scala
@@ -23,7 +23,7 @@ import org.apache.kafka.common.record.internal.RecordBatch
 import org.apache.kafka.common.requests.{AddPartitionsToTxnResponse, 
TransactionResult}
 import org.apache.kafka.common.utils.{MockTime, ProducerIdAndEpoch}
 import org.apache.kafka.common.utils.internals.LogContext
-import 
org.apache.kafka.coordinator.transaction.{CoordinatorEpochAndTxnMetadata, 
ProducerIdManager, TransactionConfig, TransactionMetadata, TransactionState, 
TransactionStateManagerConfig, TransactionalIdAndProducerIdEpoch, 
TxnTransitMetadata}
+import 
org.apache.kafka.coordinator.transaction.{CoordinatorEpochAndTxnMetadata, 
InitProducerIdResult, ProducerIdManager, TransactionConfig, 
TransactionMetadata, TransactionState, TransactionStateManagerConfig, 
TransactionalIdAndProducerIdEpoch, TxnTransitMetadata}
 import org.apache.kafka.server.common.{RequestLocal, TransactionVersion}
 import org.apache.kafka.server.common.TransactionVersion.{TV_0, TV_2}
 import org.apache.kafka.server.util.MockScheduler
@@ -94,10 +94,10 @@ class TransactionCoordinatorTest {
 
     coordinator.handleInitProducerId("", txnTimeoutMs, enableTwoPCFlag = false,
       keepPreparedTxn = false, None, initProducerIdMockCallback)
-    assertEquals(InitProducerIdResult(-1L, -1, Errors.INVALID_REQUEST), result)
+    assertEquals(new InitProducerIdResult(-1L, -1, Errors.INVALID_REQUEST), 
result)
     coordinator.handleInitProducerId("", txnTimeoutMs, enableTwoPCFlag = false,
       keepPreparedTxn = false, None, initProducerIdMockCallback)
-    assertEquals(InitProducerIdResult(-1L, -1, Errors.INVALID_REQUEST), result)
+    assertEquals(new InitProducerIdResult(-1L, -1, Errors.INVALID_REQUEST), 
result)
   }
 
   @Test
@@ -106,7 +106,7 @@ class TransactionCoordinatorTest {
 
     coordinator.handleInitProducerId("", txnTimeoutMs, enableTwoPCFlag = false,
       keepPreparedTxn = true, None, initProducerIdMockCallback)
-    assertEquals(InitProducerIdResult(-1L, -1, Errors.INVALID_REQUEST), result)
+    assertEquals(new InitProducerIdResult(-1L, -1, Errors.INVALID_REQUEST), 
result)
   }
 
   @Test
@@ -115,7 +115,7 @@ class TransactionCoordinatorTest {
 
     coordinator.handleInitProducerId("", txnTimeoutMs, enableTwoPCFlag = true,
       keepPreparedTxn = false, None, initProducerIdMockCallback)
-    assertEquals(InitProducerIdResult(-1L, -1, Errors.INVALID_REQUEST), result)
+    assertEquals(new InitProducerIdResult(-1L, -1, Errors.INVALID_REQUEST), 
result)
   }
 
   @Test
@@ -124,10 +124,10 @@ class TransactionCoordinatorTest {
 
     coordinator.handleInitProducerId(null, txnTimeoutMs, enableTwoPCFlag = 
false,
       keepPreparedTxn = false, None, initProducerIdMockCallback)
-    assertEquals(InitProducerIdResult(0L, 0, Errors.NONE), result)
+    assertEquals(new InitProducerIdResult(0L, 0, Errors.NONE), result)
     coordinator.handleInitProducerId(null, txnTimeoutMs, enableTwoPCFlag = 
false,
       keepPreparedTxn = false, None, initProducerIdMockCallback)
-    assertEquals(InitProducerIdResult(1L, 0, Errors.NONE), result)
+    assertEquals(new InitProducerIdResult(1L, 0, Errors.NONE), result)
   }
 
   @Test
@@ -159,7 +159,7 @@ class TransactionCoordinatorTest {
       None,
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(nextPid - 1, 0, Errors.NONE), result)
+    assertEquals(new InitProducerIdResult(nextPid - 1, 0, Errors.NONE), result)
   }
 
   @Test
@@ -191,7 +191,7 @@ class TransactionCoordinatorTest {
       Some(new ProducerIdAndEpoch(producerId, producerEpoch)),
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(nextPid - 1, 0, Errors.NONE), result)
+    assertEquals(new InitProducerIdResult(nextPid - 1, 0, Errors.NONE), result)
   }
 
   @Test
@@ -275,7 +275,7 @@ class TransactionCoordinatorTest {
       None,
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(-1, -1, Errors.NOT_COORDINATOR), result)
+    assertEquals(new InitProducerIdResult(-1, -1, Errors.NOT_COORDINATOR), 
result)
   }
 
   @Test
@@ -293,7 +293,7 @@ class TransactionCoordinatorTest {
       None,
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(-1, -1, 
Errors.COORDINATOR_LOAD_IN_PROGRESS), result)
+    assertEquals(new InitProducerIdResult(-1, -1, 
Errors.COORDINATOR_LOAD_IN_PROGRESS), result)
   }
 
   @Test
@@ -1063,7 +1063,7 @@ class TransactionCoordinatorTest {
       initProducerIdMockCallback
     )
 
-    assertEquals(InitProducerIdResult(-1, -1, Errors.CONCURRENT_TRANSACTIONS), 
result)
+    assertEquals(new InitProducerIdResult(-1, -1, 
Errors.CONCURRENT_TRANSACTIONS), result)
     verify(transactionManager).validateTransactionTimeoutMs(anyBoolean(), 
anyInt())
     verify(transactionManager, 
times(3)).getTransactionState(ArgumentMatchers.eq(transactionalId))
     verify(transactionManager).appendTransactionToLog(
@@ -1100,7 +1100,7 @@ class TransactionCoordinatorTest {
       initProducerIdMockCallback
     )
 
-    assertEquals(InitProducerIdResult(-1, -1, Errors.PRODUCER_FENCED), result)
+    assertEquals(new InitProducerIdResult(-1, -1, Errors.PRODUCER_FENCED), 
result)
 
     verify(transactionManager).validateTransactionTimeoutMs(anyBoolean(), 
anyInt())
     verify(transactionManager, 
times(2)).getTransactionState(ArgumentMatchers.eq(transactionalId))
@@ -1155,7 +1155,7 @@ class TransactionCoordinatorTest {
       None,
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(-1, -1, Errors.NOT_ENOUGH_REPLICAS), 
result)
+    assertEquals(new InitProducerIdResult(-1, -1, Errors.NOT_ENOUGH_REPLICAS), 
result)
 
     assertEquals((producerEpoch + 1).toShort, txnMetadata.producerEpoch)
     assertTrue(txnMetadata.hasFailedEpochFence)
@@ -1168,7 +1168,7 @@ class TransactionCoordinatorTest {
       None,
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(-1, -1, Errors.NOT_ENOUGH_REPLICAS), 
result)
+    assertEquals(new InitProducerIdResult(-1, -1, Errors.NOT_ENOUGH_REPLICAS), 
result)
 
     assertEquals((producerEpoch + 1).toShort, txnMetadata.producerEpoch)
     assertTrue(txnMetadata.hasFailedEpochFence)
@@ -1182,7 +1182,7 @@ class TransactionCoordinatorTest {
       None,
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(-1, -1, Errors.CONCURRENT_TRANSACTIONS), 
result)
+    assertEquals(new InitProducerIdResult(-1, -1, 
Errors.CONCURRENT_TRANSACTIONS), result)
 
     assertEquals((producerEpoch + 1).toShort, txnMetadata.producerEpoch)
     assertFalse(txnMetadata.hasFailedEpochFence)
@@ -1247,7 +1247,7 @@ class TransactionCoordinatorTest {
     )
     assertEquals(Short.MaxValue, txnMetadata.producerEpoch)
 
-    assertEquals(InitProducerIdResult(-1, -1, Errors.CONCURRENT_TRANSACTIONS), 
result)
+    assertEquals(new InitProducerIdResult(-1, -1, 
Errors.CONCURRENT_TRANSACTIONS), result)
     verify(transactionManager).validateTransactionTimeoutMs(anyBoolean(), 
anyInt())
     verify(transactionManager, 
times(3)).getTransactionState(ArgumentMatchers.eq(transactionalId))
     verify(transactionManager).appendTransactionToLog(
@@ -1667,7 +1667,7 @@ class TransactionCoordinatorTest {
     )
 
     // THEN1
-    assertEquals(InitProducerIdResult(-1, -1, Errors.CONCURRENT_TRANSACTIONS), 
result)
+    assertEquals(new InitProducerIdResult(-1, -1, 
Errors.CONCURRENT_TRANSACTIONS), result)
 
     val capturedTransitions = 
capturedTxnTransitMetadata.getAllValues.asScala.toList
     val firstAbortTransition = capturedTransitions.head
@@ -1873,7 +1873,7 @@ class TransactionCoordinatorTest {
     )
     
     // THEN
-    val expectedResult = InitProducerIdResult(rotatedProducerId, rotatedEpoch, 
Errors.NONE) 
+    val expectedResult = new InitProducerIdResult(rotatedProducerId, 
rotatedEpoch, Errors.NONE) 
     assertEquals(expectedResult, result)
   }
 
@@ -1898,7 +1898,7 @@ class TransactionCoordinatorTest {
       Some(new ProducerIdAndEpoch(producerId, producerEpoch)),
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(RecordBatch.NO_PRODUCER_ID, 
RecordBatch.NO_PRODUCER_EPOCH, Errors.PRODUCER_FENCED), result)
+    assertEquals(new InitProducerIdResult(RecordBatch.NO_PRODUCER_ID, 
RecordBatch.NO_PRODUCER_EPOCH, Errors.PRODUCER_FENCED), result)
   }
 
   @Test
@@ -1921,7 +1921,7 @@ class TransactionCoordinatorTest {
       Some(new ProducerIdAndEpoch(producerId, producerEpoch)),
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(RecordBatch.NO_PRODUCER_ID, 
RecordBatch.NO_PRODUCER_EPOCH, Errors.PRODUCER_FENCED), result)
+    assertEquals(new InitProducerIdResult(RecordBatch.NO_PRODUCER_ID, 
RecordBatch.NO_PRODUCER_EPOCH, Errors.PRODUCER_FENCED), result)
   }
 
   @Test
@@ -1957,7 +1957,7 @@ class TransactionCoordinatorTest {
       Some(new ProducerIdAndEpoch(producerId, 10)),
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(producerId, 11, Errors.NONE), result)
+    assertEquals(new InitProducerIdResult(producerId, 11, Errors.NONE), result)
 
     // Simulate producer retrying after successfully re-initializing but 
failing to receive the response
     coordinator.handleInitProducerId(
@@ -1968,7 +1968,7 @@ class TransactionCoordinatorTest {
       Some(new ProducerIdAndEpoch(producerId, 10)),
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(producerId, 11, Errors.NONE), result)
+    assertEquals(new InitProducerIdResult(producerId, 11, Errors.NONE), result)
   }
 
   @Test
@@ -2007,7 +2007,7 @@ class TransactionCoordinatorTest {
       None,
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(producerId, 11, Errors.NONE), result)
+    assertEquals(new InitProducerIdResult(producerId, 11, Errors.NONE), result)
 
     // Simulate old producer trying to continue from epoch 10
     coordinator.handleInitProducerId(
@@ -2018,7 +2018,7 @@ class TransactionCoordinatorTest {
       Some(new ProducerIdAndEpoch(producerId, 10)),
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(RecordBatch.NO_PRODUCER_ID, 
RecordBatch.NO_PRODUCER_EPOCH, Errors.PRODUCER_FENCED), result)
+    assertEquals(new InitProducerIdResult(RecordBatch.NO_PRODUCER_ID, 
RecordBatch.NO_PRODUCER_EPOCH, Errors.PRODUCER_FENCED), result)
   }
 
   @Test
@@ -2060,7 +2060,7 @@ class TransactionCoordinatorTest {
       Some(new ProducerIdAndEpoch(producerId, (Short.MaxValue - 1).toShort)),
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(producerId + 1, 0, Errors.NONE), result)
+    assertEquals(new InitProducerIdResult(producerId + 1, 0, Errors.NONE), 
result)
 
     // Simulate producer retrying old request after producer bump
     coordinator.handleInitProducerId(
@@ -2071,7 +2071,7 @@ class TransactionCoordinatorTest {
       Some(new ProducerIdAndEpoch(producerId, (Short.MaxValue - 1).toShort)),
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(producerId + 1, 0, Errors.NONE), result)
+    assertEquals(new InitProducerIdResult(producerId + 1, 0, Errors.NONE), 
result)
   }
 
   @Test
@@ -2113,7 +2113,7 @@ class TransactionCoordinatorTest {
       Some(new ProducerIdAndEpoch(producerId, (Short.MaxValue - 1).toShort)),
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(producerId + 1, 0, Errors.NONE), result)
+    assertEquals(new InitProducerIdResult(producerId + 1, 0, Errors.NONE), 
result)
 
     // Validate that producer with old producer ID and stale epoch is fenced
     coordinator.handleInitProducerId(
@@ -2124,7 +2124,7 @@ class TransactionCoordinatorTest {
       Some(new ProducerIdAndEpoch(producerId, (Short.MaxValue - 2).toShort)),
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(RecordBatch.NO_PRODUCER_ID, 
RecordBatch.NO_PRODUCER_EPOCH, Errors.PRODUCER_FENCED), result)
+    assertEquals(new InitProducerIdResult(RecordBatch.NO_PRODUCER_ID, 
RecordBatch.NO_PRODUCER_EPOCH, Errors.PRODUCER_FENCED), result)
   }
 
   @Test
@@ -2286,7 +2286,7 @@ class TransactionCoordinatorTest {
       Some(new ProducerIdAndEpoch(producerId, 10)),
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(RecordBatch.NO_PRODUCER_ID, 
RecordBatch.NO_PRODUCER_EPOCH, Errors.CONCURRENT_TRANSACTIONS), result)
+    assertEquals(new InitProducerIdResult(RecordBatch.NO_PRODUCER_ID, 
RecordBatch.NO_PRODUCER_EPOCH, Errors.CONCURRENT_TRANSACTIONS), result)
 
     verify(transactionManager).validateTransactionTimeoutMs(anyBoolean(), 
anyInt())
     
verify(transactionManager).getTransactionState(ArgumentMatchers.eq(transactionalId))
@@ -2368,7 +2368,7 @@ class TransactionCoordinatorTest {
     coordinator.handleInitProducerId(transactionalId, 10, enableTwoPCFlag = 
false,
       keepPreparedTxn = false, None, initProducerIdMockCallback)
 
-    assertEquals(InitProducerIdResult(-1, -1, Errors.CONCURRENT_TRANSACTIONS), 
result)
+    assertEquals(new InitProducerIdResult(-1, -1, 
Errors.CONCURRENT_TRANSACTIONS), result)
   }
 
   private def validateIncrementEpochAndUpdateMetadata(state: TransactionState, 
transactionVersion: Short): Unit = {
@@ -2401,7 +2401,7 @@ class TransactionCoordinatorTest {
     coordinator.handleInitProducerId(transactionalId, newTxnTimeoutMs, 
enableTwoPCFlag = false,
       keepPreparedTxn = false, None, initProducerIdMockCallback)
 
-    assertEquals(InitProducerIdResult(producerId, (producerEpoch + 1).toShort, 
Errors.NONE), result)
+    assertEquals(new InitProducerIdResult(producerId, (producerEpoch + 
1).toShort, Errors.NONE), result)
     assertEquals(newTxnTimeoutMs, metadata.txnTimeoutMs)
     assertEquals(time.milliseconds(), metadata.txnLastUpdateTimestamp)
     assertEquals((producerEpoch + 1).toShort, metadata.producerEpoch)
@@ -2498,7 +2498,7 @@ class TransactionCoordinatorTest {
       None,
       initProducerIdMockCallback
     )
-    assertEquals(InitProducerIdResult(-1, -1, 
Errors.COORDINATOR_NOT_AVAILABLE), result)
+    assertEquals(new InitProducerIdResult(-1, -1, 
Errors.COORDINATOR_NOT_AVAILABLE), result)
 
     // After the first failed attempt, the state should be:
     // - hasFailedEpochFence = false (NOT set for TV2)
@@ -2557,7 +2557,7 @@ class TransactionCoordinatorTest {
     )
 
     // The second attempt should return CONCURRENT_TRANSACTIONS (this is 
intentional)
-    assertEquals(InitProducerIdResult(-1, -1, Errors.CONCURRENT_TRANSACTIONS), 
result)
+    assertEquals(new InitProducerIdResult(-1, -1, 
Errors.CONCURRENT_TRANSACTIONS), result)
 
     // The transactionMarkerChannelManager mock should have completed the 
transition to COMPLETE_ABORT
     // Verify that hasFailedEpochFence was never set to true for TV2, allowing 
future epoch bumps
@@ -2598,7 +2598,7 @@ class TransactionCoordinatorTest {
 
     // The third attempt should succeed with epoch 3 (2 + 1)
     // This demonstrates that TV2 allows epoch re-bumping after failed writes
-    assertEquals(InitProducerIdResult(producerId, 3.toShort, Errors.NONE), 
result)
+    assertEquals(new InitProducerIdResult(producerId, 3.toShort, Errors.NONE), 
result)
     
     // Final verification that hasFailedEpochFence was never set to true for 
TV2
     assertFalse(txnMetadata.hasFailedEpochFence)
diff --git 
a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionMarkerRequestCompletionHandlerTest.scala
 
b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionMarkerRequestCompletionHandlerTest.scala
index 1678bc830ea..9869fcbc1ca 100644
--- 
a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionMarkerRequestCompletionHandlerTest.scala
+++ 
b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionMarkerRequestCompletionHandlerTest.scala
@@ -22,7 +22,7 @@ import org.apache.kafka.common.TopicPartition
 import org.apache.kafka.common.protocol.{ApiKeys, Errors}
 import org.apache.kafka.common.record.internal.RecordBatch
 import org.apache.kafka.common.requests.{RequestHeader, TransactionResult, 
WriteTxnMarkersRequest, WriteTxnMarkersResponse}
-import 
org.apache.kafka.coordinator.transaction.{CoordinatorEpochAndTxnMetadata, 
TransactionMetadata, TransactionState}
+import 
org.apache.kafka.coordinator.transaction.{CoordinatorEpochAndTxnMetadata, 
PendingCompleteTxn, PendingCompleteTxnAndMarkerEntry, TransactionMetadata, 
TransactionState}
 import org.apache.kafka.server.common.TransactionVersion
 import org.junit.jupiter.api.Assertions._
 import org.junit.jupiter.api.Test
@@ -44,8 +44,8 @@ class TransactionMarkerRequestCompletionHandlerTest {
   private val txnMetadata = new TransactionMetadata(transactionalId, 
producerId, producerId, RecordBatch.NO_PRODUCER_ID,
     producerEpoch, lastProducerEpoch, txnTimeoutMs, 
TransactionState.PREPARE_COMMIT, util.Set.of(topicPartition), 0L, 0L, 
TransactionVersion.TV_2)
   private val pendingCompleteTxnAndMarkers = util.List.of(
-    PendingCompleteTxnAndMarkerEntry(
-      PendingCompleteTxn(transactionalId, coordinatorEpoch, txnMetadata, 
txnMetadata.prepareComplete(42)),
+    new PendingCompleteTxnAndMarkerEntry(
+      new PendingCompleteTxn(transactionalId, coordinatorEpoch, txnMetadata, 
txnMetadata.prepareComplete(42)),
       new WriteTxnMarkersRequest.TxnMarkerEntry(producerId, producerEpoch, 
coordinatorEpoch, txnResult, util.List.of(topicPartition), 0)))
 
   private val markerChannelManager: TransactionMarkerChannelManager =
diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala 
b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
index de7ed5c902d..6310d0a2f6c 100644
--- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
+++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
@@ -19,7 +19,7 @@ package kafka.server
 
 import com.yammer.metrics.core.{Histogram, Meter}
 import kafka.cluster.Partition
-import kafka.coordinator.transaction.{InitProducerIdResult, 
TransactionCoordinator}
+import kafka.coordinator.transaction.TransactionCoordinator
 import kafka.network.RequestChannel
 import kafka.server.QuotaFactory.QuotaManagers
 import kafka.server.share.SharePartitionManager
@@ -84,7 +84,7 @@ import 
org.apache.kafka.coordinator.group.modern.share.ShareGroupConfig
 import org.apache.kafka.coordinator.group.{GroupConfig, GroupConfigManager, 
GroupCoordinator, GroupCoordinatorConfig}
 import org.apache.kafka.coordinator.group.streams.StreamsGroupHeartbeatResult
 import org.apache.kafka.coordinator.share.{ShareCoordinator, 
ShareCoordinatorTestConfig}
-import org.apache.kafka.coordinator.transaction.TransactionLogConfig
+import org.apache.kafka.coordinator.transaction.{InitProducerIdResult, 
TransactionLogConfig}
 import org.apache.kafka.image.{MetadataDelta, MetadataImage, 
MetadataProvenance}
 import org.apache.kafka.metadata.{ConfigRepository, KRaftMetadataCache, 
MetadataCache, MetadataCacheFixtures, MockConfigRepository}
 import org.apache.kafka.network.{Request, Session}
@@ -1811,7 +1811,7 @@ class KafkaApisTest extends Logging {
         ArgumentMatchers.eq(expectedProducerIdAndEpoch),
         responseCallback.capture(),
         ArgumentMatchers.eq(requestLocal)
-      )).thenAnswer(_ => 
responseCallback.getValue.apply(InitProducerIdResult(producerId, epoch, 
Errors.PRODUCER_FENCED)))
+      )).thenAnswer(_ => responseCallback.getValue.apply(new 
InitProducerIdResult(producerId, epoch, Errors.PRODUCER_FENCED)))
       val kafkaApis = createKafkaApis()
       try {
         kafkaApis.handleInitProducerIdRequest(request, requestLocal)
@@ -2055,7 +2055,7 @@ class KafkaApisTest extends Logging {
       any(),
       responseCallback.capture(),
       ArgumentMatchers.eq(requestLocal)
-    )).thenAnswer(_ => 
responseCallback.getValue.apply(InitProducerIdResult(15L, 0.toShort, 
Errors.NONE)))
+    )).thenAnswer(_ => responseCallback.getValue.apply(new 
InitProducerIdResult(15L, 0.toShort, Errors.NONE)))
 
     kafkaApis.handleInitProducerIdRequest(request, requestLocal)
 
diff --git 
a/transaction-coordinator/src/main/java/org/apache/kafka/coordinator/transaction/InitProducerIdResult.java
 
b/transaction-coordinator/src/main/java/org/apache/kafka/coordinator/transaction/InitProducerIdResult.java
new file mode 100644
index 00000000000..2d31c000d0d
--- /dev/null
+++ 
b/transaction-coordinator/src/main/java/org/apache/kafka/coordinator/transaction/InitProducerIdResult.java
@@ -0,0 +1,22 @@
+/*
+ * 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.kafka.coordinator.transaction;
+
+import org.apache.kafka.common.protocol.Errors;
+
+public record InitProducerIdResult(long producerId, short producerEpoch, 
Errors error) {
+}
diff --git 
a/transaction-coordinator/src/main/java/org/apache/kafka/coordinator/transaction/PendingCompleteTxn.java
 
b/transaction-coordinator/src/main/java/org/apache/kafka/coordinator/transaction/PendingCompleteTxn.java
new file mode 100644
index 00000000000..81bfde78a1f
--- /dev/null
+++ 
b/transaction-coordinator/src/main/java/org/apache/kafka/coordinator/transaction/PendingCompleteTxn.java
@@ -0,0 +1,25 @@
+/*
+ * 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.kafka.coordinator.transaction;
+
+public record PendingCompleteTxn(
+    String transactionalId,
+    int coordinatorEpoch,
+    TransactionMetadata txnMetadata,
+    TxnTransitMetadata newMetadata
+) {
+}
diff --git 
a/transaction-coordinator/src/main/java/org/apache/kafka/coordinator/transaction/PendingCompleteTxnAndMarkerEntry.java
 
b/transaction-coordinator/src/main/java/org/apache/kafka/coordinator/transaction/PendingCompleteTxnAndMarkerEntry.java
new file mode 100644
index 00000000000..60d1d17f232
--- /dev/null
+++ 
b/transaction-coordinator/src/main/java/org/apache/kafka/coordinator/transaction/PendingCompleteTxnAndMarkerEntry.java
@@ -0,0 +1,25 @@
+/*
+ * 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.kafka.coordinator.transaction;
+
+import org.apache.kafka.common.requests.WriteTxnMarkersRequest.TxnMarkerEntry;
+
+public record PendingCompleteTxnAndMarkerEntry(
+    PendingCompleteTxn pendingCompleteTxn,
+    TxnMarkerEntry txnMarkerEntry
+) {
+}
diff --git 
a/transaction-coordinator/src/main/java/org/apache/kafka/coordinator/transaction/TxnMarkerQueue.java
 
b/transaction-coordinator/src/main/java/org/apache/kafka/coordinator/transaction/TxnMarkerQueue.java
new file mode 100644
index 00000000000..f82091ed439
--- /dev/null
+++ 
b/transaction-coordinator/src/main/java/org/apache/kafka/coordinator/transaction/TxnMarkerQueue.java
@@ -0,0 +1,100 @@
+/*
+ * 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.kafka.coordinator.transaction;
+
+import org.apache.kafka.common.Node;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.Optional;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.function.BiConsumer;
+
+public class TxnMarkerQueue {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(TxnMarkerQueue.class);
+
+    // keep track of the requests per txn topic partition so we can easily 
clear the queue
+    // during partition emigration
+    private final ConcurrentHashMap<Integer, 
BlockingQueue<PendingCompleteTxnAndMarkerEntry>> markersPerTxnTopicPartition =
+        new ConcurrentHashMap<>();
+
+    private volatile Node destination;
+
+    public TxnMarkerQueue(Node destination) {
+        this.destination = destination;
+    }
+
+    public Node destination() {
+        return destination;
+    }
+
+    public void setDestination(Node destination) {
+        this.destination = destination;
+    }
+
+    public Optional<BlockingQueue<PendingCompleteTxnAndMarkerEntry>> 
removeMarkersForTxnTopicPartition(int partition) {
+        return 
Optional.ofNullable(markersPerTxnTopicPartition.remove(partition));
+    }
+
+    public void addMarkers(int txnTopicPartition, 
PendingCompleteTxnAndMarkerEntry pendingCompleteTxnAndMarker) {
+        BlockingQueue<PendingCompleteTxnAndMarkerEntry> queue = 
markersPerTxnTopicPartition.computeIfAbsent(txnTopicPartition, partition -> {
+            LOG.info("Creating new marker queue for txn partition {} to 
destination broker {}", txnTopicPartition, destination.id());
+            return new LinkedBlockingQueue<>();
+        });
+        queue.add(pendingCompleteTxnAndMarker);
+
+        if (markersPerTxnTopicPartition.get(txnTopicPartition) != queue) {
+            // This could happen if the queue got removed concurrently.
+            // Note that it could create an unexpected state when the queue is 
removed from
+            // removeMarkersForTxnTopicPartition, we could have:
+            //
+            // 1. [addMarkers] Retrieve queue.
+            // 2. [removeMarkersForTxnTopicPartition] Remove queue.
+            // 3. [removeMarkersForTxnTopicPartition] Iterate over queue, but 
not removeMarkersForTxn because queue is empty.
+            // 4. [addMarkers] Add markers to the queue.
+            //
+            // Now we've effectively removed the markers while 
transactionsWithPendingMarkers has an entry.
+            //
+            // While this could lead to an orphan entry in 
transactionsWithPendingMarkers, sending new markers
+            // will fix the state, so it shouldn't impact the state machine 
operation.
+            LOG.warn("Added {} to dead queue for txn partition {} to 
destination broker {}",
+                pendingCompleteTxnAndMarker, txnTopicPartition, 
destination.id());
+        }
+    }
+
+    public void forEachTxnTopicPartition(BiConsumer<Integer, 
BlockingQueue<PendingCompleteTxnAndMarkerEntry>> f) {
+        markersPerTxnTopicPartition.forEach(f);
+    }
+
+    public int totalNumMarkers() {
+        int total = 0;
+        for (BlockingQueue<PendingCompleteTxnAndMarkerEntry> queue : 
markersPerTxnTopicPartition.values()) {
+            total += queue.size();
+        }
+        return total;
+    }
+
+    // visible for testing
+    public int totalNumMarkers(int txnTopicPartition) {
+        BlockingQueue<PendingCompleteTxnAndMarkerEntry> queue = 
markersPerTxnTopicPartition.get(txnTopicPartition);
+        return queue == null ? 0 : queue.size();
+    }
+}


Reply via email to