From 1b5d07c73309ea36216e434e86d28cf9cb64bcee Mon Sep 17 00:00:00 2001 From: Manikumar Reddy Date: Sat, 21 Jul 2018 17:41:41 +0530 Subject: [PATCH 1/8] MINOR: close Zookeeper instance incase of connect timeout in ZooKeeperClient --- core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala | 1 + 1 file changed, 1 insertion(+) diff --git a/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala b/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala index 5cb127c3c7edf..91faa7dc9bd28 100755 --- a/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala +++ b/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala @@ -222,6 +222,7 @@ class ZooKeeperClient(connectString: String, var state = connectionState while (!state.isConnected && state.isAlive) { if (nanos <= 0) { + zooKeeper.close() throw new ZooKeeperClientTimeoutException(s"Timed out waiting for connection while in state: $state") } nanos = isConnectedOrExpiredCondition.awaitNanos(nanos) From cabce51e19c98da3bca00c123faf127b1ab1fbd8 Mon Sep 17 00:00:00 2001 From: Manikumar Reddy Date: Sat, 21 Jul 2018 22:45:57 +0530 Subject: [PATCH 2/8] Address review comment --- .../src/main/scala/kafka/zookeeper/ZooKeeperClient.scala | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala b/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala index 91faa7dc9bd28..bcaa06a6f98fa 100755 --- a/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala +++ b/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala @@ -92,7 +92,13 @@ class ZooKeeperClient(connectString: String, metricNames += "SessionState" expiryScheduler.startup() - waitUntilConnected(connectionTimeoutMs, TimeUnit.MILLISECONDS) + try { + waitUntilConnected(connectionTimeoutMs, TimeUnit.MILLISECONDS) + } catch { + case e: ZooKeeperClientTimeoutException => + zooKeeper.close() + throw e + } override def metricName(name: String, metricTags: scala.collection.Map[String, String]): MetricName = { explicitMetricName(metricGroup, metricType, name, metricTags) @@ -222,7 +228,6 @@ class ZooKeeperClient(connectString: String, var state = connectionState while (!state.isConnected && state.isAlive) { if (nanos <= 0) { - zooKeeper.close() throw new ZooKeeperClientTimeoutException(s"Timed out waiting for connection while in state: $state") } nanos = isConnectedOrExpiredCondition.awaitNanos(nanos) From ec05c3e7109891f221a799ff63a66365b8dd3741 Mon Sep 17 00:00:00 2001 From: Manikumar Reddy Date: Sun, 22 Jul 2018 01:47:46 +0530 Subject: [PATCH 3/8] Call ZooKeeperClient close --- core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala b/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala index bcaa06a6f98fa..197b4a1b823ee 100755 --- a/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala +++ b/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala @@ -96,7 +96,7 @@ class ZooKeeperClient(connectString: String, waitUntilConnected(connectionTimeoutMs, TimeUnit.MILLISECONDS) } catch { case e: ZooKeeperClientTimeoutException => - zooKeeper.close() + close() throw e } From 227f889fe0f588dc24d53a7804e6f964a27cf615 Mon Sep 17 00:00:00 2001 From: Manikumar Reddy Date: Sun, 22 Jul 2018 02:37:57 +0530 Subject: [PATCH 4/8] use ZooKeeperClientException --- core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala b/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala index 197b4a1b823ee..566c4b2708b24 100755 --- a/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala +++ b/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala @@ -95,7 +95,7 @@ class ZooKeeperClient(connectString: String, try { waitUntilConnected(connectionTimeoutMs, TimeUnit.MILLISECONDS) } catch { - case e: ZooKeeperClientTimeoutException => + case e: ZooKeeperClientException => close() throw e } From 40ddf9b84d1b0c671fbd9c53a3c0a0d2b78ff84d Mon Sep 17 00:00:00 2001 From: Manikumar Reddy Date: Sun, 22 Jul 2018 10:56:30 +0530 Subject: [PATCH 5/8] Catch Throwable --- core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala b/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala index 566c4b2708b24..74f9a2c1a61ef 100755 --- a/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala +++ b/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala @@ -95,7 +95,7 @@ class ZooKeeperClient(connectString: String, try { waitUntilConnected(connectionTimeoutMs, TimeUnit.MILLISECONDS) } catch { - case e: ZooKeeperClientException => + case e: Throwable => close() throw e } From 1a01f88baabd612aed60d86ae3a6096b580f8075 Mon Sep 17 00:00:00 2001 From: Manikumar Reddy Date: Sun, 22 Jul 2018 02:29:12 +0530 Subject: [PATCH 6/8] Add zookeeper SendThread count based test --- .../kafka/zookeeper/ZooKeeperClientTest.scala | 19 ++++++++++++++++--- 1 file changed, 16 insertions(+), 3 deletions(-) diff --git a/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala b/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala index fcbf699ec9614..a5ebd27a10d18 100644 --- a/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala +++ b/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala @@ -23,6 +23,7 @@ import java.util.concurrent.{ArrayBlockingQueue, ConcurrentLinkedQueue, CountDow import com.yammer.metrics.Metrics import com.yammer.metrics.core.{Gauge, Meter, MetricName} +import kafka.utils.TestUtils import kafka.zk.ZooKeeperTestHarness import org.apache.kafka.common.security.JaasUtils import org.apache.kafka.common.utils.Time @@ -57,12 +58,24 @@ class ZooKeeperClientTest extends ZooKeeperTestHarness { System.clearProperty(JaasUtils.JAVA_LOGIN_CONFIG_PARAM) } - @Test(expected = classOf[ZooKeeperClientTimeoutException]) + @Test def testUnresolvableConnectString(): Unit = { - new ZooKeeperClient("some.invalid.hostname.foo.bar.local", zkSessionTimeout, connectionTimeoutMs = 10, - Int.MaxValue, time, "testMetricGroup", "testMetricType").close() + val hostAddress = "some.invalid.hostname.foo.bar.local" + try { + new ZooKeeperClient(hostAddress, zkSessionTimeout, connectionTimeoutMs = 10, + Int.MaxValue, time, "testMetricGroup", "testMetricType") + } catch { + case e: ZooKeeperClientTimeoutException => TestUtils.waitUntilTrue(() => unexpectedZKThreadCount(hostAddress) == 0, + "ZooKeeper client threads are not closed") + } } + def unexpectedZKThreadCount(hostAddress: String): Int = Thread.getAllStackTraces.keySet.toArray + .map(_.asInstanceOf[Thread]) + .filter(_.isAlive) + .map(_.getName) + .count(t => t.contains("SendThread()") || t.contains(hostAddress)) // Verify threadName pattern for unresolvable host zk send threads + @Test(expected = classOf[ZooKeeperClientTimeoutException]) def testConnectionTimeout(): Unit = { zookeeper.shutdown() From a19ff8045dd15e5b9df9f7d445942d3014a27bc6 Mon Sep 17 00:00:00 2001 From: Manikumar Reddy Date: Sun, 22 Jul 2018 15:48:12 +0530 Subject: [PATCH 7/8] Address review comments --- .../scala/kafka/zookeeper/ZooKeeperClient.scala | 5 ++--- .../unit/kafka/zookeeper/ZooKeeperClientTest.scala | 13 +++++-------- 2 files changed, 7 insertions(+), 11 deletions(-) diff --git a/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala b/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala index 74f9a2c1a61ef..97ec9a44c36f3 100755 --- a/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala +++ b/core/src/main/scala/kafka/zookeeper/ZooKeeperClient.scala @@ -92,9 +92,8 @@ class ZooKeeperClient(connectString: String, metricNames += "SessionState" expiryScheduler.startup() - try { - waitUntilConnected(connectionTimeoutMs, TimeUnit.MILLISECONDS) - } catch { + try waitUntilConnected(connectionTimeoutMs, TimeUnit.MILLISECONDS) + catch { case e: Throwable => close() throw e diff --git a/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala b/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala index a5ebd27a10d18..affaf39625742 100644 --- a/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala +++ b/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala @@ -23,7 +23,6 @@ import java.util.concurrent.{ArrayBlockingQueue, ConcurrentLinkedQueue, CountDow import com.yammer.metrics.Metrics import com.yammer.metrics.core.{Gauge, Meter, MetricName} -import kafka.utils.TestUtils import kafka.zk.ZooKeeperTestHarness import org.apache.kafka.common.security.JaasUtils import org.apache.kafka.common.utils.Time @@ -60,21 +59,19 @@ class ZooKeeperClientTest extends ZooKeeperTestHarness { @Test def testUnresolvableConnectString(): Unit = { - val hostAddress = "some.invalid.hostname.foo.bar.local" try { - new ZooKeeperClient(hostAddress, zkSessionTimeout, connectionTimeoutMs = 10, + new ZooKeeperClient("some.invalid.hostname.foo.bar.local", zkSessionTimeout, connectionTimeoutMs = 10, Int.MaxValue, time, "testMetricGroup", "testMetricType") } catch { - case e: ZooKeeperClientTimeoutException => TestUtils.waitUntilTrue(() => unexpectedZKThreadCount(hostAddress) == 0, - "ZooKeeper client threads are not closed") + case e: ZooKeeperClientTimeoutException => + assertEquals("ZooKeeper client threads still running", Set.empty, runningZkSendThreads()) } } - def unexpectedZKThreadCount(hostAddress: String): Int = Thread.getAllStackTraces.keySet.toArray - .map(_.asInstanceOf[Thread]) + private def runningZkSendThreads(): collection.Set[String] = Thread.getAllStackTraces.keySet.asScala .filter(_.isAlive) .map(_.getName) - .count(t => t.contains("SendThread()") || t.contains(hostAddress)) // Verify threadName pattern for unresolvable host zk send threads + .filter(t => t.contains("SendThread()")) //Verify threadName pattern for unresolvable host zk send threads @Test(expected = classOf[ZooKeeperClientTimeoutException]) def testConnectionTimeout(): Unit = { From b6735f1d8b05d6e44a9c3d6ed5c85bcf89bf023e Mon Sep 17 00:00:00 2001 From: Ismael Juma Date: Sun, 22 Jul 2018 09:31:20 -0700 Subject: [PATCH 8/8] Minor tweaks --- .../scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala b/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala index affaf39625742..0088c657de0a4 100644 --- a/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala +++ b/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala @@ -64,14 +64,14 @@ class ZooKeeperClientTest extends ZooKeeperTestHarness { Int.MaxValue, time, "testMetricGroup", "testMetricType") } catch { case e: ZooKeeperClientTimeoutException => - assertEquals("ZooKeeper client threads still running", Set.empty, runningZkSendThreads()) + assertEquals("ZooKeeper client threads still running", Set.empty, runningZkSendThreads) } } - private def runningZkSendThreads(): collection.Set[String] = Thread.getAllStackTraces.keySet.asScala + private def runningZkSendThreads: collection.Set[String] = Thread.getAllStackTraces.keySet.asScala .filter(_.isAlive) .map(_.getName) - .filter(t => t.contains("SendThread()")) //Verify threadName pattern for unresolvable host zk send threads + .filter(t => t.contains("SendThread()")) @Test(expected = classOf[ZooKeeperClientTimeoutException]) def testConnectionTimeout(): Unit = {