Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 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
15 changes: 10 additions & 5 deletions core/src/main/scala/kafka/server/DynamicBrokerConfig.scala
Original file line number Diff line number Diff line change
Expand Up @@ -516,7 +516,15 @@ class DynamicBrokerConfig(private val kafkaConfig: KafkaConfig) extends Logging
newProps ++= staticBrokerConfigs
overrideProps(newProps, dynamicDefaultConfigs)
overrideProps(newProps, dynamicBrokerConfigs)
val oldConfig = currentConfig

// We need a copy of the current config since `currentConfig` is initialized with `kafkaConfig`
// which means the call to `updateCurrentConfig` would end up mutating `oldConfig`.
val oldConfig = if (kafkaConfig eq currentConfig) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@hachikuji Thanks for the PR. Do you know if there is an issue only with KRraft or is there an issue with ZK as well?

With ZK, the sequence is:

  1. Create KafkaConfig with static configs from server.properties
  2. Initialize KafkaConfig with initial configs from ZooKeeper using DynamicBrokerConfig.initialize(zkClient). This always creates a new KafkaConfig, so we don't need this check?
  3. Start DynamicConfigManager. All ZK updates after 2) are handled through change notifications.

With KRaft, there doesn't seem to be an initialize() for initializing the state from existing dynamic configs, so are all configs handled similar to 3)? In which case, we should perhaps move the initialization of currentConfig from DynamicBrokerConfig.initialize(zkClient) to somewhere common for KRaft and avoid this check here?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@rajinisivaram Hmm, good point. It looks like the bug only affects KRaft since there is no call to initialize. The code still feels a little slippery though, so maybe there is room for some defensiveness. Let me take a look.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I pushed a patch which initializes currentConfig as null in order to make the call to initialize required. Let me know if that seems like a reasonable approach.

KafkaConfig(currentConfig.values, doLog = false)
} else {
currentConfig
}

