Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 20 additions & 20 deletions core/src/main/scala/kafka/server/BrokerServer.scala
Original file line number Diff line number Diff line change
Expand Up @@ -312,23 +312,6 @@ class BrokerServer(
config, Some(clientToControllerChannelManager), None, None,
groupCoordinator, transactionCoordinator)

/* Add all reconfigurables for config change notification before starting the metadata listener */
config.dynamicConfig.addReconfigurables(this)

dynamicConfigHandlers = Map[String, ConfigHandler](
ConfigType.Topic -> new TopicConfigHandler(logManager, config, quotaManagers, None),
ConfigType.Broker -> new BrokerConfigHandler(config, quotaManagers))

if (!config.processRoles.contains(ControllerRole)) {
// If no controller is defined, we rely on the broker to generate snapshots.
metadataSnapshotter = Some(new BrokerMetadataSnapshotter(
config.nodeId,
time,
threadNamePrefix,
new BrokerSnapshotWriterBuilder(raftManager.client)
))
}

metadataListener = new BrokerMetadataListener(
config.nodeId,
time,
Expand Down Expand Up @@ -369,9 +352,6 @@ class BrokerServer(
featuresRemapped
)

// Register a listener with the Raft layer to receive metadata event notifications
raftManager.register(metadataListener)

val endpoints = new util.ArrayList[Endpoint](networkListeners.size())
var interBrokerListener: Endpoint = null
networkListeners.iterator().forEachRemaining(listener => {
Expand Down Expand Up @@ -408,6 +388,26 @@ class BrokerServer(
}.toMap
}

/* Add all reconfigurables for config change notification before starting the metadata listener */
config.dynamicConfig.addReconfigurables(this)

dynamicConfigHandlers = Map[String, ConfigHandler](
ConfigType.Topic -> new TopicConfigHandler(logManager, config, quotaManagers, None),
ConfigType.Broker -> new BrokerConfigHandler(config, quotaManagers))

if (!config.processRoles.contains(ControllerRole)) {
// If no controller is defined, we rely on the broker to generate snapshots.
metadataSnapshotter = Some(new BrokerMetadataSnapshotter(
config.nodeId,
time,
threadNamePrefix,
new BrokerSnapshotWriterBuilder(raftManager.client)
))
}

// Register a listener with the Raft layer to receive metadata event notifications
raftManager.register(metadataListener)

val fetchManager = new FetchManager(Time.SYSTEM,
new FetchSessionCache(config.maxIncrementalFetchSessionCacheSlots,
KafkaServer.MIN_INCREMENTAL_FETCH_SESSION_EVICTION_MS))
Expand Down
33 changes: 21 additions & 12 deletions core/src/main/scala/kafka/server/ControllerServer.scala
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,8 @@ import kafka.network.{DataPlaneAcceptor, SocketServer}
import kafka.raft.KafkaRaftManager
import kafka.security.CredentialProvider
import kafka.server.KafkaConfig.{AlterConfigPolicyClassNameProp, CreateTopicPolicyClassNameProp}
import kafka.server.KafkaRaftServer.BrokerRole
import kafka.server.QuotaFactory.QuotaManagers
import kafka.server.metadata.DynamicConfigPublisher
import kafka.utils.{CoreUtils, Logging}
import kafka.zk.{KafkaZkClient, ZkMigrationClient}
import org.apache.kafka.common.config.ConfigException
Expand All @@ -35,6 +35,7 @@ import org.apache.kafka.common.security.token.delegation.internals.DelegationTok
import org.apache.kafka.common.utils.LogContext
import org.apache.kafka.common.{ClusterResource, Endpoint}
import org.apache.kafka.controller.{Controller, QuorumController, QuorumFeatures}
import org.apache.kafka.image.publisher.MetadataPublisher
import org.apache.kafka.metadata.KafkaConfigSchema
import org.apache.kafka.metadata.authorizer.ClusterMetadataAuthorizer
import org.apache.kafka.metadata.bootstrap.BootstrapMetadata
Expand Down Expand Up @@ -104,6 +105,7 @@ class ControllerServer(
var quotaManagers: QuotaManagers = _
var controllerApis: ControllerApis = _
var controllerApisHandlerPool: KafkaRequestHandlerPool = _
def kafkaYammerMetrics: KafkaYammerMetrics = KafkaYammerMetrics.INSTANCE
var migrationSupport: Option[ControllerMigrationSupport] = None

private def maybeChangeStatus(from: ProcessStatus, to: ProcessStatus): Boolean = {
Expand All @@ -118,13 +120,6 @@ class ControllerServer(
true
}

private def doRemoteKraftSetup(): Unit = {
// Explicitly configure metric reporters on this remote controller.
// We do not yet support dynamic reconfiguration on remote controllers in general;
// remove this once that is implemented.
new DynamicMetricReporterState(config.nodeId, config, metrics, clusterId)
}

def clusterId: String = sharedServer.metaProps.clusterId

def startup(): Unit = {
Expand Down Expand Up @@ -242,11 +237,17 @@ class ControllerServer(
case _ => // nothing to do
}
controller = controllerBuilder.build()
val metadataPublishers = new java.util.ArrayList[MetadataPublisher]()

// Perform any setup that is done only when this node is a controller-only node.
if (!config.processRoles.contains(BrokerRole)) {
doRemoteKraftSetup()
}
val dynamicConfigHandlers = Map[String, ConfigHandler](
ConfigType.Broker -> new BrokerConfigHandler(config, quotaManagers)
)
metadataPublishers.add(new DynamicConfigPublisher(
config,
sharedServer.metadataPublishingFaultHandler,
dynamicConfigHandlers.toMap,
"controller"
))

if (config.migrationEnabled) {
val zkClient = KafkaZkClient.createZkClient("KRaft Migration", time, config, KafkaServer.zkClientConfigFromKafkaConfig(config))
Expand Down Expand Up @@ -290,6 +291,14 @@ class ControllerServer(
s"${DataPlaneAcceptor.MetricPrefix}RequestHandlerAvgIdlePercent",
DataPlaneAcceptor.ThreadPrefix)

config.dynamicConfig.addReconfigurables(this)

// Install the metadata publishers. Note that they will not actually receive any metadata
// until we catch up to the high water mark of __cluster_metadata-0. We are not waiting for
// that. We are just waiting for the installation process to complete.
FutureUtils.waitWithLogging(logger.underlying, "all of the metadata publishers to be installed",
sharedServer.loader.installPublishers(metadataPublishers), startupDeadline, time)

/**
* Enable the controller endpoint(s). If we are using an authorizer which stores
* ACLs in the metadata log, such as StandardAuthorizer, we will be able to start
Expand Down
29 changes: 29 additions & 0 deletions core/src/main/scala/kafka/server/DynamicBrokerConfig.scala
Original file line number Diff line number Diff line change
Expand Up @@ -271,6 +271,35 @@ class DynamicBrokerConfig(private val kafkaConfig: KafkaConfig) extends Logging
addBrokerReconfigurable(new DynamicProducerStateManagerConfig(kafkaServer.logManager.producerStateManagerConfig))
}

/**
* Add reconfigurables to be notified when a dynamic controller config is updated.
*/
def addReconfigurables(controller: ControllerServer): Unit = {
controller.authorizer match {
case Some(authz: Reconfigurable) => addReconfigurable(authz)
case _ =>
}
if (kafkaConfig.isKRaftCoResidentMode) {
debug("Relying on the broker to reconfigure the Yammer metrics.")
} else {
addReconfigurable(controller.kafkaYammerMetrics)
}
if (kafkaConfig.isKRaftCoResidentMode) {
debug("Relying on the broker to reconfigure the metric reporters.")
} else {
addReconfigurable(new DynamicMetricsReporters(
kafkaConfig.brokerId, controller.config, controller.metrics, controller.clusterId))
}
addBrokerReconfigurable(controller.socketServer)

// KAFKA-14349: add dynamic thread pool resizing here

// KAFKA-14350: add dynamic listener reconfiguration here

// When we implement controller mutation quotas in KRaft, as discussed in KAFKA-14351,
// we'll need to make them reconfigurable here (probably via DynamicClientQuotaCallback).
}

def addReconfigurable(reconfigurable: Reconfigurable): Unit = CoreUtils.inWriteLock(lock) {
verifyReconfigurableConfigs(reconfigurable.reconfigurableConfigs.asScala)
reconfigurables.add(reconfigurable)
Expand Down
5 changes: 1 addition & 4 deletions core/src/main/scala/kafka/server/KafkaRaftServer.scala
Original file line number Diff line number Diff line change
Expand Up @@ -77,10 +77,7 @@ class KafkaRaftServer(
)

private val broker: Option[BrokerServer] = if (config.processRoles.contains(BrokerRole)) {
Some(new BrokerServer(
sharedServer,
offlineDirs
))
Some(new BrokerServer(sharedServer, offlineDirs))
} else {
None
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -211,7 +211,7 @@ class BrokerMetadataPublisher(
}

// Apply configuration deltas.
dynamicConfigPublisher.publish(delta, newImage)
dynamicConfigPublisher.publish(delta, newImage, null)

// Apply client quotas delta.
try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import kafka.server.ConfigAdminManager.toLoggableProps
import kafka.server.{ConfigEntityName, ConfigHandler, ConfigType, KafkaConfig}
import kafka.utils.Logging
import org.apache.kafka.common.config.ConfigResource.Type.{BROKER, TOPIC}
import org.apache.kafka.image.loader.LoaderManifest
import org.apache.kafka.image.{MetadataDelta, MetadataImage}
import org.apache.kafka.server.fault.FaultHandler

Expand All @@ -31,10 +32,16 @@ class DynamicConfigPublisher(
faultHandler: FaultHandler,
dynamicConfigHandlers: Map[String, ConfigHandler],
nodeType: String
) extends Logging {
) extends org.apache.kafka.image.publisher.MetadataPublisher with Logging {
logIdent = s"[DynamicConfigPublisher nodeType=${nodeType} id=${conf.nodeId}] "

def publish(delta: MetadataDelta, newImage: MetadataImage): Unit = {
def name(): String = "DynamicConfigPublisher"

def publish(
delta: MetadataDelta,
newImage: MetadataImage,
manifest: LoaderManifest
): Unit = {
val deltaName = s"MetadataDelta up to ${newImage.highestOffsetAndEpoch().offset}"
try {
// Apply configuration deltas.
Expand Down
117 changes: 115 additions & 2 deletions core/src/test/scala/integration/kafka/server/KRaftClusterTest.scala
Original file line number Diff line number Diff line change
Expand Up @@ -24,27 +24,31 @@ import kafka.utils.TestUtils
import org.apache.kafka.clients.admin.AlterConfigOp.OpType
import org.apache.kafka.clients.admin._
import org.apache.kafka.common.acl.{AclBinding, AclBindingFilter}
import org.apache.kafka.common.config.ConfigException
import org.apache.kafka.common.config.ConfigResource
import org.apache.kafka.common.config.ConfigResource.Type
import org.apache.kafka.common.message.DescribeClusterRequestData
import org.apache.kafka.common.network.ListenerName
import org.apache.kafka.common.protocol.Errors._
import org.apache.kafka.common.quota.{ClientQuotaAlteration, ClientQuotaEntity, ClientQuotaFilter, ClientQuotaFilterComponent}
import org.apache.kafka.common.requests.{ApiError, DescribeClusterRequest, DescribeClusterResponse}
import org.apache.kafka.common.{Endpoint, TopicPartition, TopicPartitionInfo}
import org.apache.kafka.common.{Endpoint, Reconfigurable, TopicPartition, TopicPartitionInfo}
import org.apache.kafka.image.ClusterImage
import org.apache.kafka.metadata.BrokerState
import org.apache.kafka.server.authorizer._
import org.apache.kafka.server.common.MetadataVersion
import org.apache.kafka.server.log.remote.storage.RemoteLogManagerConfig
import org.junit.jupiter.api.Assertions._
import org.junit.jupiter.api.{Tag, Test, Timeout}
import org.junit.jupiter.params.ParameterizedTest
import org.junit.jupiter.params.provider.ValueSource
import org.slf4j.LoggerFactory

import java.io.File
import java.nio.file.{FileSystems, Path}
import java.util.concurrent.atomic.AtomicInteger
import java.{lang, util}
import java.util.concurrent.CompletionStage
import java.util.concurrent.{CompletableFuture, CompletionStage}
import java.util.{Arrays, Collections, Optional, OptionalLong, Properties}
import scala.annotation.nowarn
import scala.collection.mutable
Expand Down Expand Up @@ -966,6 +970,47 @@ class KRaftClusterTest {
cluster.close()
}
}

@ParameterizedTest
@ValueSource(booleans = Array(false, true))
def testReconfigureControllerAuthorizer(combinedMode: Boolean): Unit = {
val cluster = new KafkaClusterTestKit.Builder(
new TestKitNodes.Builder().
setNumBrokerNodes(1).
setCoResident(combinedMode).
setNumControllerNodes(1).build()).
setConfigProp("authorizer.class.name", classOf[FakeConfigurableAuthorizer].getName).
build()

def assertFoobarValue(expected: Int): Unit = {
TestUtils.retry(60000) {
assertEquals(expected, cluster.controllers().values().iterator().next().
authorizer.get.asInstanceOf[FakeConfigurableAuthorizer].foobar.get())
assertEquals(expected, cluster.brokers().values().iterator().next().
authorizer.get.asInstanceOf[FakeConfigurableAuthorizer].foobar.get())
}
}

try {
cluster.format()
cluster.startup()
cluster.waitForReadyBrokers()
assertFoobarValue(0)
val admin = Admin.create(cluster.clientProperties())
try {
admin.incrementalAlterConfigs(
Collections.singletonMap(new ConfigResource(Type.BROKER, ""),
Collections.singletonList(new AlterConfigOp(
new ConfigEntry(FakeConfigurableAuthorizer.foobarConfigKey, "123"), OpType.SET)))).
all().get()
} finally {
admin.close()
}
assertFoobarValue(123)
} finally {
cluster.close()
}
}
}

class BadAuthorizer() extends Authorizer {
Expand All @@ -986,3 +1031,71 @@ class BadAuthorizer() extends Authorizer {

override def deleteAcls(requestContext: AuthorizableRequestContext, aclBindingFilters: util.List[AclBindingFilter]): util.List[_ <: CompletionStage[AclDeleteResult]] = ???
}

object FakeConfigurableAuthorizer {
val foobarConfigKey = "fake.configurable.authorizer.foobar.config"

def fakeConfigurableAuthorizerConfigToInt(configs: util.Map[String, _]): Int = {
val result = configs.get(foobarConfigKey)
if (result == null) {
0
} else {
val resultString = result.toString().trim()
try {
Integer.valueOf(resultString)
} catch {
case e: NumberFormatException => throw new ConfigException(s"Bad value of ${foobarConfigKey}: ${resultString}")
}
}
}
}

class FakeConfigurableAuthorizer() extends Authorizer with Reconfigurable {
import FakeConfigurableAuthorizer._

val foobar = new AtomicInteger(0)

override def start(serverInfo: AuthorizerServerInfo): java.util.Map[Endpoint, _ <: CompletionStage[Void]] = {
serverInfo.endpoints().asScala.map(e => e -> {
val future = new CompletableFuture[Void]
future.complete(null)
future
}).toMap.asJava
}

override def reconfigurableConfigs(): java.util.Set[String] = Set(foobarConfigKey).asJava

override def validateReconfiguration(configs: util.Map[String, _]): Unit = {
fakeConfigurableAuthorizerConfigToInt(configs)
}

override def reconfigure(configs: util.Map[String, _]): Unit = {
foobar.set(fakeConfigurableAuthorizerConfigToInt(configs))
}

override def authorize(requestContext: AuthorizableRequestContext, actions: util.List[Action]): util.List[AuthorizationResult] = {
actions.asScala.map(_ => AuthorizationResult.ALLOWED).toList.asJava
}

override def acls(filter: AclBindingFilter): lang.Iterable[AclBinding] = List[AclBinding]().asJava

override def close(): Unit = {}

override def configure(configs: util.Map[String, _]): Unit = {
foobar.set(fakeConfigurableAuthorizerConfigToInt(configs))
}

override def createAcls(
requestContext: AuthorizableRequestContext,
aclBindings: util.List[AclBinding]
): util.List[_ <: CompletionStage[AclCreateResult]] = {
Collections.emptyList()
}

override def deleteAcls(
requestContext: AuthorizableRequestContext,
aclBindingFilters: util.List[AclBindingFilter]
): util.List[_ <: CompletionStage[AclDeleteResult]] = {
Collections.emptyList()
}
}
Loading