From 86c65f30b3947e01975737d0a8c7271ed79200c4 Mon Sep 17 00:00:00 2001 From: lambdaliu Date: Wed, 24 Oct 2018 14:14:59 +0800 Subject: [PATCH 1/3] FetchResponse should return lastStableOffset when no neead to convert data --- core/src/main/scala/kafka/server/KafkaApis.scala | 2 +- core/src/main/scala/kafka/server/ReplicaManager.scala | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 6a99db31a14af..86677674d044d 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -582,7 +582,7 @@ class KafkaApis(val requestChannel: RequestChannel, } } case None => new FetchResponse.PartitionData[BaseRecords](partitionData.error, partitionData.highWatermark, - FetchResponse.INVALID_LAST_STABLE_OFFSET, partitionData.logStartOffset, partitionData.abortedTransactions, + partitionData.lastStableOffset, partitionData.logStartOffset, partitionData.abortedTransactions, unconvertedRecords) } } diff --git a/core/src/main/scala/kafka/server/ReplicaManager.scala b/core/src/main/scala/kafka/server/ReplicaManager.scala index d5a3c68d5011f..e3feb7194c6f2 100644 --- a/core/src/main/scala/kafka/server/ReplicaManager.scala +++ b/core/src/main/scala/kafka/server/ReplicaManager.scala @@ -95,7 +95,7 @@ case class LogReadResult(info: FetchDataInfo, override def toString = s"Fetch Data: [$info], HW: [$highWatermark], leaderLogStartOffset: [$leaderLogStartOffset], leaderLogEndOffset: [$leaderLogEndOffset], " + - s"followerLogStartOffset: [$followerLogStartOffset], fetchTimeMs: [$fetchTimeMs], readSize: [$readSize], error: [$error]" + s"followerLogStartOffset: [$followerLogStartOffset], fetchTimeMs: [$fetchTimeMs], readSize: [$readSize], lastStableOffset: [$lastStableOffset], error: [$error]" } From d5af1ddb76322f963814d5a2f18694b4673bf9c0 Mon Sep 17 00:00:00 2001 From: lambdaliu Date: Thu, 25 Oct 2018 09:05:42 +0800 Subject: [PATCH 2/3] also return lastStatbleOffset when there is need to convert data --- core/src/main/scala/kafka/server/KafkaApis.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 86677674d044d..e3dc9217d45bf 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -573,7 +573,7 @@ class KafkaApis(val requestChannel: RequestChannel, // down-conversion always guarantees that at least one batch of messages is down-converted and sent out to the // client. new FetchResponse.PartitionData[BaseRecords](partitionData.error, partitionData.highWatermark, - FetchResponse.INVALID_LAST_STABLE_OFFSET, partitionData.logStartOffset, partitionData.abortedTransactions, + partitionData.lastStableOffset, partitionData.logStartOffset, partitionData.abortedTransactions, new LazyDownConversionRecords(tp, unconvertedRecords, magic, fetchContext.getFetchOffset(tp).get, time)) } catch { case e: UnsupportedCompressionTypeException => From 11be006acab988c250f2cd367a2523d278298bec Mon Sep 17 00:00:00 2001 From: lambdaliu Date: Thu, 25 Oct 2018 09:45:15 +0800 Subject: [PATCH 3/3] add tests --- .../kafka/api/PlaintextConsumerTest.scala | 27 +++++++++++++++++++ .../unit/kafka/server/FetchRequestTest.scala | 18 ++++++++++++- 2 files changed, 44 insertions(+), 1 deletion(-) diff --git a/core/src/test/scala/integration/kafka/api/PlaintextConsumerTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextConsumerTest.scala index 522ca49d3b944..a23513f498808 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextConsumerTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextConsumerTest.scala @@ -1586,6 +1586,33 @@ class PlaintextConsumerTest extends BaseConsumerTest { assertNull(consumer.metrics.get(new MetricName("records-lag", "consumer-fetch-manager-metrics", "", tags))) } + @Test + def testPerPartitionLagMetricsWhenReadCommitted() { + val numMessages = 1000 + // send some messages. + val producer = createProducer() + sendRecords(producer, numMessages, tp) + sendRecords(producer, numMessages, tp2) + + consumerConfig.setProperty(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed") + consumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "testPerPartitionLagMetricsCleanUpWithAssign") + consumerConfig.setProperty(ConsumerConfig.CLIENT_ID_CONFIG, "testPerPartitionLagMetricsCleanUpWithAssign") + val consumer = createConsumer() + consumer.assign(List(tp).asJava) + var records: ConsumerRecords[Array[Byte], Array[Byte]] = ConsumerRecords.empty() + TestUtils.waitUntilTrue(() => { + records = consumer.poll(100) + !records.records(tp).isEmpty + }, "Consumer did not consume any message before timeout.") + // Verify the metric exist. + val tags = new util.HashMap[String, String]() + tags.put("client-id", "testPerPartitionLagMetricsCleanUpWithAssign") + tags.put("topic", tp.topic()) + tags.put("partition", String.valueOf(tp.partition())) + val fetchLag = consumer.metrics.get(new MetricName("records-lag", "consumer-fetch-manager-metrics", "", tags)) + assertNotNull(fetchLag) + } + @Test def testPerPartitionLeadWithMaxPollRecords() { val numMessages = 1000 diff --git a/core/src/test/scala/unit/kafka/server/FetchRequestTest.scala b/core/src/test/scala/unit/kafka/server/FetchRequestTest.scala index 72a28549cf584..86f3314d62401 100644 --- a/core/src/test/scala/unit/kafka/server/FetchRequestTest.scala +++ b/core/src/test/scala/unit/kafka/server/FetchRequestTest.scala @@ -28,7 +28,7 @@ import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord} import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.protocol.{ApiKeys, Errors} import org.apache.kafka.common.record.{MemoryRecords, Record, RecordBatch} -import org.apache.kafka.common.requests.{FetchRequest, FetchResponse, FetchMetadata => JFetchMetadata} +import org.apache.kafka.common.requests.{FetchRequest, FetchResponse, IsolationLevel, FetchMetadata => JFetchMetadata} import org.apache.kafka.common.serialization.{ByteArraySerializer, StringSerializer} import org.junit.Assert._ import org.junit.Test @@ -171,6 +171,22 @@ class FetchRequestTest extends BaseRequestTest { assertEquals(0, records(partitionData).map(_.sizeInBytes).sum) } + @Test + def testFetchRequestV4WithReadCommitted(): Unit = { + initProducer() + val maxPartitionBytes = 200 + val (topicPartition, leaderId) = createTopics(numTopics = 1, numPartitions = 1).head + producer.send(new ProducerRecord(topicPartition.topic, topicPartition.partition, + "key", new String(new Array[Byte](maxPartitionBytes + 1)))).get + val fetchRequest = FetchRequest.Builder.forConsumer(Int.MaxValue, 0, createPartitionMap(maxPartitionBytes, + Seq(topicPartition))).isolationLevel(IsolationLevel.READ_COMMITTED).build(4) + val fetchResponse = sendFetchRequest(leaderId, fetchRequest) + val partitionData = fetchResponse.responseData.get(topicPartition) + assertEquals(Errors.NONE, partitionData.error) + assertTrue(partitionData.lastStableOffset > 0) + assertTrue(records(partitionData).map(_.sizeInBytes).sum > 0) + } + @Test def testFetchRequestToNonReplica(): Unit = { val topic = "topic"