diff --git a/core/src/main/scala/kafka/server/AlterIsrManager.scala b/core/src/main/scala/kafka/server/AlterIsrManager.scala index 9ad734f708c9b..e6a934df61423 100644 --- a/core/src/main/scala/kafka/server/AlterIsrManager.scala +++ b/core/src/main/scala/kafka/server/AlterIsrManager.scala @@ -80,7 +80,7 @@ object AlterIsrManager { time = time, metrics = metrics, config = config, - channelName = "alterIsrChannel", + channelName = "alterIsr", threadNamePrefix = threadNamePrefix, retryTimeoutMs = Long.MaxValue ) diff --git a/core/src/main/scala/kafka/server/BrokerServer.scala b/core/src/main/scala/kafka/server/BrokerServer.scala index 8689afd743d88..aecb71ea4b045 100644 --- a/core/src/main/scala/kafka/server/BrokerServer.scala +++ b/core/src/main/scala/kafka/server/BrokerServer.scala @@ -186,10 +186,11 @@ class BrokerServer( time, metrics, config, - channelName = "controllerForwardingChannel", + channelName = "forwarding", threadNamePrefix, retryTimeoutMs = 60000 ) + clientToControllerChannelManager.start() forwardingManager = new ForwardingManagerImpl(clientToControllerChannelManager) val apiVersionManager = ApiVersionManager( @@ -211,7 +212,7 @@ class BrokerServer( time, metrics, config, - channelName = "alterisr", + channelName = "alterIsr", threadNamePrefix, retryTimeoutMs = Long.MaxValue ) diff --git a/core/src/main/scala/kafka/server/BrokerToControllerChannelManager.scala b/core/src/main/scala/kafka/server/BrokerToControllerChannelManager.scala index 16c4a5b95d5f8..66c12dadd9cc9 100644 --- a/core/src/main/scala/kafka/server/BrokerToControllerChannelManager.scala +++ b/core/src/main/scala/kafka/server/BrokerToControllerChannelManager.scala @@ -168,7 +168,7 @@ class BrokerToControllerChannelManagerImpl( threadNamePrefix: Option[String], retryTimeoutMs: Long ) extends BrokerToControllerChannelManager with Logging { - private val logContext = new LogContext(s"[broker-${config.brokerId}-to-controller] ") + private val logContext = new LogContext(s"[BrokerToControllerChannelManager broker=${config.brokerId} name=$channelName] ") private val manualMetadataUpdater = new ManualMetadataUpdater() private val apiVersions = new ApiVersions() private val currentNodeApiVersions = NodeApiVersions.create() @@ -226,8 +226,8 @@ class BrokerToControllerChannelManagerImpl( ) } val threadName = threadNamePrefix match { - case None => s"broker-${config.brokerId}-to-controller-send-thread" - case Some(name) => s"$name:broker-${config.brokerId}-to-controller-send-thread" + case None => s"BrokerToControllerChannelManager broker=${config.brokerId} name=$channelName" + case Some(name) => s"$name:BrokerToControllerChannelManager broker=${config.brokerId} name=$channelName" } new BrokerToControllerRequestThread( @@ -295,6 +295,10 @@ class BrokerToControllerRequestThread( private val requestQueue = new LinkedBlockingDeque[BrokerToControllerQueueItem]() private val activeController = new AtomicReference[Node](null) + // Used for testing + @volatile + private[server] var started = false + def activeControllerAddress(): Option[Node] = { Option(activeController.get()) } @@ -304,6 +308,9 @@ class BrokerToControllerRequestThread( } def enqueue(request: BrokerToControllerQueueItem): Unit = { + if (!started) { + throw new IllegalStateException("Cannot enqueue a request if the request thread is not running") + } requestQueue.add(request) if (activeControllerAddress().isDefined) { wakeup() @@ -380,4 +387,9 @@ class BrokerToControllerRequestThread( } } } + + override def start(): Unit = { + super.start() + started = true + } } diff --git a/core/src/main/scala/kafka/server/KafkaServer.scala b/core/src/main/scala/kafka/server/KafkaServer.scala index 13b72ea44b5c9..947c23bf8c749 100755 --- a/core/src/main/scala/kafka/server/KafkaServer.scala +++ b/core/src/main/scala/kafka/server/KafkaServer.scala @@ -265,7 +265,7 @@ class KafkaServer( time = time, metrics = metrics, config = config, - channelName = "controllerForwardingChannel", + channelName = "forwarding", threadNamePrefix = threadNamePrefix, retryTimeoutMs = config.requestTimeoutMs.longValue) brokerToControllerManager.start() diff --git a/core/src/test/scala/kafka/server/BrokerToControllerRequestThreadTest.scala b/core/src/test/scala/kafka/server/BrokerToControllerRequestThreadTest.scala index 676eb349b6156..46f329eebb8c2 100644 --- a/core/src/test/scala/kafka/server/BrokerToControllerRequestThreadTest.scala +++ b/core/src/test/scala/kafka/server/BrokerToControllerRequestThreadTest.scala @@ -46,6 +46,7 @@ class BrokerToControllerRequestThreadTest { val retryTimeoutMs = 30000 val testRequestThread = new BrokerToControllerRequestThread(mockClient, new ManualMetadataUpdater(), controllerNodeProvider, config, time, "", retryTimeoutMs) + testRequestThread.started = true val completionHandler = new TestRequestCompletionHandler(None) val queueItem = BrokerToControllerQueueItem( @@ -82,6 +83,7 @@ class BrokerToControllerRequestThreadTest { val expectedResponse = RequestTestUtils.metadataUpdateWith(2, Collections.singletonMap("a", 2)) val testRequestThread = new BrokerToControllerRequestThread(mockClient, new ManualMetadataUpdater(), controllerNodeProvider, config, time, "", retryTimeoutMs = Long.MaxValue) + testRequestThread.started = true mockClient.prepareResponse(expectedResponse) val completionHandler = new TestRequestCompletionHandler(Some(expectedResponse)) @@ -123,6 +125,7 @@ class BrokerToControllerRequestThreadTest { val expectedResponse = RequestTestUtils.metadataUpdateWith(3, Collections.singletonMap("a", 2)) val testRequestThread = new BrokerToControllerRequestThread(mockClient, new ManualMetadataUpdater(), controllerNodeProvider, config, time, "", retryTimeoutMs = Long.MaxValue) + testRequestThread.started = true val completionHandler = new TestRequestCompletionHandler(Some(expectedResponse)) val queueItem = BrokerToControllerQueueItem( @@ -172,6 +175,7 @@ class BrokerToControllerRequestThreadTest { val expectedResponse = RequestTestUtils.metadataUpdateWith(3, Collections.singletonMap("a", 2)) val testRequestThread = new BrokerToControllerRequestThread(mockClient, new ManualMetadataUpdater(), controllerNodeProvider, config, time, "", retryTimeoutMs = Long.MaxValue) + testRequestThread.started = true val completionHandler = new TestRequestCompletionHandler(Some(expectedResponse)) val queueItem = BrokerToControllerQueueItem( @@ -226,6 +230,7 @@ class BrokerToControllerRequestThreadTest { Collections.singletonMap("a", 2)) val testRequestThread = new BrokerToControllerRequestThread(mockClient, new ManualMetadataUpdater(), controllerNodeProvider, config, time, "", retryTimeoutMs) + testRequestThread.started = true val completionHandler = new TestRequestCompletionHandler() val queueItem = BrokerToControllerQueueItem( @@ -283,6 +288,7 @@ class BrokerToControllerRequestThreadTest { val testRequestThread = new BrokerToControllerRequestThread(mockClient, new ManualMetadataUpdater(), controllerNodeProvider, config, time, "", retryTimeoutMs = Long.MaxValue) + testRequestThread.started = true testRequestThread.enqueue(queueItem) pollUntil(testRequestThread, () => callbackResponse.get != null) @@ -319,12 +325,38 @@ class BrokerToControllerRequestThreadTest { val testRequestThread = new BrokerToControllerRequestThread(mockClient, new ManualMetadataUpdater(), controllerNodeProvider, config, time, "", retryTimeoutMs = Long.MaxValue) + testRequestThread.started = true testRequestThread.enqueue(queueItem) pollUntil(testRequestThread, () => callbackResponse.get != null) assertNotNull(callbackResponse.get.authenticationException) } + @Test + def testThreadNotStarted(): Unit = { + // Make sure we throw if we enqueue anything while the thread is not running + val time = new MockTime() + val config = new KafkaConfig(TestUtils.createBrokerConfig(1, "localhost:2181")) + + val metadata = mock(classOf[Metadata]) + val mockClient = new MockClient(time, metadata) + + val controllerNodeProvider = mock(classOf[ControllerNodeProvider]) + + val testRequestThread = new BrokerToControllerRequestThread(mockClient, new ManualMetadataUpdater(), controllerNodeProvider, + config, time, "", retryTimeoutMs = Long.MaxValue) + + val completionHandler = new TestRequestCompletionHandler(None) + val queueItem = BrokerToControllerQueueItem( + time.milliseconds(), + new MetadataRequest.Builder(new MetadataRequestData()), + completionHandler + ) + + assertThrows(classOf[IllegalStateException], () => testRequestThread.enqueue(queueItem)) + assertEquals(0, testRequestThread.queueSize) + } + private def pollUntil( requestThread: BrokerToControllerRequestThread, condition: () => Boolean,