From 5fa856ba1cc8794b1af1c53c767818994a63081e Mon Sep 17 00:00:00 2001 From: Ming Liu Date: Sun, 2 Dec 2018 21:44:14 -0800 Subject: [PATCH 1/4] KAFKA-7692: Fix ProduceStateManager SequenceNumber overflow. --- .../scala/kafka/log/ProducerStateManager.scala | 6 +++--- .../unit/kafka/log/ProducerStateManagerTest.scala | 14 +++++++------- 2 files changed, 10 insertions(+), 10 deletions(-) diff --git a/core/src/main/scala/kafka/log/ProducerStateManager.scala b/core/src/main/scala/kafka/log/ProducerStateManager.scala index a5c182c2ce306..cefc8fe031b69 100644 --- a/core/src/main/scala/kafka/log/ProducerStateManager.scala +++ b/core/src/main/scala/kafka/log/ProducerStateManager.scala @@ -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..5bb09ef86807e 100644 --- a/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala +++ b/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala @@ -191,7 +191,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 +207,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 +224,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 +253,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 +263,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 +819,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) } From d16ed7691fbac2a6c353b7e21e830ed3d0f6ef6e Mon Sep 17 00:00:00 2001 From: Ming Liu Date: Fri, 11 Jan 2019 10:46:08 -0800 Subject: [PATCH 2/4] Add unittest for producerStateManager test. --- .../kafka/common/record/DefaultRecordBatch.java | 8 +++++++- .../scala/kafka/log/ProducerStateManager.scala | 4 ++-- .../kafka/log/ProducerStateManagerTest.scala | 17 +++++++++++++++++ 3 files changed, 26 insertions(+), 3 deletions(-) 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 71e668e45da01..74a3babad3594 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 @@ -523,12 +523,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 cefc8fe031b69..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 = { diff --git a/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala b/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala index 5bb09ef86807e..c8c9504790fed 100644 --- a/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala +++ b/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala @@ -131,6 +131,23 @@ 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) + + var lastEntry = stateManager.lastEntry(producerId).get + assertEquals(Int.MaxValue-10, lastEntry.firstSeq) + assertEquals(9, lastEntry.lastSeq) + assertEquals(2020L, lastEntry.lastDataOffset) + } + @Test(expected = classOf[OutOfOrderSequenceException]) def testProducerSequenceInvalidWrapAround(): Unit = { val epoch = 15.toShort From f4bc8a2342989906fbfa302220804c69ca143393 Mon Sep 17 00:00:00 2001 From: Ming Liu Date: Sun, 13 Jan 2019 18:06:55 -0800 Subject: [PATCH 3/4] Minor update. --- .../src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala | 1 + 1 file changed, 1 insertion(+) diff --git a/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala b/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala index c8c9504790fed..be309dfd84050 100644 --- a/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala +++ b/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala @@ -145,6 +145,7 @@ class ProducerStateManagerTest extends JUnitSuite { var lastEntry = stateManager.lastEntry(producerId).get assertEquals(Int.MaxValue-10, lastEntry.firstSeq) assertEquals(9, lastEntry.lastSeq) + assertEquals(2000L, lastEntry.firstOffset) assertEquals(2020L, lastEntry.lastDataOffset) } From a1f07cabe172340d9cf32fa28c0f78f803140363 Mon Sep 17 00:00:00 2001 From: Ming Liu Date: Sun, 13 Jan 2019 20:59:37 -0800 Subject: [PATCH 4/4] Use val instead of var. --- .../test/scala/unit/kafka/log/ProducerStateManagerTest.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala b/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala index be309dfd84050..a2abf7b3f5aea 100644 --- a/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala +++ b/core/src/test/scala/unit/kafka/log/ProducerStateManagerTest.scala @@ -142,7 +142,7 @@ class ProducerStateManagerTest extends JUnitSuite { stateManager.update(appendInfo) assertTrue(stateManager.lastEntry(producerId).isDefined) - var lastEntry = stateManager.lastEntry(producerId).get + val lastEntry = stateManager.lastEntry(producerId).get assertEquals(Int.MaxValue-10, lastEntry.firstSeq) assertEquals(9, lastEntry.lastSeq) assertEquals(2000L, lastEntry.firstOffset)