From c097f84cffbe41323ca0bbae5979620fcac0437f Mon Sep 17 00:00:00 2001 From: lenovo Date: Thu, 29 Nov 2018 19:27:34 +0800 Subject: [PATCH 1/2] Support JmxTool to connect to a secured RMI port. --- core/src/main/scala/kafka/tools/JmxTool.scala | 28 ++++++++++++++++++- 1 file changed, 27 insertions(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/tools/JmxTool.scala b/core/src/main/scala/kafka/tools/JmxTool.scala index c5303a9d96123..2414e9b2008fe 100644 --- a/core/src/main/scala/kafka/tools/JmxTool.scala +++ b/core/src/main/scala/kafka/tools/JmxTool.scala @@ -22,6 +22,7 @@ import java.util.Date import java.text.SimpleDateFormat import javax.management._ import javax.management.remote._ +import javax.rmi.ssl.SslRMIClientSocketFactory import joptsimple.OptionParser @@ -82,6 +83,16 @@ object JmxTool extends Logging { .describedAs("report-format") .ofType(classOf[java.lang.String]) .defaultsTo("original") + val jmxAuthPropOpt = parser.accepts("jmx-auth-prop", "A mechanism to pass property in the form 'username=password' " + + "when enabling remote JMX with password authentication.") + .withRequiredArg + .describedAs("jmx-auth-prop") + .ofType(classOf[String]) + val jmxSslEnableOpt = parser.accepts("jmx-ssl-enable", "Flag to enable remote JMX with SSL.") + .withRequiredArg + .describedAs("ssl-enable") + .ofType(classOf[java.lang.Boolean]) + .defaultsTo(false) val waitOpt = parser.accepts("wait", "Wait for requested JMX objects to become available before starting output. " + "Only supported when the list of objects is non-empty and contains no object name patterns.") val helpOpt = parser.accepts("help", "Print usage information.") @@ -109,6 +120,9 @@ object JmxTool extends Logging { val reportFormat = parseFormat(options.valueOf(reportFormatOpt).toLowerCase) val reportFormatOriginal = reportFormat.equals("original") + val enablePasswordAuth = options.has(jmxAuthPropOpt) + val enableSsl = options.has(jmxSslEnableOpt) + var jmxc: JMXConnector = null var mbsc: MBeanServerConnection = null var connected = false @@ -117,7 +131,19 @@ object JmxTool extends Logging { do { try { System.err.println(s"Trying to connect to JMX url: $url.") - jmxc = JMXConnectorFactory.connect(url, null) + val env = new java.util.HashMap[String, AnyRef] + // ssl enable + if (enableSsl) { + val csf = new SslRMIClientSocketFactory + env.put("com.sun.jndi.rmi.factory.socket", csf) + } + // password authentication enable + if (enablePasswordAuth) { + val userPassword = options.valueOf(jmxAuthPropOpt).split("=", 2) + val credentials = Array(userPassword(0), userPassword(1)) + env.put(JMXConnector.CREDENTIALS, credentials) + } + jmxc = JMXConnectorFactory.connect(url, env) mbsc = jmxc.getMBeanServerConnection connected = true } catch { From b9e0e50a60e7b9eef3af623af1ebdd0789733f19 Mon Sep 17 00:00:00 2001 From: lenovo Date: Tue, 15 Jan 2019 11:20:15 +0800 Subject: [PATCH 2/2] Remove a useless variable --- core/src/main/scala/kafka/tools/JmxTool.scala | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/tools/JmxTool.scala b/core/src/main/scala/kafka/tools/JmxTool.scala index 2414e9b2008fe..f51f1cfdad19a 100644 --- a/core/src/main/scala/kafka/tools/JmxTool.scala +++ b/core/src/main/scala/kafka/tools/JmxTool.scala @@ -139,8 +139,7 @@ object JmxTool extends Logging { } // password authentication enable if (enablePasswordAuth) { - val userPassword = options.valueOf(jmxAuthPropOpt).split("=", 2) - val credentials = Array(userPassword(0), userPassword(1)) + val credentials = options.valueOf(jmxAuthPropOpt).split("=", 2) env.put(JMXConnector.CREDENTIALS, credentials) } jmxc = JMXConnectorFactory.connect(url, env)