diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ConfigEntry.java b/clients/src/main/java/org/apache/kafka/clients/admin/ConfigEntry.java index e3426662e801d..30686c93eaeef 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ConfigEntry.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ConfigEntry.java @@ -174,11 +174,15 @@ public int hashCode() { return result; } + /** + * Override toString to redact sensitive value. + * WARNING, user should be responsible to set the correct "isSensitive" field for each config entry. + */ @Override public String toString() { return "ConfigEntry(" + "name=" + name + - ", value=" + value + + ", value=" + (isSensitive ? "Redacted" : value) + ", source=" + source + ", isSensitive=" + isSensitive + ", isReadOnly=" + isReadOnly + diff --git a/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala index 2cf24c87ceba9..9eefdd3933d3c 100755 --- a/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala +++ b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala @@ -295,7 +295,7 @@ class DynamicBrokerConfig(private val kafkaConfig: KafkaConfig) extends Logging dynamicBrokerConfigs ++= props.asScala updateCurrentConfig() } catch { - case e: Exception => error(s"Per-broker configs of $brokerId could not be applied: $persistentProps", e) + case e: Exception => error(s"Per-broker configs of $brokerId could not be applied: ${persistentProps.keys()}", e) } } @@ -306,7 +306,7 @@ class DynamicBrokerConfig(private val kafkaConfig: KafkaConfig) extends Logging dynamicDefaultConfigs ++= props.asScala updateCurrentConfig() } catch { - case e: Exception => error(s"Cluster default configs could not be applied: $persistentProps", e) + case e: Exception => error(s"Cluster default configs could not be applied: ${persistentProps.keys()}", e) } } @@ -469,7 +469,7 @@ class DynamicBrokerConfig(private val kafkaConfig: KafkaConfig) extends Logging } invalidProps.keys.foreach(props.remove) val configSource = if (perBrokerConfig) "broker" else "default cluster" - error(s"Dynamic $configSource config contains invalid values: $invalidProps, these configs will be ignored", e) + error(s"Dynamic $configSource config contains invalid values in: ${invalidProps.keys}, these configs will be ignored", e) } } @@ -555,7 +555,8 @@ class DynamicBrokerConfig(private val kafkaConfig: KafkaConfig) extends Logging } catch { case e: Exception => if (!validateOnly) - error(s"Failed to update broker configuration with configs : ${newConfig.originalsFromThisConfig}", e) + error(s"Failed to update broker configuration with configs : " + + s"${ConfigUtils.configMapToRedactedString(newConfig.originalsFromThisConfig, KafkaConfig.configDef)}", e) throw new ConfigException("Invalid dynamic configuration", e) } } diff --git a/core/src/main/scala/kafka/server/ZkAdminManager.scala b/core/src/main/scala/kafka/server/ZkAdminManager.scala index 87f522fb10d26..7be5ab08b374b 100644 --- a/core/src/main/scala/kafka/server/ZkAdminManager.scala +++ b/core/src/main/scala/kafka/server/ZkAdminManager.scala @@ -415,8 +415,12 @@ class ZkAdminManager(val config: KafkaConfig, info(message) resource -> ApiError.fromThrowable(new InvalidRequestException(message, e)) case e: Throwable => + val configProps = new Properties + config.entries.asScala.filter(_.value != null).foreach { configEntry => + configProps.setProperty(configEntry.name, configEntry.value) + } // Log client errors at a lower level than unexpected exceptions - val message = s"Error processing alter configs request for resource $resource, config $config" + val message = s"Error processing alter configs request for resource $resource, config ${toLoggableProps(resource, configProps).mkString(",")}" if (e.isInstanceOf[ApiException]) info(message, e) else