Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ import kafka.common.{KafkaException, TopicAndPartition}
import collection.Set
import collection.JavaConverters._

class ControllerChannelManager(controllerContext: ControllerContext, config: KafkaConfig, time: Time, metrics: Metrics) extends Logging {
class ControllerChannelManager(controllerContext: ControllerContext, config: KafkaConfig, time: Time, metrics: Metrics, threadNamePrefix: Option[String] = None) extends Logging {
protected val brokerStateInfo = new HashMap[Int, ControllerBrokerStateInfo]
private val brokerLock = new Object
this.logIdent = "[Channel manager on controller " + config.brokerId + "]: "
Expand Down Expand Up @@ -109,7 +109,12 @@ class ControllerChannelManager(controllerContext: ControllerContext, config: Kaf
Selectable.USE_DEFAULT_BUFFER_SIZE
)
}
val requestThread = new RequestSendThread(config.brokerId, controllerContext, broker, messageQueue, networkClient, brokerNode, config, time)
val threadName = threadNamePrefix match {
case None => "Controller-%d-to-broker-%d-send-thread".format(config.brokerId, broker.id)
case Some(name) => "%s:Controller-%d-to-broker-%d-send-thread".format(name,config.brokerId, broker.id)
}

val requestThread = new RequestSendThread(config.brokerId, controllerContext, broker, messageQueue, networkClient, brokerNode, config, time, threadName)
requestThread.setDaemon(false)
brokerStateInfo.put(broker.id, new ControllerBrokerStateInfo(networkClient, brokerNode, broker, messageQueue, requestThread))
}
Expand Down Expand Up @@ -141,8 +146,9 @@ class RequestSendThread(val controllerId: Int,
val networkClient: NetworkClient,
val brokerNode: Node,
val config: KafkaConfig,
val time: Time)
extends ShutdownableThread("Controller-%d-to-broker-%d-send-thread".format(controllerId, toBroker.id)) {
val time: Time,
name: String)
extends ShutdownableThread(name = name) {

private val lock = new Object()
private val stateChangeLogger = KafkaController.stateChangeLogger
Expand Down
4 changes: 2 additions & 2 deletions core/src/main/scala/kafka/controller/KafkaController.scala
Original file line number Diff line number Diff line change
Expand Up @@ -153,7 +153,7 @@ object KafkaController extends Logging {
}
}