val (newConfig, brokerReconfigurablesToUpdate) = processReconfiguration(newProps, validateOnly = false)
if (newConfig ne currentConfig) {
currentConfig = newConfig
Expand Down Expand Up @@ -719,7 +727,7 @@ class DynamicThreadPool(server: KafkaBroker) extends BrokerReconfigurable {
if (newConfig.numNetworkThreads != oldConfig.numNetworkThreads)
server.socketServer.resizeThreadPool(oldConfig.numNetworkThreads, newConfig.numNetworkThreads)
if (newConfig.numReplicaFetchers != oldConfig.numReplicaFetchers)
server.replicaManager.replicaFetcherManager.resizeThreadPool(newConfig.numReplicaFetchers)
server.replicaManager.resizeFetcherThreadPool(newConfig.numReplicaFetchers)
if (newConfig.numRecoveryThreadsPerDataDir != oldConfig.numRecoveryThreadsPerDataDir)
server.logManager.resizeRecoveryThreadPool(newConfig.numRecoveryThreadsPerDataDir)
if (newConfig.backgroundThreads != oldConfig.backgroundThreads)
Expand Down Expand Up @@ -955,6 +963,3 @@ class DynamicListenerConfig(server: KafkaBroker) extends BrokerReconfigurable wi
listeners.map(e => (e.listenerName, e)).toMap

}



2 changes: 1 addition & 1 deletion core/src/main/scala/kafka/server/KafkaConfig.scala
Original file line number Diff line number Diff line change
Expand Up @@ -1341,7 +1341,7 @@ object KafkaConfig {
fromProps(props, doLog)
}

def apply(props: java.util.Map[_, _]): KafkaConfig = new KafkaConfig(props, true)
def apply(props: java.util.Map[_, _], doLog: Boolean = true): KafkaConfig = new KafkaConfig(props, doLog)

private def typeOf(name: String): Option[ConfigDef.Type] = Option(configDef.configKeys.get(name)).map(_.`type`)

Expand Down
12 changes: 8 additions & 4 deletions core/src/main/scala/kafka/server/KafkaServer.scala
Original file line number Diff line number Diff line change
Expand Up @@ -114,7 +114,9 @@ class KafkaServer(
var controlPlaneRequestHandlerPool: KafkaRequestHandlerPool = null

var logDirFailureChannel: LogDirFailureChannel = null
var logManager: LogManager = null
var _logManager: LogManager = null
Comment thread
hachikuji marked this conversation as resolved.
Outdated

def logManager: LogManager = _logManager

@volatile private[this] var _replicaManager: ReplicaManager = null
var adminManager: ZkAdminManager = null
Expand All @@ -129,7 +131,9 @@ class KafkaServer(

var transactionCoordinator: TransactionCoordinator = null

var kafkaController: KafkaController = null
private var _kafkaController: KafkaController = null

def kafkaController: KafkaController = _kafkaController

var forwardingManager: Option[ForwardingManager] = None

Expand Down Expand Up @@ -250,7 +254,7 @@ class KafkaServer(
logDirFailureChannel = new LogDirFailureChannel(config.logDirs.size)

/* start log manager */
logManager = LogManager(config, initialOfflineDirs,
_logManager = LogManager(config, initialOfflineDirs,
new ZkConfigRepository(new AdminZkClient(zkClient)),
kafkaScheduler, time, brokerTopicStats, logDirFailureChannel, config.usesTopicId)
_brokerState = BrokerState.RECOVERY
Expand Down Expand Up @@ -327,7 +331,7 @@ class KafkaServer(
tokenManager.startup()

/* start kafka controller */
kafkaController = new KafkaController(config, zkClient, time, metrics, brokerInfo, brokerEpoch, tokenManager, brokerFeatures, featureCache, threadNamePrefix)
_kafkaController = new KafkaController(config, zkClient, time, metrics, brokerInfo, brokerEpoch, tokenManager, brokerFeatures, featureCache, threadNamePrefix)
kafkaController.startup()

adminManager = new ZkAdminManager(config, metrics, metadataCache, zkClient)
Expand Down
4 changes: 4 additions & 0 deletions core/src/main/scala/kafka/server/ReplicaManager.scala
Original file line number Diff line number Diff line change
Expand Up @@ -280,6 +280,10 @@ class ReplicaManager(val config: KafkaConfig,
replicaAlterLogDirsManager.shutdownIdleFetcherThreads()
}

def resizeFetcherThreadPool(newSize: Int): Unit = {
replicaFetcherManager.resizeThreadPool(newSize)
}

def getLog(topicPartition: TopicPartition): Option[UnifiedLog] = logManager.getLog(topicPartition)

def hasDelayedElectionOperations: Boolean = delayedElectLeaderPurgatory.numDelayed != 0
Expand Down
100 changes: 99 additions & 1 deletion core/src/test/scala/unit/kafka/server/DynamicBrokerConfigTest.scala
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,12 @@ package kafka.server
import java.{lang, util}
import java.util.Properties
import java.util.concurrent.CompletionStage
import kafka.utils.TestUtils
import java.util.concurrent.atomic.AtomicReference

import kafka.controller.KafkaController
import kafka.log.{LogConfig, LogManager}
import kafka.network.SocketServer
import kafka.utils.{KafkaScheduler, TestUtils}
import kafka.zk.KafkaZkClient
import org.apache.kafka.common.{Endpoint, Reconfigurable}
import org.apache.kafka.common.acl.{AclBinding, AclBindingFilter}
Expand All @@ -30,6 +35,7 @@ import org.apache.kafka.server.authorizer._
import org.easymock.EasyMock
import org.junit.jupiter.api.Assertions._
import org.junit.jupiter.api.Test
import org.mockito.{ArgumentMatchers, Mockito}

import scala.annotation.nowarn
import scala.jdk.CollectionConverters._
Expand Down Expand Up @@ -80,6 +86,98 @@ class DynamicBrokerConfigTest {
}
}

@Test
def testEnableDefaultUncleanLeaderElection(): Unit = {
val origProps = TestUtils.createBrokerConfig(0, TestUtils.MockZkConnect, port = 8181)
origProps.put(KafkaConfig.UncleanLeaderElectionEnableProp, "false")

val config = KafkaConfig(origProps)
val serverMock = Mockito.mock(classOf[KafkaServer])
val controllerMock = Mockito.mock(classOf[KafkaController])
val logManagerMock = Mockito.mock(classOf[LogManager])

Mockito.when(serverMock.config).thenReturn(config)
Mockito.when(serverMock.kafkaController).thenReturn(controllerMock)
Mockito.when(serverMock.logManager).thenReturn(logManagerMock)
Mockito.when(logManagerMock.allLogs).thenReturn(Iterable.empty)

val currentDefaultLogConfig = new AtomicReference(LogConfig())
Mockito.when(logManagerMock.currentDefaultConfig).thenAnswer(_ => currentDefaultLogConfig.get())
Mockito.when(logManagerMock.reconfigureDefaultLogConfig(ArgumentMatchers.any(classOf[LogConfig])))
.thenAnswer(invocation => currentDefaultLogConfig.set(invocation.getArgument(0)))

config.dynamicConfig.addBrokerReconfigurable(new DynamicLogConfig(logManagerMock, serverMock))

val props = new Properties()

props.put(KafkaConfig.UncleanLeaderElectionEnableProp, "true")
config.dynamicConfig.updateDefaultConfig(props)
assertTrue(config.uncleanLeaderElectionEnable)
Mockito.verify(controllerMock).enableDefaultUncleanLeaderElection()
}

@Test
def testUpdateDynamicThreadPool(): Unit = {
val origProps = TestUtils.createBrokerConfig(0, TestUtils.MockZkConnect, port = 8181)
origProps.put(KafkaConfig.NumIoThreadsProp, "4")
origProps.put(KafkaConfig.NumNetworkThreadsProp, "2")
origProps.put(KafkaConfig.NumReplicaFetchersProp, "1")
origProps.put(KafkaConfig.NumRecoveryThreadsPerDataDirProp, "1")
origProps.put(KafkaConfig.BackgroundThreadsProp, "3")

val config = KafkaConfig(origProps)
val serverMock = Mockito.mock(classOf[KafkaBroker])
val handlerPoolMock = Mockito.mock(classOf[KafkaRequestHandlerPool])
val socketServerMock = Mockito.mock(classOf[SocketServer])
val replicaManagerMock = Mockito.mock(classOf[ReplicaManager])
val logManagerMock = Mockito.mock(classOf[LogManager])
val schedulerMock = Mockito.mock(classOf[KafkaScheduler])

Mockito.when(serverMock.config).thenReturn(config)
Mockito.when(serverMock.dataPlaneRequestHandlerPool).thenReturn(handlerPoolMock)
Mockito.when(serverMock.socketServer).thenReturn(socketServerMock)
Mockito.when(serverMock.replicaManager).thenReturn(replicaManagerMock)
Mockito.when(serverMock.logManager).thenReturn(logManagerMock)
Mockito.when(serverMock.kafkaScheduler).thenReturn(schedulerMock)

config.dynamicConfig.addBrokerReconfigurable(new DynamicThreadPool(serverMock))

val props = new Properties()

props.put(KafkaConfig.NumIoThreadsProp, "8")
config.dynamicConfig.updateDefaultConfig(props)
assertEquals(8, config.numIoThreads)
Mockito.verify(handlerPoolMock).resizeThreadPool(newSize = 8)

props.put(KafkaConfig.NumNetworkThreadsProp, "4")
config.dynamicConfig.updateDefaultConfig(props)
assertEquals(4, config.numNetworkThreads)
Mockito.verify(socketServerMock).resizeThreadPool(oldNumNetworkThreads = 2, newNumNetworkThreads = 4)

props.put(KafkaConfig.NumReplicaFetchersProp, "2")
config.dynamicConfig.updateDefaultConfig(props)
assertEquals(2, config.numReplicaFetchers)
Mockito.verify(replicaManagerMock).resizeFetcherThreadPool(newSize = 2)

props.put(KafkaConfig.NumRecoveryThreadsPerDataDirProp, "2")
config.dynamicConfig.updateDefaultConfig(props)
assertEquals(2, config.numRecoveryThreadsPerDataDir)
Mockito.verify(logManagerMock).resizeRecoveryThreadPool(newSize = 2)

props.put(KafkaConfig.BackgroundThreadsProp, "6")
config.dynamicConfig.updateDefaultConfig(props)
assertEquals(6, config.backgroundThreads)
Mockito.verify(schedulerMock).resizeThreadPool(newSize = 6)

Mockito.verifyNoMoreInteractions(
handlerPoolMock,
socketServerMock,
replicaManagerMock,
logManagerMock,
schedulerMock
)
}

@nowarn("cat=deprecation")
@Test
def testConfigUpdateWithSomeInvalidConfigs(): Unit = {
Expand Down