Skip to content
Merged
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
21 changes: 10 additions & 11 deletions core/src/main/scala/kafka/server/BrokerLifecycleManager.scala
Original file line number Diff line number Diff line change
Expand Up @@ -148,10 +148,10 @@ class BrokerLifecycleManager(
private var readyToUnfence = false

/**
* List of offline directories pending to be sent.
* List of accumulated offline directories.
* This variable can only be read or written from the event queue thread.
*/
private var offlineDirsPending = Set[Uuid]()
private var offlineDirs = Set[Uuid]()

/**
* True if we sent a event queue to the active controller requesting controlled
Expand Down Expand Up @@ -299,10 +299,10 @@ class BrokerLifecycleManager(

private class OfflineDirEvent(val dir: Uuid) extends EventQueue.Event {
override def run(): Unit = {
if (offlineDirsPending.isEmpty) {
offlineDirsPending = Set(dir)
if (offlineDirs.isEmpty) {
offlineDirs = Set(dir)
} else {
offlineDirsPending = offlineDirsPending + dir
offlineDirs = offlineDirs + dir
}
if (registered) {
scheduleNextCommunicationImmediately()
Expand Down Expand Up @@ -423,15 +423,15 @@ class BrokerLifecycleManager(
setCurrentMetadataOffset(metadataOffset).
setWantFence(!readyToUnfence).
setWantShutDown(_state == BrokerState.PENDING_CONTROLLED_SHUTDOWN).
setOfflineLogDirs(offlineDirsPending.toSeq.asJava)
setOfflineLogDirs(offlineDirs.toSeq.asJava)
if (isTraceEnabled) {
trace(s"Sending broker heartbeat $data")
}
val handler = new BrokerHeartbeatResponseHandler(offlineDirsPending)
val handler = new BrokerHeartbeatResponseHandler()
_channelManager.sendRequest(new BrokerHeartbeatRequest.Builder(data), handler)
}

private class BrokerHeartbeatResponseHandler(dirsInFlight: Set[Uuid]) extends ControllerRequestCompletionHandler {
private class BrokerHeartbeatResponseHandler() extends ControllerRequestCompletionHandler {
override def onComplete(response: ClientResponse): Unit = {
if (response.authenticationException() != null) {
error(s"Unable to send broker heartbeat for $nodeId because of an " +
Expand All @@ -455,7 +455,7 @@ class BrokerLifecycleManager(
// this response handler is not invoked from the event handler thread,
// and processing a successful heartbeat response requires updating
// state, so to continue we need to schedule an event
eventQueue.prepend(new BrokerHeartbeatResponseEvent(message.data(), dirsInFlight))
eventQueue.prepend(new BrokerHeartbeatResponseEvent(message.data()))
} else {
warn(s"Broker $nodeId sent a heartbeat request but received error $errorCode.")
scheduleNextCommunicationAfterFailure()
Expand All @@ -469,10 +469,9 @@ class BrokerLifecycleManager(
}
}

private class BrokerHeartbeatResponseEvent(response: BrokerHeartbeatResponseData, dirsInFlight: Set[Uuid]) extends EventQueue.Event {
private class BrokerHeartbeatResponseEvent(response: BrokerHeartbeatResponseData) extends EventQueue.Event {
override def run(): Unit = {
failedAttempts = 0
offlineDirsPending = offlineDirsPending.diff(dirsInFlight)
_state match {
case BrokerState.STARTING =>
if (response.isCaughtUp) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -207,34 +207,35 @@ class BrokerLifecycleManagerTest {
}

@Test
def testOfflineDirsSentUntilHeartbeatSuccess(): Unit = {
def testAlwaysSendsAccumulatedOfflineDirs(): Unit = {
val ctx = new RegistrationTestContext(configProperties)
val manager = new BrokerLifecycleManager(ctx.config, ctx.time, "offline-dirs-sent-in-heartbeat-", isZkBroker = false)
val controllerNode = new Node(3000, "localhost", 8021)
ctx.controllerNodeProvider.node.set(controllerNode)

val registration = prepareResponse(ctx, new BrokerRegistrationResponse(new BrokerRegistrationResponseData().setBrokerEpoch(1000)))
val hb1 = prepareResponse[BrokerHeartbeatRequest](ctx, new BrokerHeartbeatResponse(new BrokerHeartbeatResponseData()
.setErrorCode(Errors.NOT_CONTROLLER.code())))
val hb2 = prepareResponse[BrokerHeartbeatRequest](ctx, new BrokerHeartbeatResponse(new BrokerHeartbeatResponseData()))
val hb3 = prepareResponse[BrokerHeartbeatRequest](ctx, new BrokerHeartbeatResponse(new BrokerHeartbeatResponseData()))
val heartbeats = Seq.fill(6)(prepareResponse[BrokerHeartbeatRequest](ctx, new BrokerHeartbeatResponse(new BrokerHeartbeatResponseData())))

val offlineDirs = Set(Uuid.fromString("h3sC4Yk-Q9-fd0ntJTocCA"), Uuid.fromString("ej8Q9_d2Ri6FXNiTxKFiow"))
offlineDirs.foreach(manager.propagateDirectoryFailure)

// start the manager late to prevent a race, and force expectations on the first heartbeat
manager.start(() => ctx.highestMetadataOffset.get(),
ctx.mockChannelManager, ctx.clusterId, ctx.advertisedListeners,
Collections.emptyMap(), OptionalLong.empty())

poll(ctx, manager, registration)
val dirs1 = poll(ctx, manager, hb1).data().offlineLogDirs()
val dirs2 = poll(ctx, manager, hb2).data().offlineLogDirs()
val dirs3 = poll(ctx, manager, hb3).data().offlineLogDirs()

assertEquals(offlineDirs, dirs1.asScala.toSet)
assertEquals(offlineDirs, dirs2.asScala.toSet)
assertEquals(Set.empty, dirs3.asScala.toSet)
manager.propagateDirectoryFailure(Uuid.fromString("h3sC4Yk-Q9-fd0ntJTocCA"))
poll(ctx, manager, heartbeats(0)).data()
val dirs1 = poll(ctx, manager, heartbeats(1)).data().offlineLogDirs()

manager.propagateDirectoryFailure(Uuid.fromString("ej8Q9_d2Ri6FXNiTxKFiow"))
poll(ctx, manager, heartbeats(2)).data()
val dirs2 = poll(ctx, manager, heartbeats(3)).data().offlineLogDirs()

manager.propagateDirectoryFailure(Uuid.fromString("1iF76HVNRPqC7Y4r6647eg"))
poll(ctx, manager, heartbeats(4)).data()
val dirs3 = poll(ctx, manager, heartbeats(5)).data().offlineLogDirs()

@apoorvmittal10 apoorvmittal10 Nov 23, 2023

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.

@soarez Can we please verify the flakiness in the test. I have added @RepeatedTest(100) locally for this test and can always get failures among 100 runs. Listing out 2 failures below

java.util.ConcurrentModificationException
	at java.base/java.util.ArrayDeque.nonNullElementAt(ArrayDeque.java:270)
	at java.base/java.util.ArrayDeque$DeqIterator.next(ArrayDeque.java:700)
	at kafka.server.MockNodeToControllerChannelManager.poll(MockNodeToControllerChannelManager.scala:73)
	at kafka.server.RegistrationTestContext.poll(RegistrationTestContext.scala:76)
	at kafka.server.BrokerLifecycleManagerTest.poll(BrokerLifecycleManagerTest.scala:202)
	at kafka.server.BrokerLifecycleManagerTest.testAlwaysSendsAccumulatedOfflineDirs(BrokerLifecycleManagerTest.scala:230)
	at java.base/java.lang.reflect.Method.invoke(Method.java:568)
	at java.base/java.util.stream.ForEachOps$ForEachOp$OfRef.accept(ForEachOps.java:183)
	at java.base/java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:197)
	at java.base/java.util.stream.ReferencePipeline$2$1.accept(ReferencePipeline.java:179)
	at java.base/java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:197)
	at java.base/java.util.stream.ForEachOps$ForEachOp$OfRef.accept(ForEachOps.java:183)
	at java.base/java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:197)
	at java.base/java.util.stream.IntPipeline$1$1.accept(IntPipeline.java:180)
	at java.base/java.util.stream.Streams$RangeIntSpliterator.forEachRemaining(Streams.java:104)
	at java.base/java.util.Spliterator$OfInt.forEachRemaining(Spliterator.java:711)
	at java.base/java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:509)
	at java.base/java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:499)
	at java.base/java.util.stream.ForEachOps$ForEachOp.evaluateSequential(ForEachOps.java:150)
	at java.base/java.util.stream.ForEachOps$ForEachOp$OfRef.evaluateSequential(ForEachOps.java:173)
	at java.base/java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234)
	at java.base/java.util.stream.ReferencePipeline.forEach(ReferencePipeline.java:596)
	at java.base/java.util.stream.ReferencePipeline$7$1.accept(ReferencePipeline.java:276)
rg.opentest4j.AssertionFailedError: 
Expected :Set(h3sC4Yk-Q9-fd0ntJTocCA, ej8Q9_d2Ri6FXNiTxKFiow, 1iF76HVNRPqC7Y4r6647eg)
Actual   :Set(h3sC4Yk-Q9-fd0ntJTocCA, ej8Q9_d2Ri6FXNiTxKFiow)
<Click to see difference>


	at org.junit.jupiter.api.AssertionFailureBuilder.build(AssertionFailureBuilder.java:151)
	at org.junit.jupiter.api.AssertionFailureBuilder.buildAndThrow(AssertionFailureBuilder.java:132)
	at org.junit.jupiter.api.AssertEquals.failNotEqual(AssertEquals.java:197)
	at org.junit.jupiter.api.AssertEquals.assertEquals(AssertEquals.java:182)
	at org.junit.jupiter.api.AssertEquals.assertEquals(AssertEquals.java:177)
	at org.junit.jupiter.api.Assertions.assertEquals(Assertions.java:1141)
	at kafka.server.BrokerLifecycleManagerTest.testAlwaysSendsAccumulatedOfflineDirs(BrokerLifecycleManagerTest.scala:238)
	at java.base/java.lang.reflect.Method.invoke(Method.java:568)
	at java.base/java.util.stream.ForEachOps$ForEachOp$OfRef.accept(ForEachOps.java:183)
	at java.base/java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:197)
	at java.base/java.util.stream.ReferencePipeline$2$1.accept(ReferencePipeline.java:179)
	at java.base/java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:197)
	at java.base/java.util.stream.ForEachOps$ForEachOp$OfRef.accept(ForEachOps.java:183)
	at java.base/java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:197)
	at java.base/java.util.stream.IntPipeline$1$1.accept(IntPipeline.java:180)
	at java.base/java.util.stream.Streams$RangeIntSpliterator.forEachRemaining(Streams.java:104)

cc: @cmccabe @junrao (Saw the flakiness in the run at PR build: #14699, https://ci-builds.apache.org/blue/organizations/jenkins/Kafka%2Fkafka-pr/detail/PR-14699/17/tests/)

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.

Thanks for reporting this Apoorv. Please see #14836

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.

Thanks for fixing this @soarez.


assertEquals(Set("h3sC4Yk-Q9-fd0ntJTocCA").map(Uuid.fromString), dirs1.asScala.toSet)
assertEquals(Set("h3sC4Yk-Q9-fd0ntJTocCA", "ej8Q9_d2Ri6FXNiTxKFiow").map(Uuid.fromString), dirs2.asScala.toSet)
assertEquals(Set("h3sC4Yk-Q9-fd0ntJTocCA", "ej8Q9_d2Ri6FXNiTxKFiow", "1iF76HVNRPqC7Y4r6647eg").map(Uuid.fromString), dirs3.asScala.toSet)
manager.close()
}

Expand Down