Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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 @@ -322,16 +322,22 @@ class BrokerMetadataListener(

private def publish(publisher: MetadataPublisher): Unit = {
val delta = _delta
_image = _delta.apply()
_delta = new MetadataDelta(_image)
if (isDebugEnabled) {
debug(s"Publishing new metadata delta $delta at offset ${_image.highestOffsetAndEpoch().offset}.")
}
publisher.publish(delta, _image)
try {
_image = _delta.apply()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Note that it is possible for _delta to include a lot of batches maybe even the entire log. I wonder that if the broker encounters an error applying a delta we want to instead rewind, generate and apply a delta per record batch.

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.

Rewinding and re-applying does sound useful for some kind of automatic error mitigation, but I think it will be a quite a bit of work. As it stands, I believe the broker can only process metadata going forward.

I can think of a degenerate case we have today where loadBatches is able to process all but one record, but delta.apply cannot complete and so we can't publish any new metadata. Like you mention, I think the only way to mitigate a situation like this would be to produce smaller deltas to reduce the blast radius of a bad record.

_delta = new MetadataDelta(_image)
if (isDebugEnabled) {
debug(s"Publishing new metadata delta $delta at offset ${_image.highestOffsetAndEpoch().offset}.")
}

// Update the metrics since the publisher handled the lastest image
brokerMetrics.lastAppliedRecordOffset.set(_highestOffset)
brokerMetrics.lastAppliedRecordTimestamp.set(_highestTimestamp)
// This publish call is done with its own try-catch and fault handler
publisher.publish(delta, _image)

// Update the metrics since the publisher handled the lastest image
brokerMetrics.lastAppliedRecordOffset.set(_highestOffset)
brokerMetrics.lastAppliedRecordTimestamp.set(_highestTimestamp)
} catch {
Comment thread
mumrah marked this conversation as resolved.
case t: Throwable => metadataLoadingFaultHandler.handleFault(s"Error applying metadata delta $delta", t)
}
}

override def handleLeaderChange(leaderAndEpoch: LeaderAndEpoch): Unit = {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ import org.apache.kafka.common.{Endpoint, Uuid}
import org.apache.kafka.image.{MetadataDelta, MetadataImage}
import org.apache.kafka.metadata.{BrokerRegistration, RecordTestUtils, VersionRange}
import org.apache.kafka.server.common.{ApiMessageAndVersion, MetadataVersion}
import org.apache.kafka.server.fault.MockFaultHandler
import org.apache.kafka.server.fault.{FaultHandler, MockFaultHandler}
import org.junit.jupiter.api.Assertions.{assertEquals, assertTrue}
import org.junit.jupiter.api.{AfterEach, Test}

Expand All @@ -45,6 +45,7 @@ class BrokerMetadataListenerTest {
metrics: BrokerServerMetrics = BrokerServerMetrics(new Metrics()),
snapshotter: Option[MetadataSnapshotter] = None,
maxBytesBetweenSnapshots: Long = 1000000L,
faultHandler: FaultHandler = metadataLoadingFaultHandler
): BrokerMetadataListener = {
new BrokerMetadataListener(
brokerId = 0,
Expand All @@ -53,7 +54,7 @@ class BrokerMetadataListenerTest {
maxBytesBetweenSnapshots = maxBytesBetweenSnapshots,
snapshotter = snapshotter,
brokerMetrics = metrics,
_metadataLoadingFaultHandler = metadataLoadingFaultHandler
_metadataLoadingFaultHandler = faultHandler
)
}

Expand Down Expand Up @@ -168,6 +169,7 @@ class BrokerMetadataListenerTest {
}

private val FOO_ID = Uuid.fromString("jj1G9utnTuCegi_gpnRgYw")
private val BAR_ID = Uuid.fromString("SzN5j0LvSEaRIJHrxfMAlg")

private def generateManyRecords(listener: BrokerMetadataListener,
endOffset: Long): Unit = {
Expand All @@ -192,6 +194,27 @@ class BrokerMetadataListenerTest {
listener.getImageRecords().get()
}

private def generateBadRecords(listener: BrokerMetadataListener,
endOffset: Long): Unit = {
listener.handleCommit(
RecordTestUtils.mockBatchReader(
endOffset,
0,
util.Arrays.asList(
new ApiMessageAndVersion(new PartitionChangeRecord().
setPartitionId(0).
setTopicId(BAR_ID).
setRemovingReplicas(Collections.singletonList(1)), 0.toShort),
new ApiMessageAndVersion(new PartitionChangeRecord().
setPartitionId(0).
setTopicId(BAR_ID).
setRemovingReplicas(Collections.emptyList()), 0.toShort)
)
)
)
listener.getImageRecords().get()
}

@Test
def testHandleCommitsWithNoSnapshotterDefined(): Unit = {
val listener = newBrokerMetadataListener(maxBytesBetweenSnapshots = 1000L)
Expand Down Expand Up @@ -289,6 +312,39 @@ class BrokerMetadataListenerTest {
assertEquals(endOffset, snapshotter.activeSnapshotOffset, "We should generate snapshot on feature update")
}

@Test
def testNoShapshotAfterError(): Unit = {
val snapshotter = new MockMetadataSnapshotter()
val faultHandler = new MockFaultHandler("metadata loading")

val listener = newBrokerMetadataListener(
snapshotter = Some(snapshotter),
maxBytesBetweenSnapshots = 1000L,
faultHandler = faultHandler)
try {
val brokerIds = 0 to 3

registerBrokers(listener, brokerIds, endOffset = 100L)
createTopicWithOnePartition(listener, replicas = brokerIds, endOffset = 200L)
listener.getImageRecords().get()
assertEquals(200L, listener.highestMetadataOffset)
assertEquals(-1L, snapshotter.prevCommittedOffset)
assertEquals(-1L, snapshotter.activeSnapshotOffset)

// Append invalid records that will normally trigger a snapshot
generateBadRecords(listener, 1000L)
assertEquals(-1L, snapshotter.prevCommittedOffset)
assertEquals(-1L, snapshotter.activeSnapshotOffset)

// Generate some records that will not throw an error, verify still no snapshots
generateManyRecords(listener, 2000L)
assertEquals(-1L, snapshotter.prevCommittedOffset)
assertEquals(-1L, snapshotter.activeSnapshotOffset)
} finally {
listener.close()
}
}

private def registerBrokers(
listener: BrokerMetadataListener,
brokerIds: Iterable[Int],
Expand Down