Skip to content
22 changes: 16 additions & 6 deletions core/src/main/scala/kafka/server/AlterIsrManager.scala
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@ package kafka.server
import java.util
import java.util.concurrent.atomic.{AtomicBoolean, AtomicLong}
import java.util.concurrent.{ConcurrentHashMap, TimeUnit}

import kafka.api.LeaderAndIsr
import kafka.metrics.KafkaMetricsGroup
import kafka.utils.{Logging, Scheduler}
Expand Down Expand Up @@ -62,27 +61,38 @@ class AlterIsrManagerImpl(val controllerChannelManager: BrokerToControllerChanne

private val lastIsrPropagationMs = new AtomicLong(0)

override def start(): Unit = {
scheduler.schedule("send-alter-isr", propagateIsrChanges, 50, 50, TimeUnit.MILLISECONDS)
}
override def start(): Unit = { }

override def enqueue(alterIsrItem: AlterIsrItem): Boolean = {
unsentIsrUpdates.putIfAbsent(alterIsrItem.topicPartition, alterIsrItem) == null
if (unsentIsrUpdates.putIfAbsent(alterIsrItem.topicPartition, alterIsrItem) == null) {
if (inflightRequest.compareAndSet(false, true)) {
Comment thread
mumrah marked this conversation as resolved.
Outdated
// optimistically set the inflight flag even though we haven't sent the request yet
scheduler.schedule("send-alter-isr", propagateIsrChanges, 50, -1, TimeUnit.MILLISECONDS)
Comment thread
mumrah marked this conversation as resolved.
Outdated
}
true
} else {
false
}

}

override def clearPending(topicPartition: TopicPartition): Unit = {
unsentIsrUpdates.remove(topicPartition)
}

private def propagateIsrChanges(): Unit = {
if (!unsentIsrUpdates.isEmpty && inflightRequest.compareAndSet(false, true)) {
// Updates could have been cleared by new LeaderAndIsr, so check again
if (!unsentIsrUpdates.isEmpty) {
// Copy current unsent ISRs but don't remove from the map
val inflightAlterIsrItems = new ListBuffer[AlterIsrItem]()
unsentIsrUpdates.values().forEach(item => inflightAlterIsrItems.append(item))

val now = time.milliseconds()
lastIsrPropagationMs.set(now)
sendRequest(inflightAlterIsrItems.toSeq)
} else {
// Never sent a request, so clear the flag
inflightRequest.set(false)
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -254,6 +254,10 @@ class AlterIsrManagerTest {
EasyMock.expect(brokerToController.sendRequest(EasyMock.anyObject(), EasyMock.capture(callbackCapture), EasyMock.eq(requestTimeout))).once()
EasyMock.replay(brokerToController)

// Need to re-enqueue again to trigger the thread to be scheduled
alterIsrManager.clearPending(tp2)
alterIsrManager.enqueue(AlterIsrItem(tp2, new LeaderAndIsr(1, 1, List(1,2,3), 10), _ => {}))

time.sleep(100)
scheduler.tick()
EasyMock.verify(brokerToController)
Expand Down