diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/util/clusters/EmbeddedKafkaCluster.java b/connect/runtime/src/test/java/org/apache/kafka/connect/util/clusters/EmbeddedKafkaCluster.java index 6f287d060d326..a8328c5069460 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/util/clusters/EmbeddedKafkaCluster.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/util/clusters/EmbeddedKafkaCluster.java @@ -263,7 +263,7 @@ public Set brokersInState(Predicate desiredState) { protected boolean hasState(KafkaServer server, Predicate desiredState) { try { - return desiredState.test(server.brokerState().get()); + return desiredState.test(server.brokerState()); } catch (Throwable e) { // Broker failed to respond. return false; diff --git a/core/src/main/scala/kafka/server/KafkaBroker.scala b/core/src/main/scala/kafka/server/KafkaBroker.scala index d47283e8061a6..3613076c88bae 100644 --- a/core/src/main/scala/kafka/server/KafkaBroker.scala +++ b/core/src/main/scala/kafka/server/KafkaBroker.scala @@ -18,7 +18,6 @@ package kafka.server import java.util -import java.util.concurrent.atomic.AtomicReference import com.yammer.metrics.core.MetricName import kafka.log.LogManager @@ -71,8 +70,11 @@ object KafkaBroker { } trait KafkaBroker extends KafkaMetricsGroup { + @volatile private var _brokerState: BrokerState = BrokerState.NOT_RUNNING + def authorizer: Option[Authorizer] - val brokerState = new AtomicReference[BrokerState](BrokerState.NOT_RUNNING) + def brokerState: BrokerState = _brokerState + protected def brokerState_= (brokerState: BrokerState): Unit = _brokerState = brokerState def clusterId: String def config: KafkaConfig def dataPlaneRequestHandlerPool: KafkaRequestHandlerPool @@ -90,7 +92,7 @@ trait KafkaBroker extends KafkaMetricsGroup { explicitMetricName(KafkaBroker.metricsPrefix, KafkaBroker.metricsTypeName, name, metricTags) } - newGauge("BrokerState", () => brokerState.get.value()) + newGauge("BrokerState", () => brokerState.value) newGauge("ClusterId", () => clusterId) newGauge("yammer-metrics-count", () => KafkaYammerMetrics.defaultRegistry.allMetrics.size) diff --git a/core/src/main/scala/kafka/server/KafkaServer.scala b/core/src/main/scala/kafka/server/KafkaServer.scala index 7aed40d1787fc..6bb134f9eca40 100755 --- a/core/src/main/scala/kafka/server/KafkaServer.scala +++ b/core/src/main/scala/kafka/server/KafkaServer.scala @@ -185,7 +185,7 @@ class KafkaServer( val canStartup = isStartingUp.compareAndSet(false, true) if (canStartup) { - brokerState.set(BrokerState.STARTING) + brokerState = BrokerState.STARTING /* setup zookeeper */ initZkClient(time) @@ -247,7 +247,7 @@ class KafkaServer( logManager = LogManager(config, initialOfflineDirs, new ZkConfigRepository(new AdminZkClient(zkClient)), kafkaScheduler, time, brokerTopicStats, logDirFailureChannel) - brokerState.set(BrokerState.RECOVERY) + brokerState = BrokerState.RECOVERY logManager.startup(zkClient.getAllTopicsInCluster()) metadataCache = MetadataCache.zkMetadataCache(config.brokerId) @@ -394,7 +394,7 @@ class KafkaServer( socketServer.startProcessingRequests(authorizerFutures) - brokerState.set(BrokerState.RUNNING) + brokerState = BrokerState.RUNNING shutdownLatch = new CountDownLatch(1) startupComplete.set(true) isStartingUp.set(false) @@ -632,7 +632,7 @@ class KafkaServer( // the shutdown. info("Starting controlled shutdown") - brokerState.set(BrokerState.PENDING_CONTROLLED_SHUTDOWN) + brokerState = BrokerState.PENDING_CONTROLLED_SHUTDOWN val shutdownSucceeded = doControlledShutdown(config.controlledShutdownMaxRetries.intValue) @@ -657,7 +657,7 @@ class KafkaServer( // `true` at the end of this method. if (shutdownLatch.getCount > 0 && isShuttingDown.compareAndSet(false, true)) { CoreUtils.swallow(controlledShutdown(), this) - brokerState.set(BrokerState.SHUTTING_DOWN) + brokerState = BrokerState.SHUTTING_DOWN if (dynamicConfigManager != null) CoreUtils.swallow(dynamicConfigManager.shutdown(), this) @@ -728,7 +728,7 @@ class KafkaServer( // Clear all reconfigurable instances stored in DynamicBrokerConfig config.dynamicConfig.clear() - brokerState.set(BrokerState.NOT_RUNNING) + brokerState = BrokerState.NOT_RUNNING startupComplete.set(false) isShuttingDown.set(false) diff --git a/core/src/test/scala/unit/kafka/server/BaseRequestTest.scala b/core/src/test/scala/unit/kafka/server/BaseRequestTest.scala index 720dbaf279241..7b51bfeff2384 100644 --- a/core/src/test/scala/unit/kafka/server/BaseRequestTest.scala +++ b/core/src/test/scala/unit/kafka/server/BaseRequestTest.scala @@ -53,7 +53,7 @@ abstract class BaseRequestTest extends IntegrationTestHarness { def anySocketServer: SocketServer = { servers.find { server => - val state = server.brokerState.get() + val state = server.brokerState state != BrokerState.NOT_RUNNING && state != BrokerState.SHUTTING_DOWN }.map(_.socketServer).getOrElse(throw new IllegalStateException("No live broker is available")) } diff --git a/core/src/test/scala/unit/kafka/server/MetadataRequestTest.scala b/core/src/test/scala/unit/kafka/server/MetadataRequestTest.scala index 01daf02245da5..0518c8186980b 100644 --- a/core/src/test/scala/unit/kafka/server/MetadataRequestTest.scala +++ b/core/src/test/scala/unit/kafka/server/MetadataRequestTest.scala @@ -309,7 +309,7 @@ class MetadataRequestTest extends AbstractMetadataRequestTest { @Test def testIsrAfterBrokerShutDownAndJoinsBack(): Unit = { def checkIsr(servers: Seq[KafkaServer], topic: String): Unit = { - val activeBrokers = servers.filter(_.brokerState.get() != BrokerState.NOT_RUNNING) + val activeBrokers = servers.filter(_.brokerState != BrokerState.NOT_RUNNING) val expectedIsr = activeBrokers.map(_.config.brokerId).toSet // Assert that topic metadata at new brokers is updated correctly @@ -355,7 +355,7 @@ class MetadataRequestTest extends AbstractMetadataRequestTest { val brokersInController = controllerMetadataResponse.get.brokers.asScala.toSeq.sortBy(_.id) // Assert that metadata is propagated correctly - servers.filter(_.brokerState.get() != BrokerState.NOT_RUNNING).foreach { broker => + servers.filter(_.brokerState != BrokerState.NOT_RUNNING).foreach { broker => TestUtils.waitUntilTrue(() => { val metadataResponse = sendMetadataRequest(MetadataRequest.Builder.allTopics.build, Some(brokerSocketServer(broker.config.brokerId))) diff --git a/core/src/test/scala/unit/kafka/server/ServerShutdownTest.scala b/core/src/test/scala/unit/kafka/server/ServerShutdownTest.scala index 0239465b648aa..76745b3fe8a4e 100755 --- a/core/src/test/scala/unit/kafka/server/ServerShutdownTest.scala +++ b/core/src/test/scala/unit/kafka/server/ServerShutdownTest.scala @@ -172,10 +172,10 @@ class ServerShutdownTest extends ZooKeeperTestHarness { // goes wrong so that awaitShutdown doesn't hang case e: Exception => assertTrue(exceptionClassTag.runtimeClass.isInstance(e), s"Unexpected exception $e") - assertEquals(BrokerState.NOT_RUNNING, server.brokerState.get()) + assertEquals(BrokerState.NOT_RUNNING, server.brokerState) } finally { - if (server.brokerState.get() != BrokerState.NOT_RUNNING) + if (server.brokerState != BrokerState.NOT_RUNNING) server.shutdown() server.awaitShutdown() } diff --git a/core/src/test/scala/unit/kafka/server/ServerStartupTest.scala b/core/src/test/scala/unit/kafka/server/ServerStartupTest.scala index cc5706e40c9c7..b8280a082c1e4 100755 --- a/core/src/test/scala/unit/kafka/server/ServerStartupTest.scala +++ b/core/src/test/scala/unit/kafka/server/ServerStartupTest.scala @@ -99,7 +99,7 @@ class ServerStartupTest extends ZooKeeperTestHarness { server = new KafkaServer(KafkaConfig.fromProps(props)) server.startup() - TestUtils.waitUntilTrue(() => server.brokerState.get() == BrokerState.RUNNING, + TestUtils.waitUntilTrue(() => server.brokerState == BrokerState.RUNNING, "waiting for the broker state to become RUNNING") val brokers = zkClient.getAllBrokersInCluster assertEquals(1, brokers.size)