class KafkaController(val config : KafkaConfig, zkClient: ZkClient, val brokerState: BrokerState, time: Time, metrics: Metrics) extends Logging with KafkaMetricsGroup {
class KafkaController(val config : KafkaConfig, zkClient: ZkClient, val brokerState: BrokerState, time: Time, metrics: Metrics, threadNamePrefix: Option[String] = None) extends Logging with KafkaMetricsGroup {
this.logIdent = "[Controller " + config.brokerId + "]: "
private var isRunning = true
private val stateChangeLogger = KafkaController.stateChangeLogger
Expand Down Expand Up @@ -815,7 +815,7 @@ class KafkaController(val config : KafkaConfig, zkClient: ZkClient, val brokerSt
}

private def startChannelManager() {
controllerContext.controllerChannelManager = new ControllerChannelManager(controllerContext, config, time, metrics)
controllerContext.controllerChannelManager = new ControllerChannelManager(controllerContext, config, time, metrics, threadNamePrefix)
controllerContext.controllerChannelManager.startup()
}

Expand Down
4 changes: 2 additions & 2 deletions core/src/main/scala/kafka/server/KafkaServer.scala
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,7 @@ object KafkaServer {
* Represents the lifecycle of a single Kafka broker. Handles all functionality required
* to start up and shutdown a single Kafka node.
*/
class KafkaServer(val config: KafkaConfig, time: Time = SystemTime) extends Logging with KafkaMetricsGroup {
class KafkaServer(val config: KafkaConfig, time: Time = SystemTime, threadNamePrefix: Option[String] = None) extends Logging with KafkaMetricsGroup {
private val startupComplete = new AtomicBoolean(false)
private val isShuttingDown = new AtomicBoolean(false)
private val isStartingUp = new AtomicBoolean(false)
Expand Down Expand Up @@ -183,7 +183,7 @@ class KafkaServer(val config: KafkaConfig, time: Time = SystemTime) extends Logg
replicaManager.startup()

/* start kafka controller */
kafkaController = new KafkaController(config, zkClient, brokerState, kafkaMetricsTime, metrics)
kafkaController = new KafkaController(config, zkClient, brokerState, kafkaMetricsTime, metrics, threadNamePrefix)
kafkaController.startup()

/* start kafka coordinator */
Expand Down
10 changes: 8 additions & 2 deletions core/src/main/scala/kafka/server/ReplicaFetcherManager.scala
Original file line number Diff line number Diff line change
Expand Up @@ -21,12 +21,18 @@ import kafka.cluster.BrokerEndPoint
import org.apache.kafka.common.metrics.Metrics
import org.apache.kafka.common.utils.Time

class ReplicaFetcherManager(brokerConfig: KafkaConfig, replicaMgr: ReplicaManager, metrics: Metrics, time: Time)
class ReplicaFetcherManager(brokerConfig: KafkaConfig, replicaMgr: ReplicaManager, metrics: Metrics, time: Time, threadNamePrefix: Option[String] = None)
extends AbstractFetcherManager("ReplicaFetcherManager on broker " + brokerConfig.brokerId,
"Replica", brokerConfig.numReplicaFetchers) {

override def createFetcherThread(fetcherId: Int, sourceBroker: BrokerEndPoint): AbstractFetcherThread = {
new ReplicaFetcherThread("ReplicaFetcherThread-%d-%d".format(fetcherId, sourceBroker.id), sourceBroker, brokerConfig,
val threadName = threadNamePrefix match {
case None =>
"ReplicaFetcherThread-%d-%d".format(fetcherId, sourceBroker.id)
case Some(p) =>
"%s:ReplicaFetcherThread-%d-%d".format(p, fetcherId, sourceBroker.id)
}
new ReplicaFetcherThread(threadName, sourceBroker, brokerConfig,
replicaMgr, metrics, time)
}

Expand Down
5 changes: 3 additions & 2 deletions core/src/main/scala/kafka/server/ReplicaManager.scala
Original file line number Diff line number Diff line change
Expand Up @@ -102,13 +102,14 @@ class ReplicaManager(val config: KafkaConfig,
val zkClient: ZkClient,
scheduler: Scheduler,
val logManager: LogManager,
val isShuttingDown: AtomicBoolean) extends Logging with KafkaMetricsGroup {
val isShuttingDown: AtomicBoolean,
threadNamePrefix: Option[String] = None) extends Logging with KafkaMetricsGroup {
/* epoch of the controller that last changed the leader */
@volatile var controllerEpoch: Int = KafkaController.InitialControllerEpoch - 1
private val localBrokerId = config.brokerId
private val allPartitions = new Pool[(String, Int), Partition]
private val replicaStateChangeLock = new Object
val replicaFetcherManager = new ReplicaFetcherManager(config, this, metrics, jTime)
val replicaFetcherManager = new ReplicaFetcherManager(config, this, metrics, jTime, threadNamePrefix)
private val highWatermarkCheckPointThreadStarted = new AtomicBoolean(false)
val highWatermarkCheckpoints = config.logDirs.map(dir => (new File(dir).getAbsolutePath, new OffsetCheckpoint(new File(dir, ReplicaManager.HighWatermarkFilename)))).toMap
private var hwThreadInitialized = false
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,7 @@ class ReplicaManagerTest {
val time: MockTime = new MockTime()
val jTime = new JMockTime
val rm = new ReplicaManager(config, new Metrics, time, jTime, zkClient, new MockScheduler(time), mockLogMgr,
new AtomicBoolean(false))
new AtomicBoolean(false), Option(this.getClass.getName))
val produceRequest = new ProducerRequest(1, "client 1", 3, 1000, SerializationTestUtils.topicDataProducerRequest)
def callback(responseStatus: Map[TopicAndPartition, ProducerResponseStatus]) = {
assert(responseStatus.values.head.error == Errors.INVALID_REQUIRED_ACKS.code)
Expand All @@ -93,7 +93,7 @@ class ReplicaManagerTest {

rm.shutdown(false)

TestUtils.verifyNonDaemonThreadsStatus
TestUtils.verifyNonDaemonThreadsStatus(this.getClass.getName)

}
}
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@ class ServerGenerateBrokerIdTest extends ZooKeeperTestHarness {

@Test
def testAutoGenerateBrokerId() {
var server1 = new KafkaServer(config1)
var server1 = new KafkaServer(config1, threadNamePrefix = Option(this.getClass.getName))
server1.startup()
server1.shutdown()
assertTrue(verifyBrokerMetadata(config1.logDirs, 1001))
Expand All @@ -52,14 +52,14 @@ class ServerGenerateBrokerIdTest extends ZooKeeperTestHarness {
assertEquals(server1.config.brokerId, 1001)
server1.shutdown()
CoreUtils.rm(server1.config.logDirs)
TestUtils.verifyNonDaemonThreadsStatus
TestUtils.verifyNonDaemonThreadsStatus(this.getClass.getName)
}

@Test
def testUserConfigAndGeneratedBrokerId() {
// start the server with broker.id as part of config
val server1 = new KafkaServer(config1)
val server2 = new KafkaServer(config2)
val server1 = new KafkaServer(config1, threadNamePrefix = Option(this.getClass.getName))
val server2 = new KafkaServer(config2, threadNamePrefix = Option(this.getClass.getName))
val props3 = TestUtils.createBrokerConfig(-1, zkConnect)
val config3 = KafkaConfig.fromProps(props3)
val server3 = new KafkaServer(config3)
Expand All @@ -78,7 +78,7 @@ class ServerGenerateBrokerIdTest extends ZooKeeperTestHarness {
CoreUtils.rm(server1.config.logDirs)
CoreUtils.rm(server2.config.logDirs)
CoreUtils.rm(server3.config.logDirs)
TestUtils.verifyNonDaemonThreadsStatus
TestUtils.verifyNonDaemonThreadsStatus(this.getClass.getName)
}

@Test
Expand All @@ -88,37 +88,37 @@ class ServerGenerateBrokerIdTest extends ZooKeeperTestHarness {
"," + TestUtils.tempDir().getAbsolutePath
props1.setProperty("log.dir",logDirs)
config1 = KafkaConfig.fromProps(props1)
var server1 = new KafkaServer(config1)
var server1 = new KafkaServer(config1, threadNamePrefix = Option(this.getClass.getName))
server1.startup()
server1.shutdown()
assertTrue(verifyBrokerMetadata(config1.logDirs, 1001))
// addition to log.dirs after generation of a broker.id from zk should be copied over
val newLogDirs = props1.getProperty("log.dir") + "," + TestUtils.tempDir().getAbsolutePath
props1.setProperty("log.dir",newLogDirs)
config1 = KafkaConfig.fromProps(props1)
server1 = new KafkaServer(config1)
server1 = new KafkaServer(config1, threadNamePrefix = Option(this.getClass.getName))
server1.startup()
server1.shutdown()
assertTrue(verifyBrokerMetadata(config1.logDirs, 1001))
CoreUtils.rm(server1.config.logDirs)
TestUtils.verifyNonDaemonThreadsStatus
TestUtils.verifyNonDaemonThreadsStatus(this.getClass.getName)
}

@Test
def testConsistentBrokerIdFromUserConfigAndMetaProps() {
// check if configured brokerId and stored brokerId are equal or throw InconsistentBrokerException
var server1 = new KafkaServer(config1) //auto generate broker Id
var server1 = new KafkaServer(config1, threadNamePrefix = Option(this.getClass.getName)) //auto generate broker Id
server1.startup()
server1.shutdown()
server1 = new KafkaServer(config2) // user specified broker id
server1 = new KafkaServer(config2, threadNamePrefix = Option(this.getClass.getName)) // user specified broker id
try {
server1.startup()
} catch {
case e: kafka.common.InconsistentBrokerIdException => //success
}
server1.shutdown()
CoreUtils.rm(server1.config.logDirs)
TestUtils.verifyNonDaemonThreadsStatus
TestUtils.verifyNonDaemonThreadsStatus(this.getClass.getName)
}

def verifyBrokerMetadata(logDirs: Seq[String], brokerId: Int): Boolean = {
Expand Down
12 changes: 4 additions & 8 deletions core/src/test/scala/unit/kafka/server/ServerShutdownTest.scala
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ class ServerShutdownTest extends ZooKeeperTestHarness {

@Test
def testCleanShutdown() {
var server = new KafkaServer(config)
var server = new KafkaServer(config, threadNamePrefix = Option(this.getClass.getName))
server.startup()
var producer = TestUtils.createProducer[Int, String](TestUtils.getBrokerListStrFromServers(Seq(server)),
encoder = classOf[StringEncoder].getName,
Expand Down Expand Up @@ -109,7 +109,7 @@ class ServerShutdownTest extends ZooKeeperTestHarness {
val newProps = TestUtils.createBrokerConfig(0, zkConnect)
newProps.setProperty("delete.topic.enable", "true")
val newConfig = KafkaConfig.fromProps(newProps)
val server = new KafkaServer(newConfig)
val server = new KafkaServer(newConfig, threadNamePrefix = Option(this.getClass.getName))
server.startup()
server.shutdown()
server.awaitShutdown()
Expand All @@ -122,7 +122,7 @@ class ServerShutdownTest extends ZooKeeperTestHarness {
val newProps = TestUtils.createBrokerConfig(0, zkConnect)
newProps.setProperty("zookeeper.connect", "fakehostthatwontresolve:65535")
val newConfig = KafkaConfig.fromProps(newProps)
val server = new KafkaServer(newConfig)
val server = new KafkaServer(newConfig, threadNamePrefix = Option(this.getClass.getName))
try {
server.startup()
fail("Expected KafkaServer setup to fail, throw exception")
Expand All @@ -146,11 +146,7 @@ class ServerShutdownTest extends ZooKeeperTestHarness {
}

private[this] def isNonDaemonKafkaThread(t: Thread): Boolean = {
val threadName = Option(t.getClass.getCanonicalName)
.getOrElse(t.getClass.getName())
.toLowerCase

!t.isDaemon && t.isAlive && threadName.startsWith("kafka")
!t.isDaemon && t.isAlive && t.getName.startsWith(this.getClass.getName)
}

def verifyNonDaemonThreadsStatus() {
Expand Down
4 changes: 2 additions & 2 deletions core/src/test/scala/unit/kafka/utils/TestUtils.scala
Original file line number Diff line number Diff line change
Expand Up @@ -756,10 +756,10 @@ object TestUtils extends Logging {
ZkUtils.pathExists(zkClient, ZkUtils.ReassignPartitionsPath)
}

def verifyNonDaemonThreadsStatus() {
def verifyNonDaemonThreadsStatus(threadNamePrefix: String) {
assertEquals(0, Thread.getAllStackTraces.keySet().toArray
.map(_.asInstanceOf[Thread])
.count(t => !t.isDaemon && t.isAlive && t.getClass.getCanonicalName.toLowerCase.startsWith("kafka")))
.count(t => !t.isDaemon && t.isAlive && t.getName.startsWith(threadNamePrefix)))
}

/**
Expand Down