Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 commits
Commits
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
Original file line number Diff line number Diff line change
Expand Up @@ -156,7 +156,6 @@ class TransactionCoordinator(txnConfig: TransactionConfig,
RecordBatch.NO_PRODUCER_EPOCH,
resolvedTxnTimeoutMs,
TransactionState.EMPTY,
util.Set.of(),
-1,
time.milliseconds(),
TransactionVersion.TV_0)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,6 @@ import org.apache.kafka.coordinator.transaction.{TransactionMetadata, Transactio
import org.apache.kafka.coordinator.transaction.generated.{CoordinatorRecordType, TransactionLogKey, TransactionLogValue}
import org.apache.kafka.server.common.TransactionVersion

import java.util

import scala.jdk.CollectionConverters._

/**
Expand Down Expand Up @@ -124,7 +122,6 @@ object TransactionLog {
RecordBatch.NO_PRODUCER_EPOCH,
value.transactionTimeoutMs,
TransactionState.fromId(value.transactionStatus),
util.Set.of(),
value.transactionStartTimestampMs,
value.transactionLastUpdateTimestampMs,
TransactionVersion.fromFeatureLevel(value.clientTransactionVersion))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -514,7 +514,6 @@ class TransactionCoordinatorConcurrencyTest extends AbstractCoordinatorConcurren
RecordBatch.NO_PRODUCER_EPOCH,
60000,
TransactionState.EMPTY,
new util.HashSet[TopicPartition](),
-1,
time.milliseconds(),
TransactionVersion.TV_0)
Expand Down

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -49,8 +49,7 @@ class TransactionLogTest {
val transactionalId = "transactionalId"
val producerId = 23423L

val txnMetadata = new TransactionMetadata(transactionalId, producerId, RecordBatch.NO_PRODUCER_ID, RecordBatch.NO_PRODUCER_ID, producerEpoch,
RecordBatch.NO_PRODUCER_EPOCH, transactionTimeoutMs, TransactionState.EMPTY, util.Set.of, 0, 0, TV_0)
val txnMetadata = new TransactionMetadata(transactionalId, producerId, RecordBatch.NO_PRODUCER_ID, RecordBatch.NO_PRODUCER_ID, producerEpoch, RecordBatch.NO_PRODUCER_EPOCH, transactionTimeoutMs, TransactionState.EMPTY, 0, 0, TV_0)
txnMetadata.addPartitions(topicPartitions)

assertThrows(classOf[IllegalStateException], () => TransactionLog.valueToBytes(txnMetadata.prepareNoTransit(), TV_2))
Expand All @@ -74,8 +73,7 @@ class TransactionLogTest {

// generate transaction log messages
val txnRecords = pidMappings.map { case (transactionalId, producerId) =>
val txnMetadata = new TransactionMetadata(transactionalId, producerId, RecordBatch.NO_PRODUCER_ID, RecordBatch.NO_PRODUCER_ID, producerEpoch,
RecordBatch.NO_PRODUCER_EPOCH, transactionTimeoutMs, transactionStates(producerId), util.Set.of, 0, 0, TV_0)
val txnMetadata = new TransactionMetadata(transactionalId, producerId, RecordBatch.NO_PRODUCER_ID, RecordBatch.NO_PRODUCER_ID, producerEpoch, RecordBatch.NO_PRODUCER_EPOCH, transactionTimeoutMs, transactionStates(producerId), 0, 0, TV_0)

if (!txnMetadata.state.equals(TransactionState.EMPTY))
txnMetadata.addPartitions(topicPartitions)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -65,10 +65,10 @@ class TransactionMarkerChannelManagerTest {
private val coordinatorEpoch2 = 1
private val txnTimeoutMs = 0
private val txnResult = TransactionResult.COMMIT
private val txnMetadata1 = new TransactionMetadata(transactionalId1, producerId1, producerId1, RecordBatch.NO_PRODUCER_ID,
producerEpoch, lastProducerEpoch, txnTimeoutMs, TransactionState.PREPARE_COMMIT, util.Set.of(partition1, partition2), 0L, 0L, TransactionVersion.TV_2)
private val txnMetadata2 = new TransactionMetadata(transactionalId2, producerId2, producerId2, RecordBatch.NO_PRODUCER_ID,
producerEpoch, lastProducerEpoch, txnTimeoutMs, TransactionState.PREPARE_COMMIT, util.Set.of(partition1), 0L, 0L, TransactionVersion.TV_2)
private val txnMetadata1 = new TransactionMetadata(transactionalId1, producerId1, producerId1, RecordBatch.NO_PRODUCER_ID, producerEpoch, lastProducerEpoch, txnTimeoutMs, TransactionState.PREPARE_COMMIT, 0L, 0L, TransactionVersion.TV_2)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

it seems those tests assume txnMetadata1 have some topicPartitions initially.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Ah good catch -- are we testing the loading path anywhere else? I think that might be a scenario where we start with partitions.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sorry, I missed this one. Fixed it.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sorry so for when we load the partition in the non test code, we always call this add partitions method?

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

we always call this add partitions method?

yes

transactionMetadata.addPartitions(partitionsSchema.partitionIds

we could remove the usage of addPartitions from production code. for example:

        val state = TransactionState.fromId(value.transactionStatus)
        val tps: util.Set[TopicPartition] = if (!state.equals(TransactionState.EMPTY)) value.transactionPartitions
          .stream().flatMap(partitionsSchema => partitionsSchema.partitionIds().stream().map(id => new TopicPartition(partitionsSchema.topic(), id.intValue())))
          .collect(Collectors.toSet)
        else util.Set.of()
        Some(new TransactionMetadata(
          transactionalId,
          value.producerId,
          value.previousProducerId,
          value.nextProducerId,
          value.producerEpoch,
          RecordBatch.NO_PRODUCER_EPOCH,
          value.transactionTimeoutMs,
          state,
          tps,
          value.transactionStartTimestampMs,
          value.transactionLastUpdateTimestampMs,
          TransactionVersion.fromFeatureLevel(value.clientTransactionVersion)))

If the method removePartition could adopt copy-on-write policy, the inner member topicPartitions could be an immutable collection.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I see. I don't know if this is particularly better or worse -- I don't have a strong opinion either way. Just wanted to confirm that we have a way we add partitions and do it the same way each time. 👍

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It should be acceptable to remove addPartitions in this PR; we can address removePartition in a subsequent change if it does not cause performance regression.

txnMetadata1.addPartitions(util.Set.of(partition1, partition2))
private val txnMetadata2 = new TransactionMetadata(transactionalId2, producerId2, producerId2, RecordBatch.NO_PRODUCER_ID, producerEpoch, lastProducerEpoch, txnTimeoutMs, TransactionState.PREPARE_COMMIT, 0L, 0L, TransactionVersion.TV_2)
txnMetadata2.addPartitions(util.Set.of(partition1))

private val capturedErrorsCallback: ArgumentCaptor[Errors => Unit] = ArgumentCaptor.forClass(classOf[Errors => Unit])
private val time = new MockTime
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,8 +41,8 @@ class TransactionMarkerRequestCompletionHandlerTest {
private val coordinatorEpoch = 0
private val txnResult = TransactionResult.COMMIT
private val topicPartition = new TopicPartition("topic1", 0)
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 txnMetadata = new TransactionMetadata(transactionalId, producerId, producerId, RecordBatch.NO_PRODUCER_ID, producerEpoch, lastProducerEpoch, txnTimeoutMs, TransactionState.PREPARE_COMMIT, 0L, 0L, TransactionVersion.TV_2)
txnMetadata.addPartitions(util.Set.of(topicPartition))
private val pendingCompleteTxnAndMarkers = util.List.of(
PendingCompleteTxnAndMarkerEntry(
PendingCompleteTxn(transactionalId, coordinatorEpoch, txnMetadata, txnMetadata.prepareComplete(42)),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.EMPTY,
util.Set.of,
-1,
time.milliseconds(),
TV_0)
Expand All @@ -78,7 +77,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.EMPTY,
util.Set.of,
-1,
time.milliseconds(),
TV_0)
Expand All @@ -103,7 +101,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.EMPTY,
util.Set.of,
-1,
time.milliseconds(),
TV_0)
Expand All @@ -126,7 +123,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.EMPTY,
util.Set.of,
-1,
time.milliseconds(),
TV_2)
Expand All @@ -151,7 +147,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.COMPLETE_ABORT,
util.Set.of,
time.milliseconds() - 1,
time.milliseconds(),
TV_2)
Expand All @@ -176,7 +171,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.COMPLETE_COMMIT,
util.Set.of,
time.milliseconds() - 1,
time.milliseconds(),
TV_2)
Expand All @@ -200,7 +194,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.EMPTY,
util.Set.of,
1L,
time.milliseconds(),
TV_0)
Expand Down Expand Up @@ -228,7 +221,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.EMPTY,
util.Set.of,
1L,
time.milliseconds(),
TV_0)
Expand All @@ -255,7 +247,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.EMPTY,
util.Set.of,
time.milliseconds(),
time.milliseconds(),
TV_0)
Expand Down Expand Up @@ -293,7 +284,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.ONGOING,
util.Set.of,
1L,
time.milliseconds(),
TV_0)
Expand Down Expand Up @@ -321,7 +311,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.ONGOING,
util.Set.of,
1L,
time.milliseconds(),
TV_0)
Expand Down Expand Up @@ -352,7 +341,6 @@ class TransactionMetadataTest {
lastProducerEpoch,
30000,
TransactionState.PREPARE_COMMIT,
util.Set.of(),
1L,
time.milliseconds(),
clientTransactionVersion
Expand Down Expand Up @@ -385,7 +373,6 @@ class TransactionMetadataTest {
lastProducerEpoch,
30000,
TransactionState.PREPARE_ABORT,
util.Set.of,
1L,
time.milliseconds(),
clientTransactionVersion
Expand Down Expand Up @@ -416,7 +403,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.ONGOING,
util.Set.of,
-1,
time.milliseconds(),
TV_0)
Expand Down Expand Up @@ -448,7 +434,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.COMPLETE_COMMIT,
util.Set.of,
-1,
time.milliseconds(),
TV_0)
Expand All @@ -470,7 +455,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.COMPLETE_ABORT,
util.Set.of,
-1,
time.milliseconds(),
TV_0)
Expand All @@ -492,7 +476,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.ONGOING,
util.Set.of,
-1,
time.milliseconds(),
TV_0)
Expand All @@ -513,7 +496,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.EMPTY,
util.Set.of,
-1,
time.milliseconds(),
TV_0)
Expand Down Expand Up @@ -541,7 +523,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.ONGOING,
util.Set.of,
time.milliseconds(),
time.milliseconds(),
TV_2)
Expand Down Expand Up @@ -573,7 +554,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.ONGOING,
util.Set.of,
time.milliseconds(),
time.milliseconds(),
TV_2)
Expand Down Expand Up @@ -627,7 +607,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.EMPTY,
util.Set.of,
-1,
time.milliseconds(),
TV_0)
Expand All @@ -652,7 +631,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
TransactionState.EMPTY,
util.Set.of,
-1,
time.milliseconds(),
TV_0)
Expand All @@ -678,7 +656,6 @@ class TransactionMetadataTest {
lastProducerEpoch,
30000,
TransactionState.EMPTY,
util.Set.of,
-1,
time.milliseconds(),
TV_0)
Expand All @@ -704,7 +681,6 @@ class TransactionMetadataTest {
lastProducerEpoch,
30000,
TransactionState.EMPTY,
util.Set.of,
-1,
time.milliseconds(),
TV_0)
Expand Down Expand Up @@ -770,7 +746,6 @@ class TransactionMetadataTest {
RecordBatch.NO_PRODUCER_EPOCH,
30000,
state,
util.Set.of,
-1,
time.milliseconds(),
clientTransactionVersion)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -875,8 +875,7 @@ class TransactionStateManagerTest {
// is left at it is. If the transactional id is never reused, the TransactionMetadata
// will be expired and it should succeed.
val timestamp = time.milliseconds()
val txnMetadata = new TransactionMetadata(transactionalId, 1, RecordBatch.NO_PRODUCER_ID, RecordBatch.NO_PRODUCER_ID, RecordBatch.NO_PRODUCER_EPOCH,
RecordBatch.NO_PRODUCER_EPOCH, transactionTimeoutMs, TransactionState.EMPTY, util.Set.of, timestamp, timestamp, TV_0)
val txnMetadata = new TransactionMetadata(transactionalId, 1, RecordBatch.NO_PRODUCER_ID, RecordBatch.NO_PRODUCER_ID, RecordBatch.NO_PRODUCER_EPOCH, RecordBatch.NO_PRODUCER_EPOCH, transactionTimeoutMs, TransactionState.EMPTY, timestamp, timestamp, TV_0)
transactionManager.putTransactionStateIfNotExists(txnMetadata)

time.sleep(txnConfig.transactionalIdExpirationMs + 1)
Expand Down Expand Up @@ -1221,8 +1220,7 @@ class TransactionStateManagerTest {
state: TransactionState = TransactionState.EMPTY,
txnTimeout: Int = transactionTimeoutMs): TransactionMetadata = {
val timestamp = time.milliseconds()
new TransactionMetadata(transactionalId, producerId, RecordBatch.NO_PRODUCER_ID, RecordBatch.NO_PRODUCER_ID, 0.toShort,
RecordBatch.NO_PRODUCER_EPOCH, txnTimeout, state, util.Set.of, timestamp, timestamp, TV_0)
new TransactionMetadata(transactionalId, producerId, RecordBatch.NO_PRODUCER_ID, RecordBatch.NO_PRODUCER_ID, 0.toShort, RecordBatch.NO_PRODUCER_EPOCH, txnTimeout, state, timestamp, timestamp, TV_0)
}

private def prepareTxnLog(topicPartition: TopicPartition,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,6 @@ public static boolean isEpochExhausted(short producerEpoch) {
* @param lastProducerEpoch last epoch of the producer
* @param txnTimeoutMs timeout to be used to abort long running transactions
* @param state current state of the transaction
* @param topicPartitions current set of partitions that are part of this transaction
* @param txnStartTimestamp time the transaction was started, i.e., when first partition is added
* @param txnLastUpdateTimestamp updated when any operation updates the TransactionMetadata. To be used for expiration
* @param clientTransactionVersion TransactionVersion used by the client when the state was transitioned
Expand All @@ -87,7 +86,6 @@ public TransactionMetadata(String transactionalId,
short lastProducerEpoch,
int txnTimeoutMs,
TransactionState state,
Set<TopicPartition> topicPartitions,
long txnStartTimestamp,
long txnLastUpdateTimestamp,
TransactionVersion clientTransactionVersion) {
Expand All @@ -99,7 +97,7 @@ public TransactionMetadata(String transactionalId,
this.lastProducerEpoch = lastProducerEpoch;
this.txnTimeoutMs = txnTimeoutMs;
this.state = state;
this.topicPartitions = new HashSet<>(topicPartitions);
this.topicPartitions = new HashSet<>();
this.txnStartTimestamp = txnStartTimestamp;
this.txnLastUpdateTimestamp = txnLastUpdateTimestamp;
this.clientTransactionVersion = clientTransactionVersion;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
import java.util.HashSet;

/**
* Immutable object representing the target transition of the transaction metadata
* Represent the target transition of the transaction metadata. The topicPartitions field is mutable.
*/
public record TxnTransitMetadata(
long producerId,
Expand Down