From cb3fc90943440ac61ad4694fd0039f995445b686 Mon Sep 17 00:00:00 2001 From: Manikumar Reddy Date: Tue, 15 Dec 2020 07:17:57 +0530 Subject: [PATCH 1/3] MINOR: Update log statements in alterBrokerConfigs/alterTopicConfigs methods --- .../scala/kafka/network/RequestChannel.scala | 16 ++-------------- .../main/scala/kafka/server/AdminManager.scala | 14 +++++++++++++- .../main/scala/kafka/server/KafkaConfig.scala | 14 ++++++++++++-- 3 files changed, 27 insertions(+), 17 deletions(-) diff --git a/core/src/main/scala/kafka/network/RequestChannel.scala b/core/src/main/scala/kafka/network/RequestChannel.scala index 0c03d70188165..70ade3ec167ac 100644 --- a/core/src/main/scala/kafka/network/RequestChannel.scala +++ b/core/src/main/scala/kafka/network/RequestChannel.scala @@ -24,12 +24,10 @@ import java.util.concurrent._ import com.fasterxml.jackson.databind.JsonNode import com.typesafe.scalalogging.Logger import com.yammer.metrics.core.Meter -import kafka.log.LogConfig import kafka.metrics.KafkaMetricsGroup import kafka.server.KafkaConfig import kafka.utils.{Logging, NotNothing, Pool} import kafka.utils.Implicits._ -import org.apache.kafka.common.config.types.Password import org.apache.kafka.common.config.ConfigResource import org.apache.kafka.common.memory.MemoryPool import org.apache.kafka.common.message.IncrementalAlterConfigsRequestData @@ -164,21 +162,11 @@ object RequestChannel extends Logging { def loggableRequest: AbstractRequest = { - def loggableValue(resourceType: ConfigResource.Type, name: String, value: String): String = { - val maybeSensitive = resourceType match { - case ConfigResource.Type.BROKER => KafkaConfig.maybeSensitive(KafkaConfig.configType(name)) - case ConfigResource.Type.TOPIC => KafkaConfig.maybeSensitive(LogConfig.configType(name)) - case ConfigResource.Type.BROKER_LOGGER => false - case _ => true - } - if (maybeSensitive) Password.HIDDEN else value - } - bodyAndSize.request match { case alterConfigs: AlterConfigsRequest => val loggableConfigs = alterConfigs.configs().asScala.map { case (resource, config) => val loggableEntries = new AlterConfigsRequest.Config(config.entries.asScala.map { entry => - new AlterConfigsRequest.ConfigEntry(entry.name, loggableValue(resource.`type`, entry.name, entry.value)) + new AlterConfigsRequest.ConfigEntry(entry.name, KafkaConfig.loggableValue(resource.`type`, entry.name, entry.value)) }.asJavaCollection) (resource, loggableEntries) }.asJava @@ -193,7 +181,7 @@ object RequestChannel extends Logging { resource.configs.forEach { config => newResource.configs.add(new AlterableConfig() .setName(config.name) - .setValue(loggableValue(ConfigResource.Type.forId(resource.resourceType), config.name, config.value)) + .setValue(KafkaConfig.loggableValue(ConfigResource.Type.forId(resource.resourceType), config.name, config.value)) .setConfigOperation(config.configOperation)) } resources.add(newResource) diff --git a/core/src/main/scala/kafka/server/AdminManager.scala b/core/src/main/scala/kafka/server/AdminManager.scala index f1bd1e276c4f1..f82bdee29835f 100644 --- a/core/src/main/scala/kafka/server/AdminManager.scala +++ b/core/src/main/scala/kafka/server/AdminManager.scala @@ -506,7 +506,7 @@ class AdminManager(val config: KafkaConfig, adminZkClient.validateTopicConfig(topic, configProps) validateConfigPolicy(resource, configEntriesMap) if (!validateOnly) { - info(s"Updating topic $topic with new configuration $config") + info(s"Updating topic $topic with new configuration : ${toLoggableProps(resource, configProps).mkString(",")}") adminZkClient.changeTopicConfig(topic, configProps) } @@ -522,6 +522,12 @@ class AdminManager(val config: KafkaConfig, if (!validateOnly) { if (perBrokerConfig) this.config.dynamicConfig.reloadUpdatedFilesWithoutConfigChange(configProps) + + if (perBrokerConfig) + info(s"Updating broker ${brokerId.get} with new configuration : ${toLoggableProps(resource, configProps).mkString(",")}") + else + info(s"Updating brokers with new configuration : ${toLoggableProps(resource, configProps).mkString(",")}") + adminZkClient.changeBrokerConfig(brokerId, this.config.dynamicConfig.toPersistentProps(configProps, perBrokerConfig)) } @@ -529,6 +535,12 @@ class AdminManager(val config: KafkaConfig, resource -> ApiError.NONE } + private def toLoggableProps(resource: ConfigResource, configProps: Properties): Map[String, String] = { + configProps.asScala.map { + case (key, value) => (key, KafkaConfig.loggableValue(resource.`type`, key, value)) + } + } + private def alterLogLevelConfigs(alterConfigOps: Seq[AlterConfigOp]): Unit = { alterConfigOps.foreach { alterConfigOp => val loggerName = alterConfigOp.configEntry().name() diff --git a/core/src/main/scala/kafka/server/KafkaConfig.scala b/core/src/main/scala/kafka/server/KafkaConfig.scala index 93b840576e09a..895612e2f7d06 100755 --- a/core/src/main/scala/kafka/server/KafkaConfig.scala +++ b/core/src/main/scala/kafka/server/KafkaConfig.scala @@ -24,17 +24,17 @@ import kafka.api.{ApiVersion, ApiVersionValidator, KAFKA_0_10_0_IV1, KAFKA_2_1_I import kafka.cluster.EndPoint import kafka.coordinator.group.OffsetConfig import kafka.coordinator.transaction.{TransactionLog, TransactionStateManager} +import kafka.log.LogConfig import kafka.message.{BrokerCompressionCodec, CompressionCodec, ZStdCompressionCodec} import kafka.security.authorizer.AuthorizerUtils import kafka.utils.CoreUtils import kafka.utils.Implicits._ import org.apache.kafka.clients.CommonClientConfigs import org.apache.kafka.common.Reconfigurable -import org.apache.kafka.common.config.SecurityConfig +import org.apache.kafka.common.config.{AbstractConfig, ConfigDef, ConfigException, ConfigResource, SaslConfigs, SecurityConfig, SslClientAuth, SslConfigs, TopicConfig} import org.apache.kafka.common.config.ConfigDef.{ConfigKey, ValidList} import org.apache.kafka.common.config.internals.BrokerSecurityConfigs import org.apache.kafka.common.config.types.Password -import org.apache.kafka.common.config.{AbstractConfig, ConfigDef, ConfigException, SaslConfigs, SslClientAuth, SslConfigs, TopicConfig} import org.apache.kafka.common.metrics.Sensor import org.apache.kafka.common.network.ListenerName import org.apache.kafka.common.record.{LegacyRecord, Records, TimestampType} @@ -1311,6 +1311,16 @@ object KafkaConfig { // If we can't determine the config entry type, treat it as a sensitive config to be safe configType.isEmpty || configType.contains(ConfigDef.Type.PASSWORD) } + + def loggableValue(resourceType: ConfigResource.Type, name: String, value: String): String = { + val maybeSensitive = resourceType match { + case ConfigResource.Type.BROKER => KafkaConfig.maybeSensitive(KafkaConfig.configType(name)) + case ConfigResource.Type.TOPIC => KafkaConfig.maybeSensitive(LogConfig.configType(name)) + case ConfigResource.Type.BROKER_LOGGER => false + case _ => true + } + if (maybeSensitive) Password.HIDDEN else value + } } class KafkaConfig(val props: java.util.Map[_, _], doLog: Boolean, dynamicConfigOverride: Option[DynamicBrokerConfig]) From f2df054674da77649fd8f376032d6969143fdb74 Mon Sep 17 00:00:00 2001 From: Manikumar Reddy Date: Wed, 6 Jan 2021 13:03:38 +0530 Subject: [PATCH 2/3] Address review comments --- core/src/main/scala/kafka/server/AdminManager.scala | 10 ++++++++-- 1 file changed, 8 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/server/AdminManager.scala b/core/src/main/scala/kafka/server/AdminManager.scala index f82bdee29835f..9da262f5882a7 100644 --- a/core/src/main/scala/kafka/server/AdminManager.scala +++ b/core/src/main/scala/kafka/server/AdminManager.scala @@ -546,8 +546,14 @@ class AdminManager(val config: KafkaConfig, val loggerName = alterConfigOp.configEntry().name() val logLevel = alterConfigOp.configEntry().value() alterConfigOp.opType() match { - case OpType.SET => Log4jController.logLevel(loggerName, logLevel) - case OpType.DELETE => Log4jController.unsetLogLevel(loggerName) + case OpType.SET => { + info(s"Updating the log level of $loggerName to $logLevel") + Log4jController.logLevel(loggerName, logLevel) + } + case OpType.DELETE => { + info(s"Unset the log level of $loggerName") + Log4jController.unsetLogLevel(loggerName) + } case _ => throw new IllegalArgumentException( s"Log level cannot be changed for OpType: ${alterConfigOp.opType()}") } From d7957b1093ccba53882f4a7cafb6a3d1a435a5a0 Mon Sep 17 00:00:00 2001 From: Manikumar Reddy Date: Wed, 6 Jan 2021 16:04:07 +0530 Subject: [PATCH 3/3] remove redundant braces --- core/src/main/scala/kafka/server/AdminManager.scala | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/core/src/main/scala/kafka/server/AdminManager.scala b/core/src/main/scala/kafka/server/AdminManager.scala index 9da262f5882a7..b22ccd7ccdc70 100644 --- a/core/src/main/scala/kafka/server/AdminManager.scala +++ b/core/src/main/scala/kafka/server/AdminManager.scala @@ -546,14 +546,12 @@ class AdminManager(val config: KafkaConfig, val loggerName = alterConfigOp.configEntry().name() val logLevel = alterConfigOp.configEntry().value() alterConfigOp.opType() match { - case OpType.SET => { + case OpType.SET => info(s"Updating the log level of $loggerName to $logLevel") Log4jController.logLevel(loggerName, logLevel) - } - case OpType.DELETE => { + case OpType.DELETE => info(s"Unset the log level of $loggerName") Log4jController.unsetLogLevel(loggerName) - } case _ => throw new IllegalArgumentException( s"Log level cannot be changed for OpType: ${alterConfigOp.opType()}") }