Skip to content
Merged
Show file tree
Hide file tree
Changes from 4 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
175 changes: 116 additions & 59 deletions core/src/main/scala/kafka/network/SocketServer.scala
Original file line number Diff line number Diff line change
Expand Up @@ -100,29 +100,28 @@ class SocketServer(val config: KafkaConfig,

private var nextProcessorId = 0
private var connectionQuotas: ConnectionQuotas = _
private var startedProcessingRequests = false
private var stoppedProcessingRequests = false

/**
* Start the socket server. Acceptors for all the listeners are started. Processors
* are started if `startupProcessors` is true. If not, processors are only started when
* [[kafka.network.SocketServer#startDataPlaneProcessors()]] or
* [[kafka.network.SocketServer#startControlPlaneProcessor()]] is invoked. Delayed starting of processors
* is used to delay processing client connections until server is fully initialized, e.g.
* to ensure that all credentials have been loaded before authentications are performed.
* Acceptors are always started during `startup` so that the bound port is known when this
* method completes even when ephemeral ports are used. Incoming connections on this server

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

These two lines are still true, but removed from the comment?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Partially. The acceptors are not started but start to listen. Let me rework the comment to include the part about the bound port though.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

These two lines are still true, but removed from the comment?

* are processed when processors start up and invoke [[org.apache.kafka.common.network.Selector#poll]].
* Starts the socket server and creates all the Acceptors and the Processors. The Acceptors
* start listening at this stage. Acceptors and Processors are started if `startProcessingRequests`
* is true. If not, acceptors and processors are only started when
* [[kafka.network.SocketServer#startProcessingRequests()]] is invoked. Delayed starting of
* acceptors and processors is used to delay processing client connections until server is
* fully initialized, e.g. to ensure that all credentials have been loaded before authentications
* are performed. Incoming connections on this server are processed when processors start up
* and invoke [[org.apache.kafka.common.network.Selector#poll]].
*
* @param startupProcessors Flag indicating whether `Processor`s must be started.
* @param startProcessingRequests Flag indicating whether `Processor`s must be started.
*/
def startup(startupProcessors: Boolean = true): Unit = {
def startup(startProcessingRequests: Boolean = true): Unit = {
this.synchronized {
connectionQuotas = new ConnectionQuotas(config, time)
createControlPlaneAcceptorAndProcessor(config.controlPlaneListener)
createDataPlaneAcceptorsAndProcessors(config.numNetworkThreads, config.dataPlaneListeners)
if (startupProcessors) {
startControlPlaneProcessor(Map.empty)
startDataPlaneProcessors(Map.empty)
if (startProcessingRequests) {
this.startProcessingRequests()
}
}

Expand Down Expand Up @@ -160,66 +159,101 @@ class SocketServer(val config: KafkaConfig,
Option(metrics.metric(metricName)).fold(0.0)(m => m.metricValue.asInstanceOf[Double])
}.getOrElse(0.0)
})
info(s"Started ${dataPlaneAcceptors.size} acceptor threads for data-plane")
if (controlPlaneAcceptorOpt.isDefined)
info("Started control-plane acceptor thread")
}

/**
* Starts processors of all the data-plane acceptors of this server if they have not already been started.
* This method is used for delayed starting of data-plane processors if [[kafka.network.SocketServer#startup]]
* was invoked with `startupProcessors=false`.
* Start processing requests and new connections. This method is used for delayed starting of
* data-plane processors if [[kafka.network.SocketServer#startup]] was invoked with

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this is not just data-plane processors?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Correct. Let me rework the comment.

* `startProcessingRequests=false`.
*
* Before starting processors for each endpoint, we ensure that authorizer has all the metadata
* to authorize requests on that endpoint by waiting on the provided future. We start inter-broker listener
* before other listeners. This allows authorization metadata for other listeners to be stored in Kafka topics
* in this cluster.
* to authorize requests on that endpoint by waiting on the provided future. We start inter-broker
* listener before other listeners. This allows authorization metadata for other listeners to be
* stored in Kafka topics in this cluster.
*
* @param authorizerFutures

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Add a description?

*/
def startProcessingRequests(authorizerFutures: Map[Endpoint, CompletableFuture[Void]] = Map.empty): Unit = {
info("Starting socket server acceptors and processors")
this.synchronized {
if (!startedProcessingRequests) {
startControlPlaneProcessorAndAcceptor(authorizerFutures)
startDataPlaneProcessorsAndAcceptors(authorizerFutures)
startedProcessingRequests = true
} else {
info("Socket server acceptors and processors already started")
}
}
info("Started socket server acceptors and processors")
}

/**
* Starts processors of the provided acceptor and the acceptor itself.
*
* Before starting them, we ensure that authorizer has all the metadata to authorize
* requests on that endpoint by waiting on the provided future.
*/
private def startAcceptorAndProcessors(threadPrefix: String,
endpoint: EndPoint,
acceptor: Acceptor,
authorizerFutures: Map[Endpoint, CompletableFuture[Void]] = Map.empty): Unit = {
debug(s"Wait for authorizer to complete start up on listener ${endpoint.listenerName}")
waitForAuthorizerFuture(acceptor, authorizerFutures)
debug(s"Start processors on listener ${endpoint.listenerName}")
acceptor.startProcessors(threadPrefix)
debug(s"Start acceptor thread on listener ${endpoint.listenerName}")
if (!acceptor.isStarted()) {
KafkaThread.nonDaemon(
s"${threadPrefix}-kafka-socket-acceptor-${endpoint.listenerName}-${endpoint.securityProtocol}-${endpoint.port}",
acceptor
).start()
acceptor.awaitStartup()
}
info(s"Started $threadPrefix acceptor and processor(s) for endpoint : ${endpoint.listenerName}")
}

/**
* Starts processors of all the data-plane acceptors and all the acceptors of this server.
*
* We start inter-broker listener before other listeners. This allows authorization metadata for
* other listeners to be stored in Kafka topics in this cluster.
*/
def startDataPlaneProcessors(authorizerFutures: Map[Endpoint, CompletableFuture[Void]] = Map.empty): Unit = synchronized {
private def startDataPlaneProcessorsAndAcceptors(authorizerFutures: Map[Endpoint, CompletableFuture[Void]]): Unit = {
val interBrokerListener = dataPlaneAcceptors.asScala.keySet
.find(_.listenerName == config.interBrokerListenerName)
.getOrElse(throw new IllegalStateException(s"Inter-broker listener ${config.interBrokerListenerName} not found, endpoints=${dataPlaneAcceptors.keySet}"))
val orderedAcceptors = List(dataPlaneAcceptors.get(interBrokerListener)) ++
dataPlaneAcceptors.asScala.filter { case (k, _) => k != interBrokerListener }.values
orderedAcceptors.foreach { acceptor =>
val endpoint = acceptor.endPoint
debug(s"Wait for authorizer to complete start up on listener ${endpoint.listenerName}")
waitForAuthorizerFuture(acceptor, authorizerFutures)
debug(s"Start processors on listener ${endpoint.listenerName}")
acceptor.startProcessors(DataPlaneThreadPrefix)
startAcceptorAndProcessors(DataPlaneThreadPrefix, endpoint, acceptor, authorizerFutures)
}
info(s"Started data-plane processors for ${dataPlaneAcceptors.size} acceptors")
}

/**
* Start the processor of control-plane acceptor of this server if it has not already been started.
* This method is used for delayed starting of control-plane processor if [[kafka.network.SocketServer#startup]]
* was invoked with `startupProcessors=false`.
* Start the processor of control-plane acceptor and the acceptor of this server.
*/
def startControlPlaneProcessor(authorizerFutures: Map[Endpoint, CompletableFuture[Void]] = Map.empty): Unit = synchronized {
private def startControlPlaneProcessorAndAcceptor(authorizerFutures: Map[Endpoint, CompletableFuture[Void]]): Unit = {
controlPlaneAcceptorOpt.foreach { controlPlaneAcceptor =>
waitForAuthorizerFuture(controlPlaneAcceptor, authorizerFutures)
controlPlaneAcceptor.startProcessors(ControlPlaneThreadPrefix)
info(s"Started control-plane processor for the control-plane acceptor")
val endpoint = config.controlPlaneListener.get
startAcceptorAndProcessors(ControlPlaneThreadPrefix, endpoint, controlPlaneAcceptor, authorizerFutures)
}
}

private def endpoints = config.listeners.map(l => l.listenerName -> l).toMap

private def createDataPlaneAcceptorsAndProcessors(dataProcessorsPerListener: Int,
endpoints: Seq[EndPoint]): Unit = synchronized {
endpoints: Seq[EndPoint]): Unit = {
endpoints.foreach { endpoint =>
connectionQuotas.addListener(config, endpoint.listenerName)
val dataPlaneAcceptor = createAcceptor(endpoint, DataPlaneMetricPrefix)
addDataPlaneProcessors(dataPlaneAcceptor, endpoint, dataProcessorsPerListener)
KafkaThread.nonDaemon(s"data-plane-kafka-socket-acceptor-${endpoint.listenerName}-${endpoint.securityProtocol}-${endpoint.port}", dataPlaneAcceptor).start()
dataPlaneAcceptor.awaitStartup()
dataPlaneAcceptors.put(endpoint, dataPlaneAcceptor)
info(s"Created data-plane acceptor and processors for endpoint : $endpoint")
info(s"Created data-plane acceptor and processors for endpoint : ${endpoint.listenerName}")
}
}

private def createControlPlaneAcceptorAndProcessor(endpointOpt: Option[EndPoint]): Unit = synchronized {
private def createControlPlaneAcceptorAndProcessor(endpointOpt: Option[EndPoint]): Unit = {
endpointOpt.foreach { endpoint =>
connectionQuotas.addListener(config, endpoint.listenerName)
val controlPlaneAcceptor = createAcceptor(endpoint, ControlPlaneMetricPrefix)
Expand All @@ -231,20 +265,18 @@ class SocketServer(val config: KafkaConfig,
controlPlaneRequestChannelOpt.foreach(_.addProcessor(controlPlaneProcessor))
nextProcessorId += 1
controlPlaneAcceptor.addProcessors(listenerProcessors, ControlPlaneThreadPrefix)
KafkaThread.nonDaemon(s"${ControlPlaneThreadPrefix}-kafka-socket-acceptor-${endpoint.listenerName}-${endpoint.securityProtocol}-${endpoint.port}", controlPlaneAcceptor).start()
controlPlaneAcceptor.awaitStartup()
info(s"Created control-plane acceptor and processor for endpoint : $endpoint")
info(s"Created control-plane acceptor and processor for endpoint : ${endpoint.listenerName}")
}
}

private def createAcceptor(endPoint: EndPoint, metricPrefix: String) : Acceptor = synchronized {
private def createAcceptor(endPoint: EndPoint, metricPrefix: String) : Acceptor = {
val sendBufferSize = config.socketSendBufferBytes
val recvBufferSize = config.socketReceiveBufferBytes
val brokerId = config.brokerId
new Acceptor(endPoint, sendBufferSize, recvBufferSize, brokerId, connectionQuotas, metricPrefix)
}

private def addDataPlaneProcessors(acceptor: Acceptor, endpoint: EndPoint, newProcessorsPerListener: Int): Unit = synchronized {
private def addDataPlaneProcessors(acceptor: Acceptor, endpoint: EndPoint, newProcessorsPerListener: Int): Unit = {
val listenerName = endpoint.listenerName
val securityProtocol = endpoint.securityProtocol
val listenerProcessors = new ArrayBuffer[Processor]()
Expand All @@ -261,13 +293,13 @@ class SocketServer(val config: KafkaConfig,
/**
* Stop processing requests and new connections.
*/
def stopProcessingRequests() = {
def stopProcessingRequests(): Unit = {
info("Stopping socket server request processors")
this.synchronized {
dataPlaneAcceptors.asScala.values.foreach(_.shutdown())
dataPlaneAcceptors.asScala.values.foreach(_.awaitShutdown())
controlPlaneAcceptorOpt.foreach(_.shutdown())
dataPlaneProcessors.asScala.values.foreach(_.shutdown())
controlPlaneProcessorOpt.foreach(_.shutdown())
controlPlaneAcceptorOpt.foreach(_.awaitShutdown())
dataPlaneRequestChannel.clear()
controlPlaneRequestChannelOpt.foreach(_.clear())
stoppedProcessingRequests = true
Expand All @@ -289,7 +321,7 @@ class SocketServer(val config: KafkaConfig,
* Shutdown the socket server. If still processing requests, shutdown
* acceptors and processors first.
*/
def shutdown() = {
def shutdown(): Unit = {
info("Shutting down socket server")
this.synchronized {
if (!stoppedProcessingRequests)
Expand Down Expand Up @@ -317,14 +349,20 @@ class SocketServer(val config: KafkaConfig,
def addListeners(listenersAdded: Seq[EndPoint]): Unit = synchronized {
info(s"Adding data-plane listeners for endpoints $listenersAdded")
createDataPlaneAcceptorsAndProcessors(config.numNetworkThreads, listenersAdded)
startDataPlaneProcessors()
listenersAdded.foreach { endpoint =>
val acceptor = dataPlaneAcceptors.get(endpoint)
startAcceptorAndProcessors(DataPlaneThreadPrefix, endpoint, acceptor)
}
}

def removeListeners(listenersRemoved: Seq[EndPoint]): Unit = synchronized {
info(s"Removing data-plane listeners for endpoints $listenersRemoved")
listenersRemoved.foreach { endpoint =>
connectionQuotas.removeListener(config, endpoint.listenerName)
dataPlaneAcceptors.asScala.remove(endpoint).foreach(_.shutdown())
dataPlaneAcceptors.asScala.remove(endpoint).foreach { acceptor =>
acceptor.shutdown()
acceptor.awaitShutdown()
}
}
}

Expand Down Expand Up @@ -422,14 +460,23 @@ private[kafka] abstract class AbstractServerThread(connectionQuotas: ConnectionQ
def wakeup(): Unit

/**
* Initiates a graceful shutdown by signaling to stop and waiting for the shutdown to complete
* Initiates a graceful shutdown by signaling to stop
*/
def shutdown(): Unit = {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we rename this method to be initiateShutdown() to be consistent with kafka.utils.ShutdownableThread?

if (alive.getAndSet(false))
wakeup()
shutdownLatch.await()
}

/**
* Wait for the thread to completely shutdown
*/
def awaitShutdown(): Unit = shutdownLatch.await

/**
* Returns true if the thread is completely started
*/
def isStarted(): Boolean = startupLatch.getCount == 0

/**
* Wait for the thread to completely start up
*/
Expand Down Expand Up @@ -498,8 +545,10 @@ private[kafka] class Acceptor(val endPoint: EndPoint,

private def startProcessors(processors: Seq[Processor], processorThreadPrefix: String): Unit = synchronized {
processors.foreach { processor =>
KafkaThread.nonDaemon(s"${processorThreadPrefix}-kafka-network-thread-$brokerId-${endPoint.listenerName}-${endPoint.securityProtocol}-${processor.id}",
processor).start()
KafkaThread.nonDaemon(
s"${processorThreadPrefix}-kafka-network-thread-$brokerId-${endPoint.listenerName}-${endPoint.securityProtocol}-${processor.id}",
processor
).start()
}
}

Expand All @@ -510,6 +559,7 @@ private[kafka] class Acceptor(val endPoint: EndPoint,
val toRemove = processors.takeRight(removeCount)
processors.remove(processors.size - removeCount, removeCount)
toRemove.foreach(_.shutdown())
toRemove.foreach(_.awaitShutdown())
toRemove.foreach(processor => requestChannel.removeProcessor(processor.id))
}

Expand All @@ -520,6 +570,13 @@ private[kafka] class Acceptor(val endPoint: EndPoint,
}
}

override def awaitShutdown(): Unit = {
super.awaitShutdown()
synchronized {
processors.foreach(_.awaitShutdown())
}
}

/**
* Accept loop that checks for new connection attempts
*/
Expand All @@ -530,7 +587,6 @@ private[kafka] class Acceptor(val endPoint: EndPoint,
var currentProcessorIndex = 0
while (isRunning) {
try {

val ready = nioSelector.select(500)
if (ready > 0) {
val keys = nioSelector.selectedKeys()
Expand All @@ -542,7 +598,6 @@ private[kafka] class Acceptor(val endPoint: EndPoint,

if (key.isAcceptable) {
accept(key).foreach { socketChannel =>

// Assign the channel to the next processor (using round-robin) to which the
// channel can be added without blocking. If newConnections queue is full on
// all processors, block until the last one is able to accept a connection.
Expand Down Expand Up @@ -1038,6 +1093,9 @@ private[kafka] class Processor(val id: Int,
* Close the selector and all open connections
*/
private def closeAll(): Unit = {
// Clear to unblock blocked acceptors
newConnections.asScala.foreach(_.close())

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The blocked acceptor would then add another connection to this list right? Do we close that one?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

No, we don't close that one. Let me rework this.

newConnections.clear()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

clear() is unnecessary since we would expect the loop to clear (i.e. we shouldn't have code that clears without closing).

selector.channels.asScala.foreach { channel =>
close(channel.id)
}
Expand Down Expand Up @@ -1102,7 +1160,6 @@ private[kafka] class Processor(val id: Int,
removeMetric("IdlePercent", Map("networkProcessor" -> id.toString))
metrics.removeMetric(expiredConnectionsKilledCountMetricName)
}

}

class ConnectionQuotas(config: KafkaConfig, time: Time) extends Logging {
Expand Down
6 changes: 3 additions & 3 deletions core/src/main/scala/kafka/server/KafkaServer.scala
Original file line number Diff line number Diff line change
Expand Up @@ -270,7 +270,7 @@ class KafkaServer(val config: KafkaConfig, time: Time = Time.SYSTEM, threadNameP
// Delay starting processors until the end of the initialization sequence to ensure
// that credentials have been loaded before processing authentications.
socketServer = new SocketServer(config, metrics, time, credentialProvider)
socketServer.startup(startupProcessors = false)
socketServer.startup(startProcessingRequests = false)

/* start replica manager */
replicaManager = createReplicaManager(isShuttingDown)
Expand Down Expand Up @@ -352,8 +352,8 @@ class KafkaServer(val config: KafkaConfig, time: Time = Time.SYSTEM, threadNameP
dynamicConfigManager = new DynamicConfigManager(zkClient, dynamicConfigHandlers)
dynamicConfigManager.startup()

socketServer.startControlPlaneProcessor(authorizerFutures)
socketServer.startDataPlaneProcessors(authorizerFutures)
socketServer.startProcessingRequests(authorizerFutures)

brokerState.newState(RunningAsBroker)
shutdownLatch = new CountDownLatch(1)
startupComplete.set(true)
Expand Down
Loading