From 787b4fe9550389a2ca6d1131b32aff6e4f8610d1 Mon Sep 17 00:00:00 2001 From: Josep Prat Date: Sun, 13 Jun 2021 13:00:14 +0200 Subject: [PATCH 1/9] MINOR: clean up unneeded `@SuppressWarnings` (#10855) Reviewers: Luke Chen , Matthias J. Sax , Chia-Ping Tsai --- .../java/org/apache/kafka/streams/StreamsConfig.java | 1 - .../internals/CogroupedStreamAggregateBuilder.java | 1 - .../internals/graph/KTableKTableJoinNode.java | 1 - .../kstream/internals/graph/ProcessorParameters.java | 1 - .../kstream/internals/graph/TableProcessorNode.java | 1 - .../processor/internals/ProcessorAdapter.java | 1 - .../processor/internals/ProcessorContextImpl.java | 1 - .../processor/internals/RecordDeserializer.java | 1 - .../streams/state/internals/RecordConverters.java | 1 - .../internals/RocksDBTimeOrderedWindowStore.java | 3 --- .../internals/TimestampedWindowStoreBuilder.java | 3 --- .../WindowToTimestampedWindowByteStoreAdapter.java | 3 --- .../org/apache/kafka/streams/KafkaStreamsTest.java | 1 - .../integration/EosV2UpgradeIntegrationTest.java | 1 - .../StandbyTaskCreationIntegrationTest.java | 1 - .../internals/InternalStreamsBuilderTest.java | 1 - .../kstream/internals/KGroupedTableImplTest.java | 3 --- .../streams/kstream/internals/KTableImplTest.java | 1 - .../internals/SessionWindowedKStreamImplTest.java | 1 - .../internals/SlidingWindowedKStreamImplTest.java | 1 - .../internals/TimeWindowedKStreamImplTest.java | 1 - .../internals/TransformerSupplierAdapterTest.java | 1 - .../SubscriptionResponseWrapperSerdeTest.java | 1 - .../foreignkeyjoin/SubscriptionWrapperSerdeTest.java | 2 -- .../internals/ProcessorContextImplTest.java | 8 ++------ .../internals/CachingPersistentWindowStoreTest.java | 12 ------------ .../state/internals/FilteredCacheIteratorTest.java | 1 - .../MeteredTimestampedKeyValueStoreTest.java | 1 - .../state/internals/RocksDBWindowStoreTest.java | 1 - .../kafka/streams/tests/StreamsUpgradeTest.java | 1 - .../StreamsUpgradeToCooperativeRebalanceTest.java | 1 - .../kafka/test/GenericInMemoryKeyValueStore.java | 1 - .../GenericInMemoryTimestampedKeyValueStore.java | 1 - 33 files changed, 2 insertions(+), 58 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java index 2ef29af8b552f..cb6d58b2c694d 100644 --- a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java +++ b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java @@ -135,7 +135,6 @@ * @see ConsumerConfig * @see ProducerConfig */ -@SuppressWarnings("deprecation") public class StreamsConfig extends AbstractConfig { private static final Logger log = LoggerFactory.getLogger(StreamsConfig.class); diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/CogroupedStreamAggregateBuilder.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/CogroupedStreamAggregateBuilder.java index c7585263a718a..6cb529dfcf904 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/CogroupedStreamAggregateBuilder.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/CogroupedStreamAggregateBuilder.java @@ -48,7 +48,6 @@ class CogroupedStreamAggregateBuilder { CogroupedStreamAggregateBuilder(final InternalStreamsBuilder builder) { this.builder = builder; } - @SuppressWarnings("unchecked") KTable build(final Map, Aggregator> groupPatterns, final Initializer initializer, final NamedInternal named, diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/KTableKTableJoinNode.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/KTableKTableJoinNode.java index 0ca1e35f3b9f9..ac8d82101d0eb 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/KTableKTableJoinNode.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/KTableKTableJoinNode.java @@ -209,7 +209,6 @@ public KTableKTableJoinNodeBuilder withStoreBuilder(final StoreBu return this; } - @SuppressWarnings("unchecked") public KTableKTableJoinNode build() { return new KTableKTableJoinNode<>( nodeName, diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/ProcessorParameters.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/ProcessorParameters.java index 018d2b7dc7dd3..ec2ce48f11e83 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/ProcessorParameters.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/ProcessorParameters.java @@ -60,7 +60,6 @@ public org.apache.kafka.streams.processor.ProcessorSupplier oldProcess return oldProcessorSupplier; } - @SuppressWarnings("unchecked") KTableSource kTableSourceSupplier() { // This cast always works because KTableSource hasn't been converted yet. return oldProcessorSupplier == null diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/TableProcessorNode.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/TableProcessorNode.java index 5254c5757f1b9..f13631ff53b5b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/TableProcessorNode.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/TableProcessorNode.java @@ -57,7 +57,6 @@ public String toString() { "} " + super.toString(); } - @SuppressWarnings("unchecked") @Override public void writeToTopology(final InternalTopologyBuilder topologyBuilder, final Properties props) { final String processorName = processorParameters.processorName(); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorAdapter.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorAdapter.java index f067bbda6472d..687e92f0ddb83 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorAdapter.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorAdapter.java @@ -46,7 +46,6 @@ private ProcessorAdapter(final org.apache.kafka.streams.processor.Processor context) { // It only makes sense to use this adapter internally to Streams, in which case diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java index bd7ece4411513..ce06cb189df27 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java @@ -263,7 +263,6 @@ public void commit() { streamTask.requestCommit(); } - @SuppressWarnings("deprecation") // removing #schedule(final long intervalMs,...) will fix this @Override public Cancellable schedule(final Duration interval, final PunctuationType type, diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/RecordDeserializer.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/RecordDeserializer.java index a965187228a37..b5c821ae4c0dc 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/RecordDeserializer.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/RecordDeserializer.java @@ -50,7 +50,6 @@ class RecordDeserializer { * {@link DeserializationExceptionHandler.DeserializationHandlerResponse#FAIL FAIL} * or throws an exception itself */ - @SuppressWarnings("deprecation") ConsumerRecord deserialize(final ProcessorContext processorContext, final ConsumerRecord rawRecord) { diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/RecordConverters.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/RecordConverters.java index 1f2e5930a211b..ad3c91e8073ff 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/RecordConverters.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/RecordConverters.java @@ -23,7 +23,6 @@ public final class RecordConverters { private static final RecordConverter IDENTITY_INSTANCE = record -> record; - @SuppressWarnings("deprecation") private static final RecordConverter RAW_TO_TIMESTAMED_INSTANCE = record -> { final byte[] rawValue = record.value(); final long timestamp = record.timestamp(); diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBTimeOrderedWindowStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBTimeOrderedWindowStore.java index f8ba8837258f3..37aaa27e2f89c 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBTimeOrderedWindowStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBTimeOrderedWindowStore.java @@ -66,7 +66,6 @@ public byte[] fetch(final Bytes key, final long timestamp) { throw new UnsupportedOperationException(); } - @SuppressWarnings("deprecation") @Override public WindowStoreIterator fetch(final Bytes key, final long timeFrom, final long timeTo) { throw new UnsupportedOperationException(); @@ -77,7 +76,6 @@ public WindowStoreIterator backwardFetch(final Bytes key, final long tim throw new UnsupportedOperationException(); } - @SuppressWarnings("deprecation") // note, this method must be kept if super#fetch(...) is removed @Override public KeyValueIterator, byte[]> fetch(final Bytes keyFrom, final Bytes keyTo, @@ -105,7 +103,6 @@ public KeyValueIterator, byte[]> backwardAll() { throw new UnsupportedOperationException(); } - @SuppressWarnings("deprecation") // note, this method must be kept if super#fetchAll(...) is removed @Override public KeyValueIterator, byte[]> fetchAll(final long timeFrom, final long timeTo) { throw new UnsupportedOperationException(); diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/TimestampedWindowStoreBuilder.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/TimestampedWindowStoreBuilder.java index 417b45b46cc7e..b3727f55c01bb 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/TimestampedWindowStoreBuilder.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/TimestampedWindowStoreBuilder.java @@ -135,7 +135,6 @@ public byte[] fetch(final Bytes key, return wrapped.fetch(key, time); } - @SuppressWarnings("deprecation") @Override public WindowStoreIterator fetch(final Bytes key, final long timeFrom, @@ -150,7 +149,6 @@ public WindowStoreIterator backwardFetch(final Bytes key, return wrapped.backwardFetch(key, timeFrom, timeTo); } - @SuppressWarnings("deprecation") @Override public KeyValueIterator, byte[]> fetch(final Bytes keyFrom, final Bytes keyTo, @@ -167,7 +165,6 @@ public KeyValueIterator, byte[]> backwardFetch(final Bytes keyFr return wrapped.backwardFetch(keyFrom, keyTo, timeFrom, timeTo); } - @SuppressWarnings("deprecation") @Override public KeyValueIterator, byte[]> fetchAll(final long timeFrom, final long timeTo) { diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/WindowToTimestampedWindowByteStoreAdapter.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/WindowToTimestampedWindowByteStoreAdapter.java index 8d895fc7f88e7..f7999d3a449c8 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/WindowToTimestampedWindowByteStoreAdapter.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/WindowToTimestampedWindowByteStoreAdapter.java @@ -54,7 +54,6 @@ public byte[] fetch(final Bytes key, } @Override - @SuppressWarnings("deprecation") public WindowStoreIterator fetch(final Bytes key, final long timeFrom, final long timeTo) { @@ -83,7 +82,6 @@ public WindowStoreIterator backwardFetch(final Bytes key, } @Override - @SuppressWarnings("deprecation") public KeyValueIterator, byte[]> fetch(final Bytes keyFrom, final Bytes keyTo, final long timeFrom, @@ -126,7 +124,6 @@ public KeyValueIterator, byte[]> backwardAll() { } @Override - @SuppressWarnings("deprecation") public KeyValueIterator, byte[]> fetchAll(final long timeFrom, final long timeTo) { return new KeyValueToTimestampedKeyValueIteratorAdapter<>(store.fetchAll(timeFrom, timeTo)); diff --git a/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java b/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java index 300a9e90542a9..54a15cc062023 100644 --- a/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java @@ -1079,7 +1079,6 @@ public void shouldTransitToRunningWithGlobalOnlyTopology() throws InterruptedExc } } - @SuppressWarnings("unchecked") @Deprecated // testing old PAPI private Topology getStatefulTopology(final String inputTopic, final String outputTopic, diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/EosV2UpgradeIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/EosV2UpgradeIntegrationTest.java index 4543b99b0da74..b6aab860eac85 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/EosV2UpgradeIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/EosV2UpgradeIntegrationTest.java @@ -875,7 +875,6 @@ private KafkaStreams getKafkaStreams(final String appDir, final KStream input = builder.stream(MULTI_PARTITION_INPUT_TOPIC); input.transform(new TransformerSupplier>() { - @SuppressWarnings("unchecked") @Override public Transformer> get() { return new Transformer>() { diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/StandbyTaskCreationIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/StandbyTaskCreationIntegrationTest.java index ab22ae6b605d7..f924e08b85efe 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/StandbyTaskCreationIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/StandbyTaskCreationIntegrationTest.java @@ -105,7 +105,6 @@ public void shouldNotCreateAnyStandByTasksForStateStoreWithLoggingDisabled() thr builder.addStateStore(keyValueStoreBuilder); builder.stream(INPUT_TOPIC, Consumed.with(Serdes.Integer(), Serdes.Integer())) .transform(() -> new Transformer>() { - @SuppressWarnings("unchecked") @Override public void init(final ProcessorContext context) {} diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilderTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilderTest.java index 1fffb500a74c4..76ae717a56061 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilderTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilderTest.java @@ -51,7 +51,6 @@ import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; -@SuppressWarnings("unchecked") public class InternalStreamsBuilderTest { private static final String APP_ID = "app-id"; diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedTableImplTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedTableImplTest.java index 130e299ff8835..ee4b1365e821c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedTableImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedTableImplTest.java @@ -201,7 +201,6 @@ public void shouldReduceWithInternalStoreName() { } } - @SuppressWarnings("unchecked") @Test public void shouldReduceAndMaterializeResults() { final KeyValueMapper> intProjection = @@ -235,7 +234,6 @@ public void shouldReduceAndMaterializeResults() { } } - @SuppressWarnings("unchecked") @Test public void shouldCountAndMaterializeResults() { builder @@ -265,7 +263,6 @@ public void shouldCountAndMaterializeResults() { } } - @SuppressWarnings("unchecked") @Test public void shouldAggregateAndMaterializeResults() { builder diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableImplTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableImplTest.java index b979d3ecd7734..1885d57ab5b79 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableImplTest.java @@ -574,7 +574,6 @@ public void shouldThrowNullPointerOnTransformValuesWithKeyWhenMaterializedIsNull assertThrows(NullPointerException.class, () -> table.transformValues(valueTransformerSupplier, (Materialized) null)); } - @SuppressWarnings("unchecked") @Test public void shouldThrowNullPointerOnTransformValuesWithKeyWhenStoreNamesNull() { final ValueTransformerWithKeySupplier valueTransformerSupplier = diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/SessionWindowedKStreamImplTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/SessionWindowedKStreamImplTest.java index d6e56ba5b6128..abca688d2d54c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/SessionWindowedKStreamImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/SessionWindowedKStreamImplTest.java @@ -283,7 +283,6 @@ public void shouldThrowNullPointerOnMaterializedReduceIfMaterializedIsNull() { } @Test - @SuppressWarnings("unchecked") public void shouldThrowNullPointerOnMaterializedReduceIfNamedIsNull() { assertThrows(NullPointerException.class, () -> stream.reduce(MockReducer.STRING_ADDER, (Named) null)); } diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/SlidingWindowedKStreamImplTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/SlidingWindowedKStreamImplTest.java index d6b26bf2a7c59..f012cedcf4391 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/SlidingWindowedKStreamImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/SlidingWindowedKStreamImplTest.java @@ -361,7 +361,6 @@ public void shouldThrowNullPointerOnMaterializedReduceIfMaterializedIsNull() { } @Test - @SuppressWarnings("unchecked") public void shouldThrowNullPointerOnMaterializedReduceIfNamedIsNull() { assertThrows(NullPointerException.class, () -> windowedStream.reduce(MockReducer.STRING_ADDER, (Named) null)); } diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/TimeWindowedKStreamImplTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/TimeWindowedKStreamImplTest.java index c35da00697ba5..38fda9d6da21c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/TimeWindowedKStreamImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/TimeWindowedKStreamImplTest.java @@ -300,7 +300,6 @@ public void shouldThrowNullPointerOnMaterializedReduceIfMaterializedIsNull() { } @Test - @SuppressWarnings("unchecked") public void shouldThrowNullPointerOnMaterializedReduceIfNamedIsNull() { assertThrows(NullPointerException.class, () -> windowedStream.reduce( MockReducer.STRING_ADDER, diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/TransformerSupplierAdapterTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/TransformerSupplierAdapterTest.java index 115855d964582..1eb55d0a08945 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/TransformerSupplierAdapterTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/TransformerSupplierAdapterTest.java @@ -34,7 +34,6 @@ import static org.hamcrest.core.IsNot.not; import static org.hamcrest.MatcherAssert.assertThat; -@SuppressWarnings("unchecked") public class TransformerSupplierAdapterTest extends EasyMockSupport { private ProcessorContext context; diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResponseWrapperSerdeTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResponseWrapperSerdeTest.java index 1bd1bd27cf7cc..30fc0c318519c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResponseWrapperSerdeTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResponseWrapperSerdeTest.java @@ -127,7 +127,6 @@ public void shouldSerdeWithNullsTest() { } @Test - @SuppressWarnings("unchecked") public void shouldThrowExceptionWithBadVersionTest() { final long[] hashedValue = null; assertThrows(UnsupportedVersionException.class, diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerdeTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerdeTest.java index b7ce34f0d1030..e937efe2bc092 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerdeTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerdeTest.java @@ -59,7 +59,6 @@ public void shouldSerdeNullHashTest() { } @Test - @SuppressWarnings("unchecked") public void shouldThrowExceptionOnNullKeyTest() { final String originalKey = null; final long[] hashedValue = Murmur3.hash128(new byte[] {(byte) 0xFF, (byte) 0xAA, (byte) 0x00, (byte) 0x19}); @@ -68,7 +67,6 @@ public void shouldThrowExceptionOnNullKeyTest() { } @Test - @SuppressWarnings("unchecked") public void shouldThrowExceptionOnNullInstructionTest() { final String originalKey = "originalKey"; final long[] hashedValue = Murmur3.hash128(new byte[] {(byte) 0xFF, (byte) 0xAA, (byte) 0x00, (byte) 0x19}); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorContextImplTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorContextImplTest.java index 9412e55e0ff60..13816d740a62b 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorContextImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorContextImplTest.java @@ -220,7 +220,6 @@ public void globalTimestampedKeyValueStoreShouldBeReadOnly() { } @Test - @SuppressWarnings("deprecation") public void globalWindowStoreShouldBeReadOnly() { doTest("GlobalWindowStore", (Consumer>) store -> { verifyStoreCannotBeInitializedOrClosed(store); @@ -238,7 +237,6 @@ public void globalWindowStoreShouldBeReadOnly() { @Test - @SuppressWarnings("deprecation") public void globalTimestampedWindowStoreShouldBeReadOnly() { doTest("GlobalTimestampedWindowStore", (Consumer>) store -> { verifyStoreCannotBeInitializedOrClosed(store); @@ -325,7 +323,6 @@ public void localTimestampedKeyValueStoreShouldNotAllowInitOrClose() { } @Test - @SuppressWarnings("deprecation") public void localWindowStoreShouldNotAllowInitOrClose() { doTest("LocalWindowStore", (Consumer>) store -> { verifyStoreCannotBeInitializedOrClosed(store); @@ -345,7 +342,6 @@ public void localWindowStoreShouldNotAllowInitOrClose() { } @Test - @SuppressWarnings("deprecation") public void localTimestampedWindowStoreShouldNotAllowInitOrClose() { doTest("LocalTimestampedWindowStore", (Consumer>) store -> { verifyStoreCannotBeInitializedOrClosed(store); @@ -615,7 +611,7 @@ private TimestampedKeyValueStore timestampedKeyValueStoreMock() { return timestampedKeyValueStoreMock; } - @SuppressWarnings({"unchecked", "deprecation"}) + @SuppressWarnings("unchecked") private WindowStore windowStoreMock() { final WindowStore windowStore = mock(WindowStore.class); @@ -638,7 +634,7 @@ private WindowStore windowStoreMock() { return windowStore; } - @SuppressWarnings({"unchecked", "deprecation"}) + @SuppressWarnings("unchecked") private TimestampedWindowStore timestampedWindowStoreMock() { final TimestampedWindowStore windowStore = mock(TimestampedWindowStore.class); diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/CachingPersistentWindowStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/CachingPersistentWindowStoreTest.java index 8bdf8b705f3ce..7c316f4ceb585 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/CachingPersistentWindowStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/CachingPersistentWindowStoreTest.java @@ -236,7 +236,6 @@ public void close() { } @Test - @SuppressWarnings("deprecation") public void shouldPutFetchFromCache() { cachingStore.put(bytesKey("a"), bytesValue("a"), DEFAULT_TIMESTAMP); cachingStore.put(bytesKey("b"), bytesValue("b"), DEFAULT_TIMESTAMP); @@ -275,7 +274,6 @@ private String stringFrom(final byte[] from) { } @Test - @SuppressWarnings("deprecation") public void shouldPutFetchRangeFromCache() { cachingStore.put(bytesKey("a"), bytesValue("a"), DEFAULT_TIMESTAMP); cachingStore.put(bytesKey("b"), bytesValue("b"), DEFAULT_TIMESTAMP); @@ -317,7 +315,6 @@ public void shouldGetAllFromCache() { } @Test - @SuppressWarnings("deprecation") public void shouldGetAllBackwardFromCache() { cachingStore.put(bytesKey("a"), bytesValue("a"), DEFAULT_TIMESTAMP); cachingStore.put(bytesKey("b"), bytesValue("b"), DEFAULT_TIMESTAMP); @@ -340,7 +337,6 @@ public void shouldGetAllBackwardFromCache() { } @Test - @SuppressWarnings("deprecation") public void shouldFetchAllWithinTimestampRange() { final String[] array = {"a", "b", "c", "d", "e", "f", "g", "h"}; for (int i = 0; i < array.length; i++) { @@ -382,7 +378,6 @@ public void shouldFetchAllWithinTimestampRange() { } @Test - @SuppressWarnings("deprecation") public void shouldFetchAllBackwardWithinTimestampRange() { final String[] array = {"a", "b", "c", "d", "e", "f", "g", "h"}; for (int i = 0; i < array.length; i++) { @@ -439,7 +434,6 @@ public void shouldFlushEvictedItemsIntoUnderlyingStore() { } @Test - @SuppressWarnings("deprecation") public void shouldForwardDirtyItemsWhenFlushCalled() { final Windowed windowedKey = new Windowed<>("1", new TimeWindow(DEFAULT_TIMESTAMP, DEFAULT_TIMESTAMP + WINDOW_SIZE)); @@ -456,7 +450,6 @@ public void shouldSetFlushListener() { } @Test - @SuppressWarnings("deprecation") public void shouldForwardOldValuesWhenEnabled() { cachingStore.setFlushListener(cacheListener, true); final Windowed windowedKey = @@ -485,7 +478,6 @@ public void shouldForwardOldValuesWhenEnabled() { } @Test - @SuppressWarnings("deprecation") public void shouldForwardOldValuesWhenDisabled() { final Windowed windowedKey = new Windowed<>("1", new TimeWindow(DEFAULT_TIMESTAMP, DEFAULT_TIMESTAMP + WINDOW_SIZE)); @@ -616,7 +608,6 @@ public void shouldIterateBackwardCacheAndStoreKeyRange() { } @Test - @SuppressWarnings("deprecation") public void shouldClearNamespaceCacheOnClose() { cachingStore.put(bytesKey("a"), bytesValue("a"), 0L); assertEquals(1, cache.size()); @@ -637,7 +628,6 @@ public void shouldThrowIfTryingToFetchRangeFromClosedCachingStore() { } @Test - @SuppressWarnings("deprecation") public void shouldThrowIfTryingToWriteToClosedCachingStore() { cachingStore.close(); assertThrows(InvalidStateStoreException.class, () -> cachingStore.put(bytesKey("a"), bytesValue("a"), 0L)); @@ -786,13 +776,11 @@ public void shouldReturnSameResultsForSingleKeyFetchAndEqualKeyRangeBackwardFetc } @Test - @SuppressWarnings("deprecation") public void shouldThrowNullPointerExceptionOnPutNullKey() { assertThrows(NullPointerException.class, () -> cachingStore.put(null, bytesValue("anyValue"), 0L)); } @Test - @SuppressWarnings("deprecation") public void shouldNotThrowNullPointerExceptionOnPutNullValue() { cachingStore.put(bytesKey("a"), null, 0L); } diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/FilteredCacheIteratorTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/FilteredCacheIteratorTest.java index f49d881f94236..bd794333660fc 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/FilteredCacheIteratorTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/FilteredCacheIteratorTest.java @@ -50,7 +50,6 @@ public Bytes cacheKey(final Bytes key) { } }; - @SuppressWarnings("unchecked") private final KeyValueStore store = new GenericInMemoryKeyValueStore<>("my-store"); private final KeyValue firstEntry = KeyValue.pair(Bytes.wrap("a".getBytes()), new LRUCacheEntry("1".getBytes())); diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampedKeyValueStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampedKeyValueStoreTest.java index eb48e0c5fc2fa..bce87ef2fbeec 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampedKeyValueStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampedKeyValueStoreTest.java @@ -426,7 +426,6 @@ private KafkaMetric metric(final MetricName metricName) { } @Test - @SuppressWarnings("unchecked") public void shouldNotThrowExceptionIfSerdesCorrectlySetFromProcessorContext() { final MeteredTimestampedKeyValueStore store = new MeteredTimestampedKeyValueStore<>( inner, diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/RocksDBWindowStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/RocksDBWindowStoreTest.java index 2b890f172d500..3bb9f14e66827 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/RocksDBWindowStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/RocksDBWindowStoreTest.java @@ -43,7 +43,6 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; -@SuppressWarnings("PointlessArithmeticExpression") public class RocksDBWindowStoreTest extends AbstractWindowBytesStoreTest { private static final String STORE_NAME = "rocksDB window store"; diff --git a/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java index 2ad07f2fa3e4a..2fabf97925ea6 100644 --- a/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java @@ -65,7 +65,6 @@ public class StreamsUpgradeTest { - @SuppressWarnings("unchecked") public static void main(final String[] args) throws Exception { if (args.length < 1) { System.err.println("StreamsUpgradeTest requires one argument (properties-file) but no provided: "); diff --git a/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeToCooperativeRebalanceTest.java b/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeToCooperativeRebalanceTest.java index bd5752aa30d86..19e81acd8a2f1 100644 --- a/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeToCooperativeRebalanceTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeToCooperativeRebalanceTest.java @@ -37,7 +37,6 @@ public class StreamsUpgradeToCooperativeRebalanceTest { - @SuppressWarnings("unchecked") public static void main(final String[] args) throws Exception { if (args.length < 1) { System.err.println("StreamsUpgradeToCooperativeRebalanceTest requires one argument (properties-file) but no args provided"); diff --git a/streams/src/test/java/org/apache/kafka/test/GenericInMemoryKeyValueStore.java b/streams/src/test/java/org/apache/kafka/test/GenericInMemoryKeyValueStore.java index 72e6c266fea8c..d9a2afe0f5884 100644 --- a/streams/src/test/java/org/apache/kafka/test/GenericInMemoryKeyValueStore.java +++ b/streams/src/test/java/org/apache/kafka/test/GenericInMemoryKeyValueStore.java @@ -60,7 +60,6 @@ public String name() { @Deprecated @Override - @SuppressWarnings("unchecked") /* This is a "dummy" store used for testing; it does not support restoring from changelog since we allow it to be serde-ignorant */ public void init(final ProcessorContext context, final StateStore root) { diff --git a/streams/src/test/java/org/apache/kafka/test/GenericInMemoryTimestampedKeyValueStore.java b/streams/src/test/java/org/apache/kafka/test/GenericInMemoryTimestampedKeyValueStore.java index c77cbacb1db24..2198d181dd8d3 100644 --- a/streams/src/test/java/org/apache/kafka/test/GenericInMemoryTimestampedKeyValueStore.java +++ b/streams/src/test/java/org/apache/kafka/test/GenericInMemoryTimestampedKeyValueStore.java @@ -61,7 +61,6 @@ public String name() { @Deprecated @Override - @SuppressWarnings("unchecked") /* This is a "dummy" store used for testing; it does not support restoring from changelog since we allow it to be serde-ignorant */ public void init(final ProcessorContext context, final StateStore root) { From 530224e4fe853df734e7bf7c15661fe6b1ab34fe Mon Sep 17 00:00:00 2001 From: Ismael Juma Date: Sun, 13 Jun 2021 08:14:37 -0700 Subject: [PATCH 2/9] KAFKA-12940: Enable JDK 16 builds in Jenkins (#10702) JDK 15 no longer receives updates, so we want to switch from JDK 15 to JDK 16. However, we have a number of tests that don't yet pass with JDK 16. Instead of replacing JDK 15 with JDK 16, we have both for now and we either disable (via annotations) or exclude (via gradle) the tests that don't pass with JDK 16 yet. The annotations approach is better, but it doesn't work for tests that rely on the PowerMock JUnit 4 runner. Also add `--illegal-access=permit` when building with JDK 16 to make MiniKdc work for now. This has been removed in JDK 17, so we'll have to figure out another solution when we migrate to that. Relevant JIRAs for the disabled tests: KAFKA-12790, KAFKA-12941, KAFKA-12942. Moved some assertions from `testTlsDefaults` to `testUnsupportedTlsVersion` since the former claims to test the success case while the former tests the failure case. Reviewers: Chia-Ping Tsai --- Jenkinsfile | 26 ++++++++++++-- build.gradle | 35 ++++++++++++++----- .../common/network/SslTransportLayerTest.java | 17 +++++---- 3 files changed, 57 insertions(+), 21 deletions(-) diff --git a/Jenkinsfile b/Jenkinsfile index a966bf26d94b2..85dd0e2f8bb92 100644 --- a/Jenkinsfile +++ b/Jenkinsfile @@ -142,6 +142,7 @@ pipeline { } } + // Remove this when all tests pass with JDK 16 stage('JDK 15 and Scala 2.13') { agent { label 'ubuntu' } tools { @@ -161,6 +162,25 @@ pipeline { } } + stage('JDK 16 and Scala 2.13') { + agent { label 'ubuntu' } + tools { + jdk 'jdk_16_latest' + } + options { + timeout(time: 8, unit: 'HOURS') + timestamps() + } + environment { + SCALA_VERSION=2.13 + } + steps { + doValidation() + doTest(env) + echo 'Skipping Kafka Streams archetype test for Java 16' + } + } + stage('ARM') { agent { label 'arm4' } options { @@ -231,14 +251,14 @@ pipeline { } } - stage('JDK 15 and Scala 2.12') { + stage('JDK 16 and Scala 2.12') { when { not { changeRequest() } beforeAgent true } agent { label 'ubuntu' } tools { - jdk 'jdk_15_latest' + jdk 'jdk_16_latest' } options { timeout(time: 8, unit: 'HOURS') @@ -250,7 +270,7 @@ pipeline { steps { doValidation() doTest(env) - echo 'Skipping Kafka Streams archetype test for Java 15' + echo 'Skipping Kafka Streams archetype test for Java 16' } } } diff --git a/build.gradle b/build.gradle index b25914bb4b335..e579e7ac18eb0 100644 --- a/build.gradle +++ b/build.gradle @@ -102,6 +102,8 @@ ext { defaultMaxHeapSize = "2g" defaultJvmArgs = ["-Xss4m", "-XX:+UseParallelGC"] + if (JavaVersion.current() == JavaVersion.VERSION_16) + defaultJvmArgs.add("--illegal-access=permit") userMaxForks = project.hasProperty('maxParallelForks') ? maxParallelForks.toInteger() : null userIgnoreFailures = project.hasProperty('ignoreFailures') ? ignoreFailures : false @@ -353,6 +355,27 @@ subprojects { } } + // The suites are for running sets of tests in IDEs. + // Gradle will run each test class, so we exclude the suites to avoid redundantly running the tests twice. + def testsToExclude = ['**/*Suite.class'] + // Exclude PowerMock tests when running with Java 16 until a version of PowerMock that supports Java 16 is released + // The relevant issues are https://github.com/powermock/powermock/issues/1094 and https://github.com/powermock/powermock/issues/1099 + if (JavaVersion.current().isCompatibleWith(JavaVersion.VERSION_16)) { + testsToExclude.addAll([ + // connect tests + "**/AbstractHerderTest.*", "**/ConnectClusterStateImplTest.*", "**/ConnectorPluginsResourceTest.*", + "**/ConnectorsResourceTest.*", "**/DistributedHerderTest.*", "**/FileOffsetBakingStoreTest.*", + "**/ErrorHandlingTaskTest.*", "**/KafkaConfigBackingStoreTest.*", "**/KafkaOffsetBackingStoreTest.*", + "**/KafkaBasedLogTest.*", "**/OffsetStorageWriterTest.*", "**/StandaloneHerderTest.*", + "**/SourceTaskOffsetCommitterTest.*", "**/WorkerConfigTransformerTest.*", "**/WorkerGroupMemberTest.*", + "**/WorkerSinkTaskTest.*", "**/WorkerSinkTaskThreadedTest.*", "**/WorkerSourceTaskTest.*", + "**/WorkerTaskTest.*", "**/WorkerTest.*", "**/RestServerTest.*", + // streams tests + "**/KafkaStreamsTest.*", "**/RepartitionTopicsTest.*", "**/RocksDBMetricsRecorderTest.*", + "**/StreamsMetricsImplTest.*", "**/StateManagerUtilTest.*", "**/TableSourceNodeTest.*" + ]) + } + test { maxParallelForks = userMaxForks ?: Runtime.runtime.availableProcessors() ignoreFailures = userIgnoreFailures @@ -367,9 +390,7 @@ subprojects { } logTestStdout.rehydrate(delegate, owner, this)() - // The suites are for running sets of tests in IDEs. - // Gradle will run each test class, so we exclude the suites to avoid redundantly running the tests twice. - exclude '**/*Suite.class' + exclude testsToExclude if (shouldUseJUnit5) useJUnitPlatform() @@ -395,9 +416,7 @@ subprojects { } logTestStdout.rehydrate(delegate, owner, this)() - // The suites are for running sets of tests in IDEs. - // Gradle will run each test class, so we exclude the suites to avoid redundantly running the tests twice. - exclude '**/*Suite.class' + exclude testsToExclude if (shouldUseJUnit5) { useJUnitPlatform { @@ -429,9 +448,7 @@ subprojects { } logTestStdout.rehydrate(delegate, owner, this)() - // The suites are for running sets of tests in IDEs. - // Gradle will run each test class, so we exclude the suites to avoid redundantly running the tests twice. - exclude '**/*Suite.class' + exclude testsToExclude if (shouldUseJUnit5) { useJUnitPlatform { diff --git a/clients/src/test/java/org/apache/kafka/common/network/SslTransportLayerTest.java b/clients/src/test/java/org/apache/kafka/common/network/SslTransportLayerTest.java index 44187134225fc..f9a64f6fdb6fc 100644 --- a/clients/src/test/java/org/apache/kafka/common/network/SslTransportLayerTest.java +++ b/clients/src/test/java/org/apache/kafka/common/network/SslTransportLayerTest.java @@ -36,6 +36,8 @@ import org.apache.kafka.test.TestSslUtils; import org.apache.kafka.test.TestUtils; import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.condition.DisabledOnJre; +import org.junit.jupiter.api.condition.JRE; import org.junit.jupiter.api.extension.ExtensionContext; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; @@ -591,7 +593,7 @@ public void testInvalidKeyPassword(Args args) throws Exception { } /** - * Tests that connection success with the default TLS version. + * Tests that connection succeeds with the default TLS version. */ @ParameterizedTest @ArgumentsSource(SslTransportLayerArgumentsProvider.class) @@ -611,12 +613,6 @@ public void testTlsDefaults(Args args) throws Exception { NetworkTestUtils.checkClientConnection(selector, "0", 10, 100); server.verifyAuthenticationMetrics(1, 0); selector.close(); - - checkAuthenticationFailed(args, "1", "TLSv1.1"); - server.verifyAuthenticationMetrics(1, 1); - - checkAuthenticationFailed(args, "2", "TLSv1"); - server.verifyAuthenticationMetrics(1, 2); } /** Checks connection failed using the specified {@code tlsVersion}. */ @@ -636,12 +632,15 @@ private void checkAuthenticationFailed(Args args, String node, String tlsVersion */ @ParameterizedTest @ArgumentsSource(SslTransportLayerArgumentsProvider.class) - public void testUnsupportedTLSVersion(Args args) throws Exception { - args.sslServerConfigs.put(SslConfigs.SSL_ENABLED_PROTOCOLS_CONFIG, Arrays.asList("TLSv1.2")); + @DisabledOnJre(JRE.JAVA_16) + public void testUnsupportedTlsVersion(Args args) throws Exception { server = createEchoServer(args, SecurityProtocol.SSL); checkAuthenticationFailed(args, "0", "TLSv1.1"); server.verifyAuthenticationMetrics(0, 1); + + checkAuthenticationFailed(args, "0", "TLSv1"); + server.verifyAuthenticationMetrics(0, 2); } /** From 39b9df50909ecd92c4e1427e476ca5996e46df9a Mon Sep 17 00:00:00 2001 From: David Christle Date: Sun, 13 Jun 2021 11:14:24 -0500 Subject: [PATCH 3/9] KAFKA-12921: Upgrade zstd-jni to 1.5.0-2 (#10847) This PR aims to upgrade `zstd-jni` from `1.4.9-1` to `1.5.0-2`. This change will incorporate a number of bug fixes and performance improvements made in `1.5.0` of `zstd`: - https://github.com/facebook/zstd/releases/tag/v1.5.0 - https://github.com/luben/zstd-jni/releases/tag/v1.5.0-1 - https://github.com/luben/zstd-jni/releases/tag/v1.5.0-2 The most recent `1.5.0` release offers +25%-140% (compression) and +15% (decompression) performance improvements under certain conditions. Those conditions are unlikely to apply to Kafka with the default configuration, however. Since this is a dependency change, this should pass all the existing CIs. Reviewers: Lee Dongjin , Ismael Juma --- gradle/dependencies.gradle | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/gradle/dependencies.gradle b/gradle/dependencies.gradle index 4315677955db0..069d3f6235365 100644 --- a/gradle/dependencies.gradle +++ b/gradle/dependencies.gradle @@ -112,7 +112,7 @@ versions += [ spotbugs: "4.2.2", zinc: "1.3.5", zookeeper: "3.5.9", - zstd: "1.4.9-1" + zstd: "1.5.0-2" ] libs += [ activation: "javax.activation:activation:$versions.activation", From 4724083a321a034ce03f708507dd81b668c0a344 Mon Sep 17 00:00:00 2001 From: Luke Chen Date: Mon, 14 Jun 2021 00:49:05 +0800 Subject: [PATCH 4/9] KAFKA-8940: decrease session timeout to make test faster and reliable (#10871) While there might still be some issue about the test as described here by @ableegoldman , but I found the reason why this test failed quite frequently recently. It's because we increased the session timeout to 45 sec in KIP-735. The reason why increasing session timeout affected this test is because in this test, we will keep adding new stream clients and remove old one, to maintain only 3 stream clients alive. The problem here is, when old stream closed, we won't trigger rebalance immediately due to the stream clients are all static members as described in KIP-345, which means, we will trigger trigger group rebalance only when session.timeout expired. That said, when old client closed, we'll have at least 45 sec with some tasks not working. Also, in this test, we have 2 timeout conditions to fail this test before verification passed: 1. 6 minutes timeout 2. polling 30 times (each with 5 seconds) without getting any data. (that is, 5 * 30 = 150 sec without consuming any data) For (1), in my test under 45 session timeout, we'll create 8 stream clients, which means, we'll have 5 clients got closed. And each closed client need 45 sec to trigger rebalance, so we'll have 45 * 5 = 225 sec (~4 mins) of the time having some tasks not working. For (2), during new client created and old client closed, it need some time to do rebalance. With 45 session timeout, we only got ~100 sec left. In slow jenkins env, it might reach the 30 retries without getting any data timeout. Therefore, decreasing session timeout can make this test completes faster and more reliable. Reviewers: Guozhang Wang --- .../integration/SmokeTestDriverIntegrationTest.java | 7 +++++++ .../org/apache/kafka/streams/tests/SmokeTestClient.java | 7 ++++--- 2 files changed, 11 insertions(+), 3 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/SmokeTestDriverIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/SmokeTestDriverIntegrationTest.java index 6c20fb311741d..22d773595aacd 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/SmokeTestDriverIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/SmokeTestDriverIntegrationTest.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.streams.integration; +import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.common.utils.Exit; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster; @@ -89,6 +90,10 @@ SmokeTestDriver.VerificationResult result() { } + // In this test, we try to keep creating new stream, and closing the old one, to maintain only 3 streams alive. + // During the new stream added and old stream left, the stream process should still complete without issue. + // We set 2 timeout condition to fail the test before passing the verification: + // (1) 6 min timeout, (2) 30 tries of polling without getting any data @Test public void shouldWorkWithRebalance() throws InterruptedException { Exit.setExitProcedure((statusCode, message) -> { @@ -110,6 +115,8 @@ public void shouldWorkWithRebalance() throws InterruptedException { final Properties props = new Properties(); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); + // decrease the session timeout so that we can trigger the rebalance soon after old client left closed + props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 10000); // cycle out Streams instances as long as the test is running. while (driver.isAlive()) { diff --git a/streams/src/test/java/org/apache/kafka/streams/tests/SmokeTestClient.java b/streams/src/test/java/org/apache/kafka/streams/tests/SmokeTestClient.java index 2e53d580d5292..e39b3c0a741b0 100644 --- a/streams/src/test/java/org/apache/kafka/streams/tests/SmokeTestClient.java +++ b/streams/src/test/java/org/apache/kafka/streams/tests/SmokeTestClient.java @@ -126,13 +126,14 @@ public void start(final Properties streamsProperties) { try { if (!countDownLatch.await(1, TimeUnit.MINUTES)) { System.out.println(name + ": SMOKE-TEST-CLIENT-EXCEPTION: Didn't start in one minute"); + } else { + System.out.println(name + ": SMOKE-TEST-CLIENT-STARTED"); + System.out.println(name + " started at " + Instant.now()); } } catch (final InterruptedException e) { System.out.println(name + ": SMOKE-TEST-CLIENT-EXCEPTION: " + e); e.printStackTrace(System.out); } - System.out.println(name + ": SMOKE-TEST-CLIENT-STARTED"); - System.out.println(name + " started at " + Instant.now()); } public void closeAsync() { @@ -145,7 +146,7 @@ public void close() { if (wasClosed && !uncaughtException) { System.out.println(name + ": SMOKE-TEST-CLIENT-CLOSED"); } else if (wasClosed) { - System.out.println(name + ": SMOKE-TEST-CLIENT-EXCEPTION"); + System.out.println(name + ": SMOKE-TEST-CLIENT-EXCEPTION: Got an uncaught exception"); } else { System.out.println(name + ": SMOKE-TEST-CLIENT-EXCEPTION: Didn't close in time."); } From 987391958dd201c5f033a88a83300a0a0db47834 Mon Sep 17 00:00:00 2001 From: John Roesler Date: Sun, 13 Jun 2021 21:35:02 -0500 Subject: [PATCH 5/9] MINOR: enable EOS during smoke test IT (#10870) This IT has been failing on trunk recently. Enabling EOS during the integration test makes it easier to be sure that the test's assumptions are really true during verification and should make the test more reliable. I also noticed that in the actual system test file, we are using the deprecated property name "beta" instead of "v2". Reviewers: Boyang Chen --- .../java/org/apache/kafka/streams/tests/SmokeTestClient.java | 1 + tests/kafkatest/tests/streams/streams_smoke_test.py | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/tests/SmokeTestClient.java b/streams/src/test/java/org/apache/kafka/streams/tests/SmokeTestClient.java index e39b3c0a741b0..8bd746d1ed423 100644 --- a/streams/src/test/java/org/apache/kafka/streams/tests/SmokeTestClient.java +++ b/streams/src/test/java/org/apache/kafka/streams/tests/SmokeTestClient.java @@ -157,6 +157,7 @@ private Properties getStreamsConfig(final Properties props) { fullProps.put(StreamsConfig.APPLICATION_ID_CONFIG, "SmokeTest"); fullProps.put(StreamsConfig.CLIENT_ID_CONFIG, "SmokeTest-" + name); fullProps.put(StreamsConfig.STATE_DIR_CONFIG, tempDirectory().getAbsolutePath()); + fullProps.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2); fullProps.putAll(props); return fullProps; } diff --git a/tests/kafkatest/tests/streams/streams_smoke_test.py b/tests/kafkatest/tests/streams/streams_smoke_test.py index 29955cd0ca38f..6cd2b7611a30e 100644 --- a/tests/kafkatest/tests/streams/streams_smoke_test.py +++ b/tests/kafkatest/tests/streams/streams_smoke_test.py @@ -48,7 +48,7 @@ def __init__(self, test_context): @cluster(num_nodes=8) @matrix(processing_guarantee=['at_least_once'], crash=[True, False], metadata_quorum=quorum.all_non_upgrade) - @matrix(processing_guarantee=['exactly_once', 'exactly_once_beta'], crash=[True, False]) + @matrix(processing_guarantee=['exactly_once', 'exactly_once_v2'], crash=[True, False]) def test_streams(self, processing_guarantee, crash, metadata_quorum=quorum.zk): processor1 = StreamsSmokeTestJobRunnerService(self.test_context, self.kafka, processing_guarantee) processor2 = StreamsSmokeTestJobRunnerService(self.test_context, self.kafka, processing_guarantee) From 3944f9664783e0cfbf5a00d49ac91192404df35b Mon Sep 17 00:00:00 2001 From: YiDing-Duke Date: Mon, 14 Jun 2021 23:11:19 -0700 Subject: [PATCH 6/9] MINOR: Log formatting for exceptions during configuration related operations (#10843) Format configuration logging during exceptions or errors. Also make sure it redacts sensitive information or unknown values. Reviewers: Luke Chen , David Jacot --- .../java/org/apache/kafka/clients/admin/ConfigEntry.java | 6 +++++- .../main/scala/kafka/server/DynamicBrokerConfig.scala | 9 +++++---- core/src/main/scala/kafka/server/ZkAdminManager.scala | 6 +++++- 3 files changed, 15 insertions(+), 6 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ConfigEntry.java b/clients/src/main/java/org/apache/kafka/clients/admin/ConfigEntry.java index e3426662e801d..30686c93eaeef 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ConfigEntry.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ConfigEntry.java @@ -174,11 +174,15 @@ public int hashCode() { return result; } + /** + * Override toString to redact sensitive value. + * WARNING, user should be responsible to set the correct "isSensitive" field for each config entry. + */ @Override public String toString() { return "ConfigEntry(" + "name=" + name + - ", value=" + value + + ", value=" + (isSensitive ? "Redacted" : value) + ", source=" + source + ", isSensitive=" + isSensitive + ", isReadOnly=" + isReadOnly + diff --git a/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala index 2cf24c87ceba9..9eefdd3933d3c 100755 --- a/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala +++ b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala @@ -295,7 +295,7 @@ class DynamicBrokerConfig(private val kafkaConfig: KafkaConfig) extends Logging dynamicBrokerConfigs ++= props.asScala updateCurrentConfig() } catch { - case e: Exception => error(s"Per-broker configs of $brokerId could not be applied: $persistentProps", e) + case e: Exception => error(s"Per-broker configs of $brokerId could not be applied: ${persistentProps.keys()}", e) } } @@ -306,7 +306,7 @@ class DynamicBrokerConfig(private val kafkaConfig: KafkaConfig) extends Logging dynamicDefaultConfigs ++= props.asScala updateCurrentConfig() } catch { - case e: Exception => error(s"Cluster default configs could not be applied: $persistentProps", e) + case e: Exception => error(s"Cluster default configs could not be applied: ${persistentProps.keys()}", e) } } @@ -469,7 +469,7 @@ class DynamicBrokerConfig(private val kafkaConfig: KafkaConfig) extends Logging } invalidProps.keys.foreach(props.remove) val configSource = if (perBrokerConfig) "broker" else "default cluster" - error(s"Dynamic $configSource config contains invalid values: $invalidProps, these configs will be ignored", e) + error(s"Dynamic $configSource config contains invalid values in: ${invalidProps.keys}, these configs will be ignored", e) } } @@ -555,7 +555,8 @@ class DynamicBrokerConfig(private val kafkaConfig: KafkaConfig) extends Logging } catch { case e: Exception => if (!validateOnly) - error(s"Failed to update broker configuration with configs : ${newConfig.originalsFromThisConfig}", e) + error(s"Failed to update broker configuration with configs : " + + s"${ConfigUtils.configMapToRedactedString(newConfig.originalsFromThisConfig, KafkaConfig.configDef)}", e) throw new ConfigException("Invalid dynamic configuration", e) } } diff --git a/core/src/main/scala/kafka/server/ZkAdminManager.scala b/core/src/main/scala/kafka/server/ZkAdminManager.scala index 87f522fb10d26..7be5ab08b374b 100644 --- a/core/src/main/scala/kafka/server/ZkAdminManager.scala +++ b/core/src/main/scala/kafka/server/ZkAdminManager.scala @@ -415,8 +415,12 @@ class ZkAdminManager(val config: KafkaConfig, info(message) resource -> ApiError.fromThrowable(new InvalidRequestException(message, e)) case e: Throwable => + val configProps = new Properties + config.entries.asScala.filter(_.value != null).foreach { configEntry => + configProps.setProperty(configEntry.name, configEntry.value) + } // Log client errors at a lower level than unexpected exceptions - val message = s"Error processing alter configs request for resource $resource, config $config" + val message = s"Error processing alter configs request for resource $resource, config ${toLoggableProps(resource, configProps).mkString(",")}" if (e.isInstanceOf[ApiException]) info(message, e) else From 01967e48a279e83f836dca557b9e503b73c17106 Mon Sep 17 00:00:00 2001 From: "Matthias J. Sax" Date: Tue, 15 Jun 2021 00:59:55 -0700 Subject: [PATCH 7/9] KAFKA-12914: StreamSourceNode should return `null` topic name for pattern subscription (#10846) Reviewers: Luke Chen , Bruno Cadonna , Guozhang Wang --- .../kstream/internals/InternalStreamsBuilder.java | 11 ++++++----- .../kstream/internals/graph/SourceGraphNode.java | 9 +++++---- .../kstream/internals/graph/StreamSourceNode.java | 11 +++++------ .../kstream/internals/graph/TableSourceNode.java | 12 +++++++++++- 4 files changed, 27 insertions(+), 16 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilder.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilder.java index 1527263af5cef..fbfd9e173a44b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilder.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilder.java @@ -314,16 +314,17 @@ private void mergeDuplicateSourceNodes() { if (graphNode instanceof StreamSourceNode) { final StreamSourceNode currentSourceNode = (StreamSourceNode) graphNode; - if (currentSourceNode.topicPattern() != null) { - if (!patternsToSourceNodes.containsKey(currentSourceNode.topicPattern())) { - patternsToSourceNodes.put(currentSourceNode.topicPattern(), currentSourceNode); + if (currentSourceNode.topicPattern().isPresent()) { + final Pattern topicPattern = currentSourceNode.topicPattern().get(); + if (!patternsToSourceNodes.containsKey(topicPattern)) { + patternsToSourceNodes.put(topicPattern, currentSourceNode); } else { - final StreamSourceNode mainSourceNode = patternsToSourceNodes.get(currentSourceNode.topicPattern()); + final StreamSourceNode mainSourceNode = patternsToSourceNodes.get(topicPattern); mainSourceNode.merge(currentSourceNode); root.removeChild(graphNode); } } else { - for (final String topic : currentSourceNode.topicNames()) { + for (final String topic : currentSourceNode.topicNames().get()) { if (!topicsToSourceNodes.containsKey(topic)) { topicsToSourceNodes.put(topic, currentSourceNode); } else { diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/SourceGraphNode.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/SourceGraphNode.java index affce1bc0e68d..05052279b5721 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/SourceGraphNode.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/SourceGraphNode.java @@ -19,6 +19,7 @@ import java.util.Collection; import java.util.Collections; import java.util.HashSet; +import java.util.Optional; import java.util.Set; import java.util.regex.Pattern; import org.apache.kafka.common.serialization.Serde; @@ -51,12 +52,12 @@ public SourceGraphNode(final String nodeName, this.consumedInternal = consumedInternal; } - public Set topicNames() { - return Collections.unmodifiableSet(topicNames); + public Optional> topicNames() { + return topicNames == null ? Optional.empty() : Optional.of(Collections.unmodifiableSet(topicNames)); } - public Pattern topicPattern() { - return topicPattern; + public Optional topicPattern() { + return Optional.ofNullable(topicPattern); } public ConsumedInternal consumedInternal() { diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/StreamSourceNode.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/StreamSourceNode.java index d4adc894de5e0..e68a9c6039b4a 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/StreamSourceNode.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/StreamSourceNode.java @@ -63,8 +63,8 @@ public void merge(final StreamSourceNode other) { @Override public String toString() { return "StreamSourceNode{" + - "topicNames=" + topicNames() + - ", topicPattern=" + topicPattern() + + "topicNames=" + (topicNames().isPresent() ? topicNames().get() : null) + + ", topicPattern=" + (topicPattern().isPresent() ? topicPattern().get() : null) + ", consumedInternal=" + consumedInternal() + "} " + super.toString(); } @@ -72,21 +72,20 @@ public String toString() { @Override public void writeToTopology(final InternalTopologyBuilder topologyBuilder, final Properties props) { - if (topicPattern() != null) { + if (topicPattern().isPresent()) { topologyBuilder.addSource(consumedInternal().offsetResetPolicy(), nodeName(), consumedInternal().timestampExtractor(), consumedInternal().keyDeserializer(), consumedInternal().valueDeserializer(), - topicPattern()); + topicPattern().get()); } else { topologyBuilder.addSource(consumedInternal().offsetResetPolicy(), nodeName(), consumedInternal().timestampExtractor(), consumedInternal().keyDeserializer(), consumedInternal().valueDeserializer(), - topicNames().toArray(new String[0])); - + topicNames().get().toArray(new String[0])); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/TableSourceNode.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/TableSourceNode.java index b708961073401..77ab7ad6b103a 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/TableSourceNode.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/TableSourceNode.java @@ -17,6 +17,7 @@ package org.apache.kafka.streams.kstream.internals.graph; +import java.util.Iterator; import org.apache.kafka.common.utils.Bytes; import org.apache.kafka.streams.kstream.internals.ConsumedInternal; import org.apache.kafka.streams.kstream.internals.KTableSource; @@ -83,7 +84,16 @@ public static TableSourceNodeBuilder tableSourceNodeBuilder() { @Override @SuppressWarnings("unchecked") public void writeToTopology(final InternalTopologyBuilder topologyBuilder, final Properties props) { - final String topicName = topicNames().iterator().next(); + final String topicName; + if (topicNames().isPresent()) { + final Iterator topicNames = topicNames().get().iterator(); + topicName = topicNames.next(); + if (topicNames.hasNext()) { + throw new IllegalStateException("A table source node must have a single topic as input"); + } + } else { + throw new IllegalStateException("A table source node must have a single topic as input"); + } // TODO: we assume source KTables can only be timestamped-key-value stores for now. // should be expanded for other types of stores as well. From 1e88c758dea8c1f1cd5f5724e98a6af93637644c Mon Sep 17 00:00:00 2001 From: Rajini Sivaram Date: Tue, 15 Jun 2021 09:18:30 +0100 Subject: [PATCH 8/9] KAFKA-12948: Remove node from ClusterConnectionStates.connectingNodes when node is removed (#10882) NetworkClient.poll() throws IllegalStateException when checking isConnectionSetupTimeout if all nodes in ClusterConnectionStates.connectingNodes aren't present in ClusterConnectionStates.nodeState. This commit ensures that when we remove a node from nodeState, we also remove from connectingNodes. Reviewers: David Jacot --- .../clients/ClusterConnectionStates.java | 1 + .../kafka/clients/NetworkClientTest.java | 33 +++++++++++++++++++ 2 files changed, 34 insertions(+) diff --git a/clients/src/main/java/org/apache/kafka/clients/ClusterConnectionStates.java b/clients/src/main/java/org/apache/kafka/clients/ClusterConnectionStates.java index b9d2b13a602d1..524d54b0923b5 100644 --- a/clients/src/main/java/org/apache/kafka/clients/ClusterConnectionStates.java +++ b/clients/src/main/java/org/apache/kafka/clients/ClusterConnectionStates.java @@ -387,6 +387,7 @@ private void updateConnectionSetupTimeout(NodeConnectionState nodeState) { */ public void remove(String id) { nodeState.remove(id); + connectingNodes.remove(id); } /** diff --git a/clients/src/test/java/org/apache/kafka/clients/NetworkClientTest.java b/clients/src/test/java/org/apache/kafka/clients/NetworkClientTest.java index 47b5b201141fd..4fbfd4293409a 100644 --- a/clients/src/test/java/org/apache/kafka/clients/NetworkClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/NetworkClientTest.java @@ -1064,6 +1064,39 @@ public void testFailedConnectionToFirstAddressAfterReconnect() { assertEquals(2, mockHostResolver.resolutionCount()); } + @Test + public void testCloseConnectingNode() { + Cluster cluster = TestUtils.clusterWith(2); + Node node0 = cluster.nodeById(0); + Node node1 = cluster.nodeById(1); + client.ready(node0, time.milliseconds()); + selector.serverConnectionBlocked(node0.idString()); + client.poll(1, time.milliseconds()); + client.close(node0.idString()); + + // Poll without any connections should return without exceptions + client.poll(0, time.milliseconds()); + assertFalse(NetworkClientUtils.isReady(client, node0, time.milliseconds())); + assertFalse(NetworkClientUtils.isReady(client, node1, time.milliseconds())); + + // Connection to new node should work + client.ready(node1, time.milliseconds()); + ByteBuffer buffer = RequestTestUtils.serializeResponseWithHeader(defaultApiVersionsResponse(), ApiKeys.API_VERSIONS.latestVersion(), 0); + selector.delayedReceive(new DelayedReceive(node1.idString(), new NetworkReceive(node1.idString(), buffer))); + while (!client.ready(node1, time.milliseconds())) + client.poll(1, time.milliseconds()); + assertTrue(client.isReady(node1, time.milliseconds())); + selector.clear(); + + // New connection to node closed earlier should work + client.ready(node0, time.milliseconds()); + buffer = RequestTestUtils.serializeResponseWithHeader(defaultApiVersionsResponse(), ApiKeys.API_VERSIONS.latestVersion(), 1); + selector.delayedReceive(new DelayedReceive(node0.idString(), new NetworkReceive(node0.idString(), buffer))); + while (!client.ready(node0, time.milliseconds())) + client.poll(1, time.milliseconds()); + assertTrue(client.isReady(node0, time.milliseconds())); + } + private RequestHeader parseHeader(ByteBuffer buffer) { buffer.getInt(); // skip size return RequestHeader.parse(buffer.slice()); From c16711cb8e0d1c03f123e3e9d7e3d810796bf315 Mon Sep 17 00:00:00 2001 From: Justine Olshan Date: Tue, 15 Jun 2021 06:09:46 -0700 Subject: [PATCH 9/9] KAFKA-12701: NPE in MetadataRequest when using topic IDs (#10584) We prevent handling MetadataRequests where the topic name is null (to prevent NPE) as well as prevent requests that set topic IDs since this functionality has not yet been implemented. When we do implement it in https://github.com/apache/kafka/pull/9769, we should bump the request/response version. Added tests to ensure the error is thrown. Reviewers: dengziming , Ismael Juma --- .../common/requests/MetadataRequest.java | 21 ++++++++-- .../common/message/MetadataRequest.json | 3 +- .../common/requests/MetadataRequestTest.java | 25 ++++++++++++ .../main/scala/kafka/server/KafkaApis.scala | 11 ++++++ .../unit/kafka/server/KafkaApisTest.scala | 38 +++++++++++++++++++ 5 files changed, 94 insertions(+), 4 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/requests/MetadataRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/MetadataRequest.java index 816f600061568..d38e9ac47b80a 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/MetadataRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/MetadataRequest.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.common.requests; +import org.apache.kafka.common.Uuid; import org.apache.kafka.common.errors.UnsupportedVersionException; import org.apache.kafka.common.message.MetadataRequestData; import org.apache.kafka.common.message.MetadataRequestData.MetadataRequestTopic; @@ -92,6 +93,16 @@ public MetadataRequest build(short version) { if (!data.allowAutoTopicCreation() && version < 4) throw new UnsupportedVersionException("MetadataRequest versions older than 4 don't support the " + "allowAutoTopicCreation field"); + if (data.topics() != null) { + data.topics().forEach(topic -> { + if (topic.name() == null) + throw new UnsupportedVersionException("MetadataRequest version " + version + + " does not support null topic names."); + if (topic.topicId() != Uuid.ZERO_UUID) + throw new UnsupportedVersionException("MetadataRequest version " + version + + " does not support non-zero topic IDs."); + }); + } return new MetadataRequest(data, version); } @@ -117,13 +128,17 @@ public MetadataRequestData data() { public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) { Errors error = Errors.forException(e); MetadataResponseData responseData = new MetadataResponseData(); - if (topics() != null) { - for (String topic : topics()) + if (data.topics() != null) { + for (MetadataRequestTopic topic : data.topics()) { + // the response does not allow null, so convert to empty string if necessary + String topicName = topic.name() == null ? "" : topic.name(); responseData.topics().add(new MetadataResponseData.MetadataResponseTopic() - .setName(topic) + .setName(topicName) + .setTopicId(topic.topicId()) .setErrorCode(error.code()) .setIsInternal(false) .setPartitions(Collections.emptyList())); + } } responseData.setThrottleTimeMs(throttleTimeMs); diff --git a/clients/src/main/resources/common/message/MetadataRequest.json b/clients/src/main/resources/common/message/MetadataRequest.json index e5083a84ea343..a1634b19970d7 100644 --- a/clients/src/main/resources/common/message/MetadataRequest.json +++ b/clients/src/main/resources/common/message/MetadataRequest.json @@ -33,7 +33,8 @@ // // Version 9 is the first flexible version. // - // Version 10 adds topicId. + // Version 10 adds topicId and allows name field to be null. However, this functionality was not implemented on the server. + // Versions 10 and 11 should not use the topicId field or set topic name to null. // // Version 11 deprecates IncludeClusterAuthorizedOperations field. This is now exposed // by the DescribeCluster API (KIP-700). diff --git a/clients/src/test/java/org/apache/kafka/common/requests/MetadataRequestTest.java b/clients/src/test/java/org/apache/kafka/common/requests/MetadataRequestTest.java index e51523297597d..74c217df91f86 100644 --- a/clients/src/test/java/org/apache/kafka/common/requests/MetadataRequestTest.java +++ b/clients/src/test/java/org/apache/kafka/common/requests/MetadataRequestTest.java @@ -16,16 +16,21 @@ */ package org.apache.kafka.common.requests; +import org.apache.kafka.common.Uuid; +import org.apache.kafka.common.errors.UnsupportedVersionException; import org.apache.kafka.common.message.MetadataRequestData; import org.apache.kafka.common.protocol.ApiKeys; import org.junit.jupiter.api.Test; +import java.util.Arrays; import java.util.Collections; +import java.util.List; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assertions.assertThrows; public class MetadataRequestTest { @@ -65,4 +70,24 @@ public void testMetadataRequestVersion() { assertEquals(minVersion, builder3.oldestAllowedVersion()); assertEquals(maxVersion, builder3.latestAllowedVersion()); } + + @Test + public void testTopicIdAndNullTopicNameRequests() { + // Construct invalid MetadataRequestTopics. We will build each one separately and ensure the error is thrown. + List topics = Arrays.asList( + new MetadataRequestData.MetadataRequestTopic().setName(null).setTopicId(Uuid.randomUuid()), + new MetadataRequestData.MetadataRequestTopic().setName(null), + new MetadataRequestData.MetadataRequestTopic().setTopicId(Uuid.randomUuid()), + new MetadataRequestData.MetadataRequestTopic().setName("topic").setTopicId(Uuid.randomUuid())); + + // if version is 10 or 11, the invalid topic metadata should return an error + List invalidVersions = Arrays.asList((short) 10, (short) 11); + invalidVersions.forEach(version -> + topics.forEach(topic -> { + MetadataRequestData metadataRequestData = new MetadataRequestData().setTopics(Collections.singletonList(topic)); + MetadataRequest.Builder builder = new MetadataRequest.Builder(metadataRequestData); + assertThrows(UnsupportedVersionException.class, () -> builder.build(version)); + }) + ); + } } diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index bafc0411253b4..2724de39a19a7 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -1142,6 +1142,17 @@ class KafkaApis(val requestChannel: RequestChannel, val metadataRequest = request.body[MetadataRequest] val requestVersion = request.header.apiVersion + // Topic IDs are not supported for versions 10 and 11. Topic names can not be null in these versions. + if (!metadataRequest.isAllTopics) { + metadataRequest.data.topics.forEach{ topic => + if (topic.name == null) { + throw new InvalidRequestException(s"Topic name can not be null for version ${metadataRequest.version}") + } else if (topic.topicId != Uuid.ZERO_UUID) { + throw new InvalidRequestException(s"Topic IDs are not supported in requests for version ${metadataRequest.version}") + } + } + } + val topics = if (metadataRequest.isAllTopics) metadataCache.getAllTopics() else diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index 9a0bdb04211ca..a6da170a426ad 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -1052,6 +1052,44 @@ class KafkaApisTest { numBrokersNeeded - 1) } + @Test + def testInvalidMetadataRequestReturnsError(): Unit = { + // Construct invalid MetadataRequestTopics. We will try each one separately and ensure the error is thrown. + val topics = List(new MetadataRequestData.MetadataRequestTopic().setName(null).setTopicId(Uuid.randomUuid()), + new MetadataRequestData.MetadataRequestTopic().setName(null), + new MetadataRequestData.MetadataRequestTopic().setTopicId(Uuid.randomUuid()), + new MetadataRequestData.MetadataRequestTopic().setName("topic1").setTopicId(Uuid.randomUuid())) + + EasyMock.replay(replicaManager, clientRequestQuotaManager, + autoTopicCreationManager, forwardingManager, clientControllerQuotaManager, groupCoordinator, txnCoordinator) + + // if version is 10 or 11, the invalid topic metadata should return an error + val invalidVersions = Set(10, 11) + invalidVersions.foreach( version => + topics.foreach(topic => { + val metadataRequestData = new MetadataRequestData().setTopics(Collections.singletonList(topic)) + val request = buildRequest(new MetadataRequest(metadataRequestData, version.toShort)) + val kafkaApis = createKafkaApis() + + val capturedResponse = EasyMock.newCapture[AbstractResponse]() + EasyMock.expect(requestChannel.sendResponse( + EasyMock.eq(request), + EasyMock.capture(capturedResponse), + EasyMock.anyObject() + )) + + EasyMock.replay(requestChannel) + kafkaApis.handle(request, RequestLocal.withThreadConfinedCaching) + + val response = capturedResponse.getValue.asInstanceOf[MetadataResponse] + assertEquals(1, response.topicMetadata.size) + assertEquals(1, response.errorCounts.get(Errors.INVALID_REQUEST)) + response.data.topics.forEach(topic => assertNotEquals(null, topic.name)) + reset(requestChannel) + }) + ) + } + @Test def testOffsetCommitWithInvalidPartition(): Unit = { val topic = "topic"