diff --git a/clients/src/main/java/org/apache/kafka/common/record/DefaultRecordBatch.java b/clients/src/main/java/org/apache/kafka/common/record/DefaultRecordBatch.java index 5156c64b5eb42..19ddb0ef3fa1a 100644 --- a/clients/src/main/java/org/apache/kafka/common/record/DefaultRecordBatch.java +++ b/clients/src/main/java/org/apache/kafka/common/record/DefaultRecordBatch.java @@ -522,12 +522,18 @@ static int estimateBatchSizeUpperBound(ByteBuffer key, ByteBuffer value, Header[ return RECORD_BATCH_OVERHEAD + DefaultRecord.recordSizeUpperBound(key, value, headers); } - static int incrementSequence(int baseSequence, int increment) { + public static int incrementSequence(int baseSequence, int increment) { if (baseSequence > Integer.MAX_VALUE - increment) return increment - (Integer.MAX_VALUE - baseSequence) - 1; return baseSequence + increment; } + public static int decrementSequence(int baseSequence, int decrement) { + if (baseSequence < decrement) + return Integer.MAX_VALUE - (decrement - baseSequence) + 1; + return baseSequence - decrement; + } + private abstract class RecordIterator implements CloseableIterator { private final Long logAppendTime; private final long baseOffset; diff --git a/core/src/main/scala/kafka/log/ProducerStateManager.scala b/core/src/main/scala/kafka/log/ProducerStateManager.scala index a5c182c2ce306..a3b03d063f256 100644 --- a/core/src/main/scala/kafka/log/ProducerStateManager.scala +++ b/core/src/main/scala/kafka/log/ProducerStateManager.scala @@ -28,7 +28,7 @@ import org.apache.kafka.common.{KafkaException, TopicPartition} import org.apache.kafka.common.errors._ import org.apache.kafka.common.internals.Topic import org.apache.kafka.common.protocol.types._ -import org.apache.kafka.common.record.{ControlRecordType, EndTransactionMarker, RecordBatch} +import org.apache.kafka.common.record.{ControlRecordType, DefaultRecordBatch, EndTransactionMarker, RecordBatch} import org.apache.kafka.common.utils.{ByteUtils, Crc32C} import scala.collection.mutable.ListBuffer @@ -76,7 +76,7 @@ private[log] object ProducerStateEntry { } private[log] case class BatchMetadata(lastSeq: Int, lastOffset: Long, offsetDelta: Int, timestamp: Long) { - def firstSeq = lastSeq - offsetDelta + def firstSeq = DefaultRecordBatch.decrementSequence(lastSeq, offsetDelta) def firstOffset = lastOffset - offsetDelta override def toString: String = { @@ -265,7 +265,7 @@ private[log] class ProducerAppendInfo(val producerId: Long, None } } else { - append(batch.producerEpoch, batch.baseSequence, batch.lastSequence, batch.maxTimestamp, batch.lastOffset, + append(batch.producerEpoch, batch.baseSequence, batch.lastSequence, batch.maxTimestamp, batch.baseOffset, batch.lastOffset, batch.isTransactional) None } @@ -275,10 +275,11 @@ private[log] class ProducerAppendInfo(val producerId: Long, firstSeq: Int, lastSeq: Int, lastTimestamp: Long, + firstOffset: Long, lastOffset: Long, isTransactional: Boolean): Unit = { maybeValidateAppend(epoch, firstSeq) - updatedEntry.addBatch(epoch, lastSeq, lastOffset, lastSeq - firstSeq, lastTimestamp) + updatedEntry.addBatch(epoch, lastSeq, lastOffset, (lastOffset - firstOffset).toInt, lastTimestamp) updatedEntry.currentTxnFirstOffset match { case Some(_) if !isTransactional => @@ -287,7 +288,6 @@ private[log] class ProducerAppendInfo(val producerId: Long, case None if isTransactional => // Began a new transaction - val firstOffset = lastOffset - (lastSeq - firstSeq) updatedEntry.currentTxnFirstOffset = Some(firstOffset) transactions += new TxnMetadata(producerId, firstOffset) diff --git a/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala b/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala index b49b5e15d877b..a2abf7b3f5aea 100644 --- a/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala +++ b/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala @@ -131,6 +131,24 @@ class ProducerStateManagerTest extends JUnitSuite { assertEquals(0, lastEntry.lastSeq) } + @Test + def testProducerSequenceWithWrapAroundBatchRecord(): Unit = { + val epoch = 15.toShort + + val appendInfo = stateManager.prepareUpdate(producerId, isFromClient = false) + // Sequence number wrap around + appendInfo.append(epoch, Int.MaxValue-10, 9, time.milliseconds(), 2000L, 2020L, isTransactional = false) + assertEquals(None, stateManager.lastEntry(producerId)) + stateManager.update(appendInfo) + assertTrue(stateManager.lastEntry(producerId).isDefined) + + val lastEntry = stateManager.lastEntry(producerId).get + assertEquals(Int.MaxValue-10, lastEntry.firstSeq) + assertEquals(9, lastEntry.lastSeq) + assertEquals(2000L, lastEntry.firstOffset) + assertEquals(2020L, lastEntry.lastDataOffset) + } + @Test(expected = classOf[OutOfOrderSequenceException]) def testProducerSequenceInvalidWrapAround(): Unit = { val epoch = 15.toShort @@ -191,7 +209,7 @@ class ProducerStateManagerTest extends JUnitSuite { val offset = 992342L val seq = 0 val producerAppendInfo = new ProducerAppendInfo(producerId, ProducerStateEntry.empty(producerId), ValidationType.Full) - producerAppendInfo.append(producerEpoch, seq, seq, time.milliseconds(), offset, isTransactional = true) + producerAppendInfo.append(producerEpoch, seq, seq, time.milliseconds(), offset, offset, isTransactional = true) val logOffsetMetadata = new LogOffsetMetadata(messageOffset = offset, segmentBaseOffset = 990000L, relativePositionInSegment = 234224) @@ -207,7 +225,7 @@ class ProducerStateManagerTest extends JUnitSuite { val offset = 992342L val seq = 0 val producerAppendInfo = new ProducerAppendInfo(producerId, ProducerStateEntry.empty(producerId), ValidationType.Full) - producerAppendInfo.append(producerEpoch, seq, seq, time.milliseconds(), offset, isTransactional = true) + producerAppendInfo.append(producerEpoch, seq, seq, time.milliseconds(), offset, offset, isTransactional = true) // use some other offset to simulate a follower append where the log offset metadata won't typically // match any of the transaction first offsets @@ -224,13 +242,13 @@ class ProducerStateManagerTest extends JUnitSuite { val producerEpoch = 0.toShort val appendInfo = stateManager.prepareUpdate(producerId, isFromClient = true) - appendInfo.append(producerEpoch, 0, 5, time.milliseconds(), 20L, isTransactional = false) + appendInfo.append(producerEpoch, 0, 5, time.milliseconds(), 15L, 20L, isTransactional = false) assertEquals(None, stateManager.lastEntry(producerId)) stateManager.update(appendInfo) assertTrue(stateManager.lastEntry(producerId).isDefined) val nextAppendInfo = stateManager.prepareUpdate(producerId, isFromClient = true) - nextAppendInfo.append(producerEpoch, 6, 10, time.milliseconds(), 30L, isTransactional = false) + nextAppendInfo.append(producerEpoch, 6, 10, time.milliseconds(), 26L, 30L, isTransactional = false) assertTrue(stateManager.lastEntry(producerId).isDefined) var lastEntry = stateManager.lastEntry(producerId).get @@ -253,7 +271,7 @@ class ProducerStateManagerTest extends JUnitSuite { append(stateManager, producerId, producerEpoch, 0, offset) val appendInfo = stateManager.prepareUpdate(producerId, isFromClient = true) - appendInfo.append(producerEpoch, 1, 5, time.milliseconds(), 20L, isTransactional = true) + appendInfo.append(producerEpoch, 1, 5, time.milliseconds(), 16L, 20L, isTransactional = true) var lastEntry = appendInfo.toEntry assertEquals(producerEpoch, lastEntry.producerEpoch) assertEquals(1, lastEntry.firstSeq) @@ -263,7 +281,7 @@ class ProducerStateManagerTest extends JUnitSuite { assertEquals(Some(16L), lastEntry.currentTxnFirstOffset) assertEquals(List(new TxnMetadata(producerId, 16L)), appendInfo.startedTransactions) - appendInfo.append(producerEpoch, 6, 10, time.milliseconds(), 30L, isTransactional = true) + appendInfo.append(producerEpoch, 6, 10, time.milliseconds(), 26L, 30L, isTransactional = true) lastEntry = appendInfo.toEntry assertEquals(producerEpoch, lastEntry.producerEpoch) assertEquals(1, lastEntry.firstSeq) @@ -819,7 +837,7 @@ class ProducerStateManagerTest extends JUnitSuite { isTransactional: Boolean = false, isFromClient : Boolean = true): Unit = { val producerAppendInfo = stateManager.prepareUpdate(producerId, isFromClient) - producerAppendInfo.append(producerEpoch, seq, seq, timestamp, offset, isTransactional) + producerAppendInfo.append(producerEpoch, seq, seq, timestamp, offset, offset, isTransactional) stateManager.update(producerAppendInfo) stateManager.updateMapEndOffset(offset + 1) }