Skip to content
Merged
Show file tree
Hide file tree
Changes from 34 commits
Commits
Show all changes
36 commits
Select commit Hold shift + click to select a range
26e1a49
KAFKA-18884: Move TransactionMetadata to transaction-coordinator module
FrankYang0529 May 13, 2025
ece6ef9
KAFAK-18884: remove some java / scala transformation
FrankYang0529 May 17, 2025
8e9afc2
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 May 17, 2025
0bee440
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 May 20, 2025
c1f2485
address comment
FrankYang0529 May 20, 2025
180881c
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 May 23, 2025
703ea8a
address comment
FrankYang0529 May 23, 2025
53c61b7
remove unused function
FrankYang0529 May 24, 2025
ff335a1
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 May 24, 2025
abaa206
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 May 25, 2025
c39b9b5
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 May 26, 2025
97bc2a5
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 May 28, 2025
3edd06e
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 May 30, 2025
1d12df5
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 May 30, 2025
e65d6e3
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Jun 2, 2025
962934c
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Jun 6, 2025
119d7d0
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Jun 6, 2025
c7f95e6
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Jun 10, 2025
3e70223
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Jun 13, 2025
7dc62a4
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Jun 15, 2025
94b3703
add TransitionData
FrankYang0529 Jun 15, 2025
663ddca
chnage TransactionData to TransitionData
FrankYang0529 Jun 15, 2025
53b87f7
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Jun 20, 2025
d34d2bf
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Jul 2, 2025
11b175f
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Jul 7, 2025
e4dd2c2
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Jul 8, 2025
66b58c5
address comment
FrankYang0529 Jul 8, 2025
7eff967
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Jul 17, 2025
c5a3791
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Jul 31, 2025
f80e9ab
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Aug 7, 2025
a0f379c
fix: change TransactionMetadata class path in spotbugs-exclude.xml
FrankYang0529 Aug 7, 2025
f7f0687
feat: change topicPartitions from Set to HashSet
FrankYang0529 Aug 8, 2025
356b47e
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Aug 9, 2025
2ade564
feat: change topicPartitions from Set to HashSet
FrankYang0529 Aug 9, 2025
eacf843
Merge branch 'trunk' into KAFKA-18884
FrankYang0529 Aug 11, 2025
fd10a55
address comment
FrankYang0529 Aug 11, 2025
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -21,11 +21,12 @@ import org.apache.kafka.common.compress.Compression
import org.apache.kafka.common.protocol.{ByteBufferAccessor, MessageUtil}
import org.apache.kafka.common.record.RecordBatch
import org.apache.kafka.common.TopicPartition
import org.apache.kafka.coordinator.transaction.{TransactionState, TxnTransitMetadata}
import org.apache.kafka.coordinator.transaction.{TransactionMetadata, TransactionState, TxnTransitMetadata}
import org.apache.kafka.coordinator.transaction.generated.{CoordinatorRecordType, TransactionLogKey, TransactionLogValue}
import org.apache.kafka.server.common.TransactionVersion

import scala.collection.mutable
import java.util

import scala.jdk.CollectionConverters._

