From ec68962e2132032b0750f1c6492a5a1e2c50b4f5 Mon Sep 17 00:00:00 2001 From: Owen Leung Date: Sat, 27 May 2023 18:56:42 +0800 Subject: [PATCH 1/6] KAFKA-14712: Produce correct error msg with correct metadataversion Fix the confusing error message --- .../kafka/image/writer/ImageWriterOptions.java | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java b/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java index 0d4550932bc1a..b77d52dcc10b1 100644 --- a/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java +++ b/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java @@ -29,6 +29,7 @@ public final class ImageWriterOptions { public static class Builder { private MetadataVersion metadataVersion; + private MetadataVersion orgMetadataVersion; private Consumer lossHandler = e -> { throw e; }; @@ -45,6 +46,7 @@ public Builder setMetadataVersion(MetadataVersion metadataVersion) { if (metadataVersion.isLessThan(MetadataVersion.MINIMUM_BOOTSTRAP_VERSION)) { // When writing an image, all versions less than 3.3-IV0 are treated as 3.0-IV1. // This is because those versions don't support FeatureLevelRecord. + setOrgMetadataVersion(metadataVersion); setRawMetadataVersion(MetadataVersion.MINIMUM_KRAFT_VERSION); } else { setRawMetadataVersion(metadataVersion); @@ -58,6 +60,11 @@ public Builder setRawMetadataVersion(MetadataVersion metadataVersion) { return this; } + public Builder setOrgMetadataVersion(MetadataVersion orgMetadataVersion) { + this.orgMetadataVersion = orgMetadataVersion; + return this; + } + public MetadataVersion metadataVersion() { return metadataVersion; } @@ -68,7 +75,11 @@ public Builder setLossHandler(Consumer lossHandler) } public ImageWriterOptions build() { - return new ImageWriterOptions(metadataVersion, lossHandler); + if (orgMetadataVersion == null) { + return new ImageWriterOptions(metadataVersion, lossHandler); + } else { + return new ImageWriterOptions(orgMetadataVersion, lossHandler); + } } } From b94ee405d41786de97b6cc0e62ba5764b5cf1a4f Mon Sep 17 00:00:00 2001 From: Owen Leung Date: Mon, 29 May 2023 21:33:10 +0800 Subject: [PATCH 2/6] Keep the original build() logic --- .../kafka/image/writer/ImageWriterOptions.java | 17 ++++++++++------- 1 file changed, 10 insertions(+), 7 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java b/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java index b77d52dcc10b1..b908c38e6b808 100644 --- a/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java +++ b/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java @@ -75,23 +75,22 @@ public Builder setLossHandler(Consumer lossHandler) } public ImageWriterOptions build() { - if (orgMetadataVersion == null) { - return new ImageWriterOptions(metadataVersion, lossHandler); - } else { - return new ImageWriterOptions(orgMetadataVersion, lossHandler); - } + return new ImageWriterOptions(metadataVersion, lossHandler, orgMetadataVersion); } } private final MetadataVersion metadataVersion; + private final MetadataVersion orgMetadataVersion; private final Consumer lossHandler; private ImageWriterOptions( MetadataVersion metadataVersion, - Consumer lossHandler + Consumer lossHandler, + MetadataVersion orgMetadataVersion ) { this.metadataVersion = metadataVersion; this.lossHandler = lossHandler; + this.orgMetadataVersion = orgMetadataVersion; } public MetadataVersion metadataVersion() { @@ -99,7 +98,11 @@ public MetadataVersion metadataVersion() { } public void handleLoss(String loss) { - lossHandler.accept(new UnwritableMetadataException(metadataVersion, loss)); + if (orgMetadataVersion != null) { + lossHandler.accept(new UnwritableMetadataException(orgMetadataVersion, loss)); + } else { + lossHandler.accept(new UnwritableMetadataException(metadataVersion, loss)); + } } } From 59135fc9967929e6e14232cebcbfed04e89fdabb Mon Sep 17 00:00:00 2001 From: Owen Leung Date: Thu, 15 Jun 2023 14:31:28 +0800 Subject: [PATCH 3/6] add unittest for handleLoss function --- .../image/writer/ImageWriterOptions.java | 7 +++-- .../image/writer/ImageWriterOptionsTest.java | 28 +++++++++++++++++++ 2 files changed, 33 insertions(+), 2 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java b/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java index b908c38e6b808..46d923ada4cb7 100644 --- a/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java +++ b/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java @@ -60,15 +60,18 @@ public Builder setRawMetadataVersion(MetadataVersion metadataVersion) { return this; } - public Builder setOrgMetadataVersion(MetadataVersion orgMetadataVersion) { + public void setOrgMetadataVersion(MetadataVersion orgMetadataVersion) { this.orgMetadataVersion = orgMetadataVersion; - return this; } public MetadataVersion metadataVersion() { return metadataVersion; } + public MetadataVersion orgmetadataVersion() { + return orgMetadataVersion; + } + public Builder setLossHandler(Consumer lossHandler) { this.lossHandler = lossHandler; return this; diff --git a/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java b/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java index 07ddd4ff8cbba..759756edbbb31 100644 --- a/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java @@ -21,6 +21,10 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; +import java.io.ByteArrayOutputStream; +import java.io.PrintStream; +import java.util.function.Consumer; + import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertThrows; @@ -44,9 +48,33 @@ public void testSetMetadataVersion() { setMetadataVersion(version); if (i < MetadataVersion.MINIMUM_BOOTSTRAP_VERSION.ordinal()) { assertEquals(MetadataVersion.MINIMUM_KRAFT_VERSION, options.metadataVersion()); + assertEquals(version, options.orgmetadataVersion()); } else { assertEquals(version, options.metadataVersion()); } } } + + @Test + public void testHandleLoss() { + PrintStream originalOut = System.out; + String expectedMessage = "stuff"; + Consumer customLossHandler = e -> System.out.println(e.getMessage()); + + for (int i = MetadataVersion.MINIMUM_KRAFT_VERSION.ordinal(); + i < MetadataVersion.VERSIONS.length; + i++) { + ByteArrayOutputStream outContent = new ByteArrayOutputStream(); + MetadataVersion version = MetadataVersion.VERSIONS[i]; + ImageWriterOptions options = new ImageWriterOptions.Builder() + .setMetadataVersion(version) + .setLossHandler(customLossHandler) + .build(); + System.setOut(new PrintStream(outContent)); + options.handleLoss(expectedMessage); + System.setOut(originalOut); + String formattedMessage = String.format("Metadata has been lost because the following could not be represented in metadata version %s: %s", version, expectedMessage); + assertEquals(formattedMessage, outContent.toString().trim()); + } + } } From 9608077944236270064cf4091e267e5ee9c09e20 Mon Sep 17 00:00:00 2001 From: Owen Leung Date: Fri, 16 Jun 2023 00:09:22 +0800 Subject: [PATCH 4/6] refactor variable names --- .../image/writer/ImageWriterOptions.java | 24 ++++++++----------- .../kafka/image/ImageDowngradeTest.java | 2 +- .../image/writer/ImageWriterOptionsTest.java | 2 +- 3 files changed, 12 insertions(+), 16 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java b/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java index 46d923ada4cb7..7c7d1a3c0e042 100644 --- a/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java +++ b/metadata/src/main/java/org/apache/kafka/image/writer/ImageWriterOptions.java @@ -29,7 +29,7 @@ public final class ImageWriterOptions { public static class Builder { private MetadataVersion metadataVersion; - private MetadataVersion orgMetadataVersion; + private MetadataVersion requestedMetadataVersion; private Consumer lossHandler = e -> { throw e; }; @@ -43,10 +43,10 @@ public Builder(MetadataImage image) { } public Builder setMetadataVersion(MetadataVersion metadataVersion) { + setRequestedMetadataVersion(metadataVersion); if (metadataVersion.isLessThan(MetadataVersion.MINIMUM_BOOTSTRAP_VERSION)) { // When writing an image, all versions less than 3.3-IV0 are treated as 3.0-IV1. // This is because those versions don't support FeatureLevelRecord. - setOrgMetadataVersion(metadataVersion); setRawMetadataVersion(MetadataVersion.MINIMUM_KRAFT_VERSION); } else { setRawMetadataVersion(metadataVersion); @@ -60,16 +60,16 @@ public Builder setRawMetadataVersion(MetadataVersion metadataVersion) { return this; } - public void setOrgMetadataVersion(MetadataVersion orgMetadataVersion) { - this.orgMetadataVersion = orgMetadataVersion; + public void setRequestedMetadataVersion(MetadataVersion orgMetadataVersion) { + this.requestedMetadataVersion = orgMetadataVersion; } public MetadataVersion metadataVersion() { return metadataVersion; } - public MetadataVersion orgmetadataVersion() { - return orgMetadataVersion; + public MetadataVersion requestedMetadataVersion() { + return requestedMetadataVersion; } public Builder setLossHandler(Consumer lossHandler) { @@ -78,12 +78,12 @@ public Builder setLossHandler(Consumer lossHandler) } public ImageWriterOptions build() { - return new ImageWriterOptions(metadataVersion, lossHandler, orgMetadataVersion); + return new ImageWriterOptions(metadataVersion, lossHandler, requestedMetadataVersion); } } private final MetadataVersion metadataVersion; - private final MetadataVersion orgMetadataVersion; + private final MetadataVersion requestedMetadataVersion; private final Consumer lossHandler; private ImageWriterOptions( @@ -93,7 +93,7 @@ private ImageWriterOptions( ) { this.metadataVersion = metadataVersion; this.lossHandler = lossHandler; - this.orgMetadataVersion = orgMetadataVersion; + this.requestedMetadataVersion = orgMetadataVersion; } public MetadataVersion metadataVersion() { @@ -101,11 +101,7 @@ public MetadataVersion metadataVersion() { } public void handleLoss(String loss) { - if (orgMetadataVersion != null) { - lossHandler.accept(new UnwritableMetadataException(orgMetadataVersion, loss)); - } else { - lossHandler.accept(new UnwritableMetadataException(metadataVersion, loss)); - } + lossHandler.accept(new UnwritableMetadataException(requestedMetadataVersion, loss)); } } diff --git a/metadata/src/test/java/org/apache/kafka/image/ImageDowngradeTest.java b/metadata/src/test/java/org/apache/kafka/image/ImageDowngradeTest.java index 86eca51b4c131..b417bebf8acbc 100644 --- a/metadata/src/test/java/org/apache/kafka/image/ImageDowngradeTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/ImageDowngradeTest.java @@ -138,7 +138,7 @@ private static void writeWithExpectedLosses( MetadataImage image = delta.apply(MetadataProvenance.EMPTY); RecordListWriter writer = new RecordListWriter(); image.write(writer, new ImageWriterOptions.Builder(). - setRawMetadataVersion(metadataVersion). + setMetadataVersion(metadataVersion). setLossHandler(lossConsumer). build()); assertEquals(expectedLosses, lossConsumer.losses, "Failed to get expected metadata losses."); diff --git a/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java b/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java index 759756edbbb31..e60bd8b901ae0 100644 --- a/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java @@ -48,7 +48,7 @@ public void testSetMetadataVersion() { setMetadataVersion(version); if (i < MetadataVersion.MINIMUM_BOOTSTRAP_VERSION.ordinal()) { assertEquals(MetadataVersion.MINIMUM_KRAFT_VERSION, options.metadataVersion()); - assertEquals(version, options.orgmetadataVersion()); + assertEquals(version, options.requestedMetadataVersion()); } else { assertEquals(version, options.metadataVersion()); } From 094e53ea2fa764a9005f8df0d3a1517a8c1474e5 Mon Sep 17 00:00:00 2001 From: Owen Leung Date: Mon, 17 Jul 2023 18:49:00 +0800 Subject: [PATCH 5/6] refactor testHandleLoss --- .../kafka/image/writer/ImageWriterOptionsTest.java | 13 ++++--------- 1 file changed, 4 insertions(+), 9 deletions(-) diff --git a/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java b/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java index e60bd8b901ae0..80f5a71e5a331 100644 --- a/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java @@ -21,8 +21,6 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; -import java.io.ByteArrayOutputStream; -import java.io.PrintStream; import java.util.function.Consumer; import static org.junit.jupiter.api.Assertions.assertEquals; @@ -57,24 +55,21 @@ public void testSetMetadataVersion() { @Test public void testHandleLoss() { - PrintStream originalOut = System.out; String expectedMessage = "stuff"; - Consumer customLossHandler = e -> System.out.println(e.getMessage()); for (int i = MetadataVersion.MINIMUM_KRAFT_VERSION.ordinal(); i < MetadataVersion.VERSIONS.length; i++) { - ByteArrayOutputStream outContent = new ByteArrayOutputStream(); MetadataVersion version = MetadataVersion.VERSIONS[i]; + String formattedMessage = String.format("Metadata has been lost because the following could not be represented in metadata version %s: %s", version, expectedMessage); + Consumer customLossHandler = e -> { + assertEquals(formattedMessage, e.getMessage()); + }; ImageWriterOptions options = new ImageWriterOptions.Builder() .setMetadataVersion(version) .setLossHandler(customLossHandler) .build(); - System.setOut(new PrintStream(outContent)); options.handleLoss(expectedMessage); - System.setOut(originalOut); - String formattedMessage = String.format("Metadata has been lost because the following could not be represented in metadata version %s: %s", version, expectedMessage); - assertEquals(formattedMessage, outContent.toString().trim()); } } } From 864cab1adc9e1d4452808e17a15b835fdfd4992f Mon Sep 17 00:00:00 2001 From: Owen Leung Date: Mon, 17 Jul 2023 19:32:34 +0800 Subject: [PATCH 6/6] Removed curly brackets --- .../org/apache/kafka/image/writer/ImageWriterOptionsTest.java | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java b/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java index 80f5a71e5a331..c8653c46d8aec 100644 --- a/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/writer/ImageWriterOptionsTest.java @@ -62,9 +62,7 @@ public void testHandleLoss() { i++) { MetadataVersion version = MetadataVersion.VERSIONS[i]; String formattedMessage = String.format("Metadata has been lost because the following could not be represented in metadata version %s: %s", version, expectedMessage); - Consumer customLossHandler = e -> { - assertEquals(formattedMessage, e.getMessage()); - }; + Consumer customLossHandler = e -> assertEquals(formattedMessage, e.getMessage()); ImageWriterOptions options = new ImageWriterOptions.Builder() .setMetadataVersion(version) .setLossHandler(customLossHandler)