From bb098299ec30745ba89ad2729d57ba4562f89beb Mon Sep 17 00:00:00 2001 From: David Arthur Date: Sat, 19 Aug 2023 10:23:10 -0400 Subject: [PATCH 1/2] Don't publish until we have replayed at least one record --- .../image/loader/MetadataBatchLoader.java | 12 +++++ .../kafka/image/loader/MetadataLoader.java | 8 ++- .../image/loader/MetadataLoaderTest.java | 53 ++++++++++++++++++- 3 files changed, 70 insertions(+), 3 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/image/loader/MetadataBatchLoader.java b/metadata/src/main/java/org/apache/kafka/image/loader/MetadataBatchLoader.java index 33f4846ac7084..4ac3c44ea023f 100644 --- a/metadata/src/main/java/org/apache/kafka/image/loader/MetadataBatchLoader.java +++ b/metadata/src/main/java/org/apache/kafka/image/loader/MetadataBatchLoader.java @@ -67,6 +67,7 @@ public interface MetadataUpdater { private int numBatches; private long totalBatchElapsedNs; private TransactionState transactionState; + private boolean empty; public MetadataBatchLoader( LogContext logContext, @@ -78,6 +79,15 @@ public MetadataBatchLoader( this.time = time; this.faultHandler = faultHandler; this.callback = callback; + this.resetToImage(MetadataImage.EMPTY); + this.empty = true; + } + + /** + * @return True if this batch loader has seen at least one record. + */ + public boolean isEmpty() { + return empty; } /** @@ -88,6 +98,7 @@ public MetadataBatchLoader( */ public void resetToImage(MetadataImage image) { this.image = image; + this.empty = false; this.delta = new MetadataDelta.Builder().setImage(image).build(); this.transactionState = TransactionState.NO_TRANSACTION; this.lastOffset = image.provenance().lastContainedOffset(); @@ -241,6 +252,7 @@ private void replay(ApiMessageAndVersion record) { default: break; } + empty = false; delta.replay(record.message()); } } diff --git a/metadata/src/main/java/org/apache/kafka/image/loader/MetadataLoader.java b/metadata/src/main/java/org/apache/kafka/image/loader/MetadataLoader.java index c5ba2af0b503c..6f78a1340a868 100644 --- a/metadata/src/main/java/org/apache/kafka/image/loader/MetadataLoader.java +++ b/metadata/src/main/java/org/apache/kafka/image/loader/MetadataLoader.java @@ -212,7 +212,6 @@ private MetadataLoader( time, faultHandler, this::maybePublishMetadata); - this.batchLoader.resetToImage(this.image); this.eventQueue = new KafkaEventQueue( Time.SYSTEM, logContext, @@ -241,6 +240,11 @@ private boolean stillNeedToCatchUp(String where, long offset) { offset + ", but the high water mark is {}", where, highWaterMark.getAsLong()); return true; } + if (batchLoader.isEmpty()) { + log.info("{}: The loader is still catching up because we have not loaded a controller record as of offset " + + offset + " and high water mark is {}", where, highWaterMark.getAsLong()); + return true; + } log.info("{}: The loader finished catching up to the current high water mark of {}", where, highWaterMark.getAsLong()); catchingUp = false; @@ -387,8 +391,8 @@ public void handleLoadSnapshot(SnapshotReader reader) { image.provenance().lastContainedOffset(), NANOSECONDS.toMicros(manifest.elapsedNs())); MetadataImage image = delta.apply(manifest.provenance()); - maybePublishMetadata(delta, image, manifest); batchLoader.resetToImage(image); + maybePublishMetadata(delta, image, manifest); } catch (Throwable e) { // This is a general catch-all block where we don't expect to end up; // failure-prone operations should have individual try/catch blocks around them. diff --git a/metadata/src/test/java/org/apache/kafka/image/loader/MetadataLoaderTest.java b/metadata/src/test/java/org/apache/kafka/image/loader/MetadataLoaderTest.java index 199d907a44773..62d974b8b5c38 100644 --- a/metadata/src/test/java/org/apache/kafka/image/loader/MetadataLoaderTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/loader/MetadataLoaderTest.java @@ -18,9 +18,11 @@ package org.apache.kafka.image.loader; import org.apache.kafka.common.Uuid; +import org.apache.kafka.common.config.ConfigResource; import org.apache.kafka.common.message.SnapshotHeaderRecord; import org.apache.kafka.common.metadata.AbortTransactionRecord; import org.apache.kafka.common.metadata.BeginTransactionRecord; +import org.apache.kafka.common.metadata.ConfigRecord; import org.apache.kafka.common.metadata.EndTransactionRecord; import org.apache.kafka.common.metadata.FeatureLevelRecord; import org.apache.kafka.common.metadata.PartitionRecord; @@ -48,6 +50,7 @@ import org.junit.jupiter.params.provider.CsvSource; import org.junit.jupiter.params.provider.ValueSource; +import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.Iterator; @@ -338,8 +341,8 @@ public void testLoadEmptySnapshot() throws Exception { setHighWaterMarkAccessor(() -> OptionalLong.of(0L)). build()) { loader.installPublishers(publishers).get(); - publishers.get(0).firstPublish.get(10, TimeUnit.SECONDS); loadEmptySnapshot(loader, 200); + publishers.get(0).firstPublish.get(10, TimeUnit.SECONDS); assertEquals(200L, loader.lastAppliedOffset()); loadEmptySnapshot(loader, 300); assertEquals(300L, loader.lastAppliedOffset()); @@ -668,6 +671,7 @@ public void testPublishTransaction(boolean abortTxn) throws Exception { .setTopicId(Uuid.fromString("dMCqhcK4T5miGH5wEX7NsQ")), (short) 0) ))); loader.waitForAllEventsToBeHandled(); + publisher.firstPublish.get(30, TimeUnit.SECONDS); assertNull(publisher.latestImage.topics().getTopic("foo"), "Topic should not be visible since we started transaction"); @@ -732,6 +736,7 @@ public void testPublishTransactionWithinBatch() throws Exception { // After MetadataLoader is fixed to handle arbitrary transactions, we would expect "foo" // to be visible at this point. + publisher.firstPublish.get(30, TimeUnit.SECONDS); assertNotNull(publisher.latestImage.topics().getTopic("foo")); } faultHandler.maybeRethrowFirstException(); @@ -758,6 +763,7 @@ public void testSnapshotDuringTransaction() throws Exception { .setTopicId(Uuid.fromString("HQSM3ccPQISrHqYK_C8GpA")), (short) 0) ))); loader.waitForAllEventsToBeHandled(); + publisher.firstPublish.get(30, TimeUnit.SECONDS); assertNull(publisher.latestImage.topics().getTopic("foo")); // loading a snapshot discards any in-flight transaction @@ -773,4 +779,49 @@ public void testSnapshotDuringTransaction() throws Exception { } faultHandler.maybeRethrowFirstException(); } + + @Test + public void testNoPublishEmptyImage() throws Exception { + MockFaultHandler faultHandler = new MockFaultHandler("testNoPublishEmptyImage"); + List capturedImages = new ArrayList<>(); + CompletableFuture firstPublish = new CompletableFuture<>(); + MetadataPublisher capturingPublisher = new MetadataPublisher() { + @Override + public String name() { + return "testNoPublishEmptyImage"; + } + + @Override + public void onMetadataUpdate(MetadataDelta delta, MetadataImage newImage, LoaderManifest manifest) { + if (!firstPublish.isDone()) { + firstPublish.complete(null); + } + capturedImages.add(newImage); + } + }; + + try (MetadataLoader loader = new MetadataLoader.Builder(). + setFaultHandler(faultHandler). + setHighWaterMarkAccessor(() -> OptionalLong.of(1)). + build()) { + loader.installPublishers(Collections.singletonList(capturingPublisher)).get(); + loader.handleCommit( + MockBatchReader.newSingleBatchReader(0, 1, Collections.singletonList( + // Any record will work here + new ApiMessageAndVersion(new ConfigRecord() + .setResourceType(ConfigResource.Type.BROKER.id()) + .setResourceName("3000") + .setName("foo") + .setValue("bar"), (short) 0) + ))); + firstPublish.get(30, TimeUnit.SECONDS); + + assertFalse(capturedImages.isEmpty()); + capturedImages.forEach(metadataImage -> { + assertFalse(metadataImage.isEmpty()); + }); + + } + faultHandler.maybeRethrowFirstException(); + } } From d19fca4348137dcebb652d6f4368dd251c6fd5ae Mon Sep 17 00:00:00 2001 From: David Arthur Date: Wed, 23 Aug 2023 21:21:04 -0400 Subject: [PATCH 2/2] rename new field and method --- .../kafka/image/loader/MetadataBatchLoader.java | 15 ++++++++------- .../apache/kafka/image/loader/MetadataLoader.java | 2 +- 2 files changed, 9 insertions(+), 8 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/image/loader/MetadataBatchLoader.java b/metadata/src/main/java/org/apache/kafka/image/loader/MetadataBatchLoader.java index 4ac3c44ea023f..b1d22364cc058 100644 --- a/metadata/src/main/java/org/apache/kafka/image/loader/MetadataBatchLoader.java +++ b/metadata/src/main/java/org/apache/kafka/image/loader/MetadataBatchLoader.java @@ -67,7 +67,7 @@ public interface MetadataUpdater { private int numBatches; private long totalBatchElapsedNs; private TransactionState transactionState; - private boolean empty; + private boolean hasSeenRecord; public MetadataBatchLoader( LogContext logContext, @@ -80,25 +80,26 @@ public MetadataBatchLoader( this.faultHandler = faultHandler; this.callback = callback; this.resetToImage(MetadataImage.EMPTY); - this.empty = true; + this.hasSeenRecord = false; } /** * @return True if this batch loader has seen at least one record. */ - public boolean isEmpty() { - return empty; + public boolean hasSeenRecord() { + return hasSeenRecord; } /** * Reset the state of this batch loader to the given image. Any un-flushed state will be - * discarded. + * discarded. This is called after applying a delta and passing it back to MetadataLoader, or + * when MetadataLoader loads a snapshot. * * @param image Metadata image to reset this batch loader's state to. */ public void resetToImage(MetadataImage image) { this.image = image; - this.empty = false; + this.hasSeenRecord = true; this.delta = new MetadataDelta.Builder().setImage(image).build(); this.transactionState = TransactionState.NO_TRANSACTION; this.lastOffset = image.provenance().lastContainedOffset(); @@ -252,7 +253,7 @@ private void replay(ApiMessageAndVersion record) { default: break; } - empty = false; + hasSeenRecord = true; delta.replay(record.message()); } } diff --git a/metadata/src/main/java/org/apache/kafka/image/loader/MetadataLoader.java b/metadata/src/main/java/org/apache/kafka/image/loader/MetadataLoader.java index 6f78a1340a868..307fda8d218bd 100644 --- a/metadata/src/main/java/org/apache/kafka/image/loader/MetadataLoader.java +++ b/metadata/src/main/java/org/apache/kafka/image/loader/MetadataLoader.java @@ -240,7 +240,7 @@ private boolean stillNeedToCatchUp(String where, long offset) { offset + ", but the high water mark is {}", where, highWaterMark.getAsLong()); return true; } - if (batchLoader.isEmpty()) { + if (!batchLoader.hasSeenRecord()) { log.info("{}: The loader is still catching up because we have not loaded a controller record as of offset " + offset + " and high water mark is {}", where, highWaterMark.getAsLong()); return true;