/**
Expand Down Expand Up @@ -115,26 +116,26 @@ object TransactionLog {
if (version >= TransactionLogValue.LOWEST_SUPPORTED_VERSION && version <= TransactionLogValue.HIGHEST_SUPPORTED_VERSION) {
val value = new TransactionLogValue(new ByteBufferAccessor(buffer), version)
val transactionMetadata = new TransactionMetadata(
transactionalId = transactionalId,
producerId = value.producerId,
prevProducerId = value.previousProducerId,
nextProducerId = value.nextProducerId,
producerEpoch = value.producerEpoch,
lastProducerEpoch = RecordBatch.NO_PRODUCER_EPOCH,
txnTimeoutMs = value.transactionTimeoutMs,
state = TransactionState.fromId(value.transactionStatus),
topicPartitions = mutable.Set.empty[TopicPartition],
txnStartTimestamp = value.transactionStartTimestampMs,
txnLastUpdateTimestamp = value.transactionLastUpdateTimestampMs,
clientTransactionVersion = TransactionVersion.fromFeatureLevel(value.clientTransactionVersion))
transactionalId,
value.producerId,
value.previousProducerId,
value.nextProducerId,
value.producerEpoch,
RecordBatch.NO_PRODUCER_EPOCH,
value.transactionTimeoutMs,
TransactionState.fromId(value.transactionStatus),
util.Set.of(),
value.transactionStartTimestampMs,
value.transactionLastUpdateTimestampMs,
TransactionVersion.fromFeatureLevel(value.clientTransactionVersion))

if (!transactionMetadata.state.equals(TransactionState.EMPTY))
value.transactionPartitions.forEach(partitionsSchema =>
value.transactionPartitions.forEach(partitionsSchema => {
transactionMetadata.addPartitions(partitionsSchema.partitionIds
.asScala
.map(partitionId => new TopicPartition(partitionsSchema.topic, partitionId))
.toSet)
)
.stream
.map(partitionId => new TopicPartition(partitionsSchema.topic, partitionId.intValue()))
.toList)
})
Some(transactionMetadata)
} else throw new IllegalStateException(s"Unknown version $version from the transaction log message value")
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ import org.apache.kafka.common.requests.{TransactionResult, WriteTxnMarkersReque
import org.apache.kafka.common.security.JaasContext
import org.apache.kafka.common.utils.{LogContext, Time}
import org.apache.kafka.common.{Node, Reconfigurable, TopicPartition}
import org.apache.kafka.coordinator.transaction.TxnTransitMetadata
import org.apache.kafka.coordinator.transaction.{TransactionMetadata, TxnTransitMetadata}
import org.apache.kafka.metadata.MetadataCache
import org.apache.kafka.server.common.RequestLocal
import org.apache.kafka.server.metrics.KafkaMetricsGroup
Expand Down Expand Up @@ -326,16 +326,16 @@ class TransactionMarkerChannelManager(
info(s"Replaced an existing pending complete txn $prev with $pendingCompleteTxn while adding markers to send.")
}
addTxnMarkersToBrokerQueue(txnMetadata.producerId,
txnMetadata.producerEpoch, txnResult, pendingCompleteTxn, txnMetadata.topicPartitions.toSet)
txnMetadata.producerEpoch, txnResult, pendingCompleteTxn, txnMetadata.topicPartitions.asScala.toSet)
maybeWriteTxnCompletion(transactionalId)
}

def numTxnsWithPendingMarkers: Int = transactionsWithPendingMarkers.size

private def hasPendingMarkersToWrite(txnMetadata: TransactionMetadata): Boolean = {
txnMetadata.inLock {
txnMetadata.topicPartitions.nonEmpty
}
txnMetadata.inLock(() =>
!txnMetadata.topicPartitions.isEmpty
)
}

def maybeWriteTxnCompletion(transactionalId: String): Unit = {
Expand Down Expand Up @@ -422,9 +422,9 @@ class TransactionMarkerChannelManager(

val txnMetadata = epochAndMetadata.transactionMetadata

txnMetadata.inLock {
txnMetadata.inLock(() =>
topicPartitions.foreach(txnMetadata.removePartition)
}
)

maybeWriteTxnCompletion(transactionalId)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -131,7 +131,7 @@ class TransactionMarkerRequestCompletionHandler(brokerId: Int,
txnMarkerChannelManager.removeMarkersForTxn(pendingCompleteTxn)
abortSending = true
} else {
txnMetadata.inLock {
txnMetadata.inLock(() => {
for ((topicPartition, error) <- errors.asScala) {
error match {
case Errors.NONE =>
Expand Down Expand Up @@ -178,7 +178,7 @@ class TransactionMarkerRequestCompletionHandler(brokerId: Int,
throw new IllegalStateException(s"Unexpected error ${other.exceptionName} while sending txn marker for $transactionalId")
}
}
}
})
}

if (!abortSending) {
Expand Down
Loading