Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -1243,15 +1243,16 @@ static String niceMemoryUnits(long bytes) {
break;
}
}
String resultFormat = " (" + value + " %s" + (value == 1 ? ")" : "s)");
switch (i) {
case 1:
return " (" + value + " kibibyte" + (value == 1 ? ")" : "s)");
return String.format(resultFormat, "kibibyte");
case 2:
return " (" + value + " mebibyte" + (value == 1 ? ")" : "s)");
return String.format(resultFormat, "mebibyte");
case 3:
return " (" + value + " gibibyte" + (value == 1 ? ")" : "s)");
return String.format(resultFormat, "gibibyte");
case 4:
return " (" + value + " tebibyte" + (value == 1 ? ")" : "s)");
return String.format(resultFormat, "tebibyte");
default:
return "";
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -502,6 +502,7 @@ public void addMetric(MetricName metricName, MetricConfig config, Measurable mea
*
* @param metricName The name of the metric
* @param metricValueProvider The metric value provider associated with this metric
* @throws IllegalArgumentException if a metric with same name already exists.
*/
public void addMetric(MetricName metricName, MetricConfig config, MetricValueProvider<?> metricValueProvider) {
KafkaMetric m = new KafkaMetric(new Object(),
Expand Down Expand Up @@ -587,18 +588,19 @@ public synchronized void removeReporter(MetricsReporter reporter) {
}

/**
* Register a metric if not present or return an already existing metric otherwise.
* Register a metric if not present or return the already existing metric with the same name.
* When a metric is newly registered, this method returns null
*
* @param metric The KafkaMetric to register
* @return KafkaMetric if the metric already exists, null otherwise
* @return the existing metric with the same name or null
*/
synchronized KafkaMetric registerMetric(KafkaMetric metric) {
MetricName metricName = metric.metricName();
if (this.metrics.containsKey(metricName)) {
return this.metrics.get(metricName);
KafkaMetric existingMetric = this.metrics.putIfAbsent(metricName, metric);
if (existingMetric != null) {
return existingMetric;
}
this.metrics.put(metricName, metric);
// newly added metric
for (MetricsReporter reporter : reporters) {
try {
reporter.metricChange(metric);
Expand Down
15 changes: 13 additions & 2 deletions core/src/main/scala/kafka/log/LogCleaner.scala
Original file line number Diff line number Diff line change
Expand Up @@ -104,7 +104,7 @@ class LogCleaner(initialConfig: CleanerConfig,
private[log] val cleanerManager = new LogCleanerManager(logDirs, logs, logDirFailureChannel)

/* a throttle used to limit the I/O of all the cleaner threads to a user-specified maximum rate */
private val throttler = new Throttler(desiredRatePerSec = config.maxIoBytesPerSecond,
private[log] val throttler = new Throttler(desiredRatePerSec = config.maxIoBytesPerSecond,
checkIntervalMs = 300,
throttleDown = true,
"cleaner-io",
Expand Down Expand Up @@ -186,11 +186,20 @@ class LogCleaner(initialConfig: CleanerConfig,
}

/**
* Reconfigure log clean config. This simply stops current log cleaners and creates new ones.
* Reconfigure log clean config. The will:
* 1. update desiredRatePerSec in Throttler with logCleanerIoMaxBytesPerSecond, if necessary
* 2. stop current log cleaners and create new ones.
* That ensures that if any of the cleaners had failed, new cleaners are created to match the new config.
*/
override def reconfigure(oldConfig: KafkaConfig, newConfig: KafkaConfig): Unit = {
config = LogCleaner.cleanerConfig(newConfig)

val maxIoBytesPerSecond = config.maxIoBytesPerSecond;
if (maxIoBytesPerSecond != oldConfig.logCleanerIoMaxBytesPerSecond) {
info(s"Updating logCleanerIoMaxBytesPerSecond: $maxIoBytesPerSecond")
throttler.updateDesiredRatePerSec(maxIoBytesPerSecond)
}

shutdown()
startup()
}
Expand Down Expand Up @@ -692,6 +701,8 @@ private[log] class Cleaner(val id: Int,
if (discardBatchRecords)
// The batch is only retained to preserve producer sequence information; the records can be removed
false
else if (batch.isControlBatch)
true
else
Cleaner.this.shouldRetainRecord(map, retainLegacyDeletesAndTxnMarkers, batch, record, stats, currentTime = currentTime)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,8 @@ object AclAuthorizer {
private def validateAclBinding(aclBinding: AclBinding): Unit = {
if (aclBinding.isUnknown)
throw new IllegalArgumentException("ACL binding contains unknown elements")
if (aclBinding.pattern().name().contains("/"))
throw new IllegalArgumentException(s"ACL binding contains invalid resource name: ${aclBinding.pattern().name()}")
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,12 +19,11 @@ package kafka.server

import java.util.concurrent.LinkedBlockingDeque
import java.util.concurrent.atomic.AtomicReference

import kafka.common.{InterBrokerSendThread, RequestAndCompletionHandler}
import kafka.raft.RaftManager
import kafka.utils.Logging
import org.apache.kafka.clients._
import org.apache.kafka.common.Node
import org.apache.kafka.common.{Node, Reconfigurable}
import org.apache.kafka.common.metrics.Metrics
import org.apache.kafka.common.network._
import org.apache.kafka.common.protocol.Errors
Expand Down Expand Up @@ -188,6 +187,10 @@ class BrokerToControllerChannelManagerImpl(
config.saslInterBrokerHandshakeRequestEnable,
logContext
)
channelBuilder match {
case reconfigurable: Reconfigurable => config.addReconfigurable(reconfigurable)
case _ =>
}
val selector = new Selector(
NetworkReceive.UNLIMITED,
Selector.NO_IDLE_TIMEOUT_MS,
Expand Down
12 changes: 8 additions & 4 deletions core/src/main/scala/kafka/utils/Throttler.scala
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ import scala.math._
* @param time: The time implementation to use
*/
@threadsafe
class Throttler(desiredRatePerSec: Double,
class Throttler(@volatile var desiredRatePerSec: Double,
checkIntervalMs: Long = 100L,
throttleDown: Boolean = true,
metricName: String = "throttler",
Expand All @@ -52,6 +52,7 @@ class Throttler(desiredRatePerSec: Double,
def maybeThrottle(observed: Double): Unit = {
val msPerSec = TimeUnit.SECONDS.toMillis(1)
val nsPerSec = TimeUnit.SECONDS.toNanos(1)
val currentDesiredRatePerSec = desiredRatePerSec;

meter.mark(observed.toLong)
lock synchronized {
Expand All @@ -62,14 +63,14 @@ class Throttler(desiredRatePerSec: Double,
// we should take a little nap
if (elapsedNs > checkIntervalNs && observedSoFar > 0) {
val rateInSecs = (observedSoFar * nsPerSec) / elapsedNs
val needAdjustment = !(throttleDown ^ (rateInSecs > desiredRatePerSec))
val needAdjustment = !(throttleDown ^ (rateInSecs > currentDesiredRatePerSec))
if (needAdjustment) {
// solve for the amount of time to sleep to make us hit the desired rate
val desiredRateMs = desiredRatePerSec / msPerSec.toDouble
val desiredRateMs = currentDesiredRatePerSec / msPerSec.toDouble
val elapsedMs = TimeUnit.NANOSECONDS.toMillis(elapsedNs)
val sleepTime = round(observedSoFar / desiredRateMs - elapsedMs)
if (sleepTime > 0) {
trace("Natural rate is %f per second but desired rate is %f, sleeping for %d ms to compensate.".format(rateInSecs, desiredRatePerSec, sleepTime))
trace("Natural rate is %f per second but desired rate is %f, sleeping for %d ms to compensate.".format(rateInSecs, currentDesiredRatePerSec, sleepTime))
time.sleep(sleepTime)
}
}
Expand All @@ -79,6 +80,9 @@ class Throttler(desiredRatePerSec: Double,
}
}

def updateDesiredRatePerSec(updatedDesiredRatePerSec: Double): Unit = {
desiredRatePerSec = updatedDesiredRatePerSec;
}
}

object Throttler {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,6 @@ import java.time.Duration
import java.util
import java.util.{Collections, Properties}
import java.util.concurrent._

import javax.management.ObjectName
import com.yammer.metrics.core.MetricName
import kafka.admin.ConfigCommand
Expand All @@ -35,9 +34,9 @@ import kafka.controller.{ControllerBrokerStateInfo, ControllerChannelManager}
import kafka.log.{CleanerConfig, LogConfig}
import kafka.message.ProducerCompressionCodec
import kafka.network.{Processor, RequestChannel}
import kafka.server.QuorumTestHarness
import kafka.utils._
import kafka.utils.Implicits._
import kafka.utils.TestUtils.TestControllerRequestCompletionHandler
import kafka.zk.ConfigEntityChangeNotificationZNode
import org.apache.kafka.clients.CommonClientConfigs
import org.apache.kafka.clients.admin.AlterConfigOp.OpType
Expand All @@ -52,10 +51,12 @@ import org.apache.kafka.common.config.types.Password
import org.apache.kafka.common.config.provider.FileConfigProvider
import org.apache.kafka.common.errors.{AuthenticationException, InvalidRequestException}
import org.apache.kafka.common.internals.Topic
import org.apache.kafka.common.message.MetadataRequestData
import org.apache.kafka.common.metrics.{KafkaMetric, MetricsContext, MetricsReporter, Quota}
import org.apache.kafka.common.network.{ListenerName, Mode}
import org.apache.kafka.common.network.CertStores.{KEYSTORE_PROPS, TRUSTSTORE_PROPS}
import org.apache.kafka.common.record.TimestampType
import org.apache.kafka.common.requests.MetadataRequest
import org.apache.kafka.common.security.auth.SecurityProtocol
import org.apache.kafka.common.security.scram.ScramCredential
import org.apache.kafka.common.serialization.{StringDeserializer, StringSerializer}
Expand Down Expand Up @@ -429,6 +430,23 @@ class DynamicBrokerReconfigurationTest extends QuorumTestHarness with SaslSetup
verifyProduceConsume(producer, consumer, 10, topic)
}

def verifyBrokerToControllerCall(controller: KafkaServer): Unit = {
val nonControllerBroker = servers.find(_.config.brokerId != controller.config.brokerId).get
val brokerToControllerManager = nonControllerBroker.clientToControllerChannelManager
val completionHandler = new TestControllerRequestCompletionHandler()
brokerToControllerManager.sendRequest(new MetadataRequest.Builder(new MetadataRequestData()), completionHandler)
TestUtils.waitUntilTrue(() => {
completionHandler.completed.get() || completionHandler.timedOut.get()
}, "Timed out while waiting for broker to controller API call")
// we do not expect a timeout from broker to controller request
assertFalse(completionHandler.timedOut.get(), "broker to controller request is timeout")
assertTrue(completionHandler.actualResponse.isDefined, "No response recorded even though request is completed")
val response = completionHandler.actualResponse.get
assertNull(response.authenticationException(), s"Request failed due to authentication error ${response.authenticationException}")
assertNull(response.versionMismatch(), s"Request failed due to unsupported version error ${response.versionMismatch}")
assertFalse(response.wasDisconnected(), "Request failed because broker is not available")
}

// Produce/consume should work with old as well as new client keystore
verifySslProduceConsume(sslProperties1, "alter-truststore-1")
verifySslProduceConsume(sslProperties2, "alter-truststore-2")
Expand Down Expand Up @@ -469,6 +487,9 @@ class DynamicBrokerReconfigurationTest extends QuorumTestHarness with SaslSetup
JTestUtils.fieldValue(controllerChannelManager, classOf[ControllerChannelManager], "brokerStateInfo")
brokerStateInfo(0).networkClient.disconnect("0")
TestUtils.createTopic(zkClient, "testtopic2", numPartitions, replicationFactor = numServers, servers)

// validate that the brokerToController request works fine
verifyBrokerToControllerCall(controller)
}

@Test
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,13 +19,14 @@ package kafka.server

import java.nio.ByteBuffer
import java.util.Collections
import java.util.concurrent.atomic.{AtomicBoolean, AtomicReference}
import java.util.concurrent.atomic.AtomicReference
import kafka.utils.TestUtils
import kafka.utils.TestUtils.TestControllerRequestCompletionHandler
import org.apache.kafka.clients.{ClientResponse, ManualMetadataUpdater, Metadata, MockClient, NodeApiVersions}
import org.apache.kafka.common.Node
import org.apache.kafka.common.message.{EnvelopeResponseData, MetadataRequestData}
import org.apache.kafka.common.protocol.{ApiKeys, Errors}
import org.apache.kafka.common.requests.{AbstractRequest, EnvelopeRequest, EnvelopeResponse, MetadataRequest, MetadataResponse, RequestTestUtils}
import org.apache.kafka.common.requests.{AbstractRequest, EnvelopeRequest, EnvelopeResponse, MetadataRequest, RequestTestUtils}
import org.apache.kafka.common.security.auth.KafkaPrincipal
import org.apache.kafka.common.security.authenticator.DefaultKafkaPrincipalBuilder
import org.apache.kafka.common.utils.MockTime
Expand All @@ -51,7 +52,7 @@ class BrokerToControllerRequestThreadTest {
config, time, "", retryTimeoutMs)
testRequestThread.started = true

val completionHandler = new TestRequestCompletionHandler(None)
val completionHandler = new TestControllerRequestCompletionHandler(None)
val queueItem = BrokerToControllerQueueItem(
time.milliseconds(),
new MetadataRequest.Builder(new MetadataRequestData()),
Expand Down Expand Up @@ -89,7 +90,7 @@ class BrokerToControllerRequestThreadTest {
testRequestThread.started = true
mockClient.prepareResponse(expectedResponse)

val completionHandler = new TestRequestCompletionHandler(Some(expectedResponse))
val completionHandler = new TestControllerRequestCompletionHandler(Some(expectedResponse))
val queueItem = BrokerToControllerQueueItem(
time.milliseconds(),
new MetadataRequest.Builder(new MetadataRequestData()),
Expand Down Expand Up @@ -130,7 +131,7 @@ class BrokerToControllerRequestThreadTest {
controllerNodeProvider, config, time, "", retryTimeoutMs = Long.MaxValue)
testRequestThread.started = true

val completionHandler = new TestRequestCompletionHandler(Some(expectedResponse))
val completionHandler = new TestControllerRequestCompletionHandler(Some(expectedResponse))
val queueItem = BrokerToControllerQueueItem(
time.milliseconds(),
new MetadataRequest.Builder(new MetadataRequestData()),
Expand Down Expand Up @@ -180,7 +181,7 @@ class BrokerToControllerRequestThreadTest {
config, time, "", retryTimeoutMs = Long.MaxValue)
testRequestThread.started = true

val completionHandler = new TestRequestCompletionHandler(Some(expectedResponse))
val completionHandler = new TestControllerRequestCompletionHandler(Some(expectedResponse))
val queueItem = BrokerToControllerQueueItem(
time.milliseconds(),
new MetadataRequest.Builder(new MetadataRequestData()
Expand Down Expand Up @@ -243,7 +244,7 @@ class BrokerToControllerRequestThreadTest {
config, time, "", retryTimeoutMs = Long.MaxValue)
testRequestThread.started = true

val completionHandler = new TestRequestCompletionHandler(Some(expectedResponse))
val completionHandler = new TestControllerRequestCompletionHandler(Some(expectedResponse))
val kafkaPrincipal = new KafkaPrincipal(KafkaPrincipal.USER_TYPE, "principal", true)
val kafkaPrincipalBuilder = new DefaultKafkaPrincipalBuilder(null, null)

Expand Down Expand Up @@ -305,7 +306,7 @@ class BrokerToControllerRequestThreadTest {
config, time, "", retryTimeoutMs)
testRequestThread.started = true

val completionHandler = new TestRequestCompletionHandler()
val completionHandler = new TestControllerRequestCompletionHandler()
val queueItem = BrokerToControllerQueueItem(
time.milliseconds(),
new MetadataRequest.Builder(new MetadataRequestData()
Expand Down Expand Up @@ -419,7 +420,7 @@ class BrokerToControllerRequestThreadTest {
val testRequestThread = new BrokerToControllerRequestThread(mockClient, new ManualMetadataUpdater(), controllerNodeProvider,
config, time, "", retryTimeoutMs = Long.MaxValue)

val completionHandler = new TestRequestCompletionHandler(None)
val completionHandler = new TestControllerRequestCompletionHandler(None)
val queueItem = BrokerToControllerQueueItem(
time.milliseconds(),
new MetadataRequest.Builder(new MetadataRequestData()),
Expand All @@ -445,22 +446,4 @@ class BrokerToControllerRequestThreadTest {
fail(s"Condition failed to be met after polling $tries times")
}
}

class TestRequestCompletionHandler(
expectedResponse: Option[MetadataResponse] = None
) extends ControllerRequestCompletionHandler {
val completed: AtomicBoolean = new AtomicBoolean(false)
val timedOut: AtomicBoolean = new AtomicBoolean(false)

override def onComplete(response: ClientResponse): Unit = {
expectedResponse.foreach { expected =>
assertEquals(expected, response.responseBody())
}
completed.set(true)
}

override def onTimeout(): Unit = {
timedOut.set(true)
}
}
}
Loading