diff --git a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java index 9e92fc4386a57..fdf38e21ad94b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java +++ b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java @@ -764,7 +764,7 @@ private KafkaStreams(final InternalTopologyBuilder internalTopologyBuilder, delegatingStateRestoreListener, i + 1); threadState.put(threads[i].getId(), threads[i].state()); - storeProviders.add(new StreamThreadStateStoreProvider(threads[i])); + storeProviders.add(new StreamThreadStateStoreProvider(threads[i], internalTopologyBuilder)); } final StreamStateListener streamStateListener = new StreamStateListener(threadState, globalThreadState); @@ -1160,47 +1160,29 @@ public KeyQueryMetadata queryMetadataForKey(final String storeName, return streamsMetadataState.getKeyQueryMetadataForKey(storeName, key, partitioner); } + /** - * Get a facade wrapping the local {@link StateStore} instances with the provided {@code storeName} if the Store's - * type is accepted by the provided {@link QueryableStoreType#accepts(StateStore) queryableStoreType}. - * The returned object can be used to query the {@link StateStore} instances. - * - * Only permits queries on active replicas of the store (no standbys or restoring replicas). - * See {@link KafkaStreams#store(java.lang.String, org.apache.kafka.streams.state.QueryableStoreType, boolean)} - * for the option to set {@code includeStaleStores} to true and trade off consistency in favor of availability. - * - * @param storeName name of the store to find - * @param queryableStoreType accept only stores that are accepted by {@link QueryableStoreType#accepts(StateStore)} - * @param return type - * @return A facade wrapping the local {@link StateStore} instances - * @throws InvalidStateStoreException if Kafka Streams is (re-)initializing or a store with {@code storeName} and - * {@code queryableStoreType} doesn't exist + * @deprecated since 2.5 release; use {@link #store(StoreQueryParams)} instead */ + @Deprecated public T store(final String storeName, final QueryableStoreType queryableStoreType) { - return store(storeName, queryableStoreType, false); + return store(StoreQueryParams.fromNameAndType(storeName, queryableStoreType)); } /** - * Get a facade wrapping the local {@link StateStore} instances with the provided {@code storeName} if the Store's + * Get a facade wrapping the local {@link StateStore} instances with the provided {@link StoreQueryParams}. + * StoreQueryParams need required parameters to be set, which are {@code storeName} and if * type is accepted by the provided {@link QueryableStoreType#accepts(StateStore) queryableStoreType}. * The returned object can be used to query the {@link StateStore} instances. * - * @param storeName name of the store to find - * @param queryableStoreType accept only stores that are accepted by {@link QueryableStoreType#accepts(StateStore)} - * @param includeStaleStores If false, only permit queries on the active replica for a partition, and only if the - * task for that partition is running. I.e., the state store is not a standby replica, - * and it is not restoring from the changelog. - * If true, allow queries on standbys and restoring replicas in addition to active ones. - * @param return type + * @param storeQueryParams to set the optional parameters to fetch type of stores user wants to fetch when a key is queried * @return A facade wrapping the local {@link StateStore} instances * @throws InvalidStateStoreException if Kafka Streams is (re-)initializing or a store with {@code storeName} and * {@code queryableStoreType} doesn't exist */ - public T store(final String storeName, - final QueryableStoreType queryableStoreType, - final boolean includeStaleStores) { + public T store(final StoreQueryParams storeQueryParams) { validateIsRunningOrRebalancing(); - return queryableStoreProvider.getStore(storeName, queryableStoreType, includeStaleStores); + return queryableStoreProvider.getStore(storeQueryParams); } /** diff --git a/streams/src/main/java/org/apache/kafka/streams/StoreQueryParams.java b/streams/src/main/java/org/apache/kafka/streams/StoreQueryParams.java new file mode 100644 index 0000000000000..2a6bbd8136359 --- /dev/null +++ b/streams/src/main/java/org/apache/kafka/streams/StoreQueryParams.java @@ -0,0 +1,138 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.streams; + +import org.apache.kafka.streams.state.QueryableStoreType; + +import java.util.Objects; + +/** + * Represents all the query options that a user can provide to state what kind of stores it is expecting. + * The options would be whether a user would want to enable/disable stale stores + * or whether it knows the list of partitions that it specifically wants to fetch. + * If this information is not provided the default behavior is to fetch the stores for all the partitions + * available on that instance for that particular store name. + * It contains a partition, which for a point queries can be populated from the {@link KeyQueryMetadata}. + */ +public class StoreQueryParams { + + private Integer partition; + private boolean staleStores; + private final String storeName; + private final QueryableStoreType queryableStoreType; + + private StoreQueryParams(final String storeName, final QueryableStoreType queryableStoreType) { + this.storeName = storeName; + this.queryableStoreType = queryableStoreType; + } + + public static StoreQueryParams fromNameAndType(final String storeName, + final QueryableStoreType queryableStoreType) { + return new StoreQueryParams(storeName, queryableStoreType); + } + + /** + * Set a specific partition that should be queried exclusively. + * + * @param partition The specific integer partition to be fetched from the stores list by using {@link StoreQueryParams}. + * + * @return String storeName + */ + public StoreQueryParams withPartition(final Integer partition) { + final StoreQueryParams storeQueryParams = StoreQueryParams.fromNameAndType(this.storeName(), this.queryableStoreType()); + storeQueryParams.partition = partition; + storeQueryParams.staleStores = this.staleStores; + return storeQueryParams; + } + + /** + * Enable querying of stale state stores, i.e., allow to query active tasks during restore as well as standby tasks. + * + * @return String storeName + */ + public StoreQueryParams enableStaleStores() { + final StoreQueryParams storeQueryParams = StoreQueryParams.fromNameAndType(this.storeName(), this.queryableStoreType()); + storeQueryParams.partition = this.partition; + storeQueryParams.staleStores = true; + return storeQueryParams; + } + + /** + * Get the store name for which key is queried by the user. + * + * @return String storeName + */ + public String storeName() { + return storeName; + } + + /** + * Get the queryable store type for which key is queried by the user. + * + * @return QueryableStoreType queryableStoreType + */ + public QueryableStoreType queryableStoreType() { + return queryableStoreType; + } + + /** + * Get the partition to be used to fetch list of stores. + * If the method returns {@code null}, it would mean that no specific partition has been requested, + * so all the local partitions for the store will be returned. + * + * @return Integer partition + */ + public Integer partition() { + return partition; + } + + /** + * Get the flag staleStores. If {@code true}, include standbys and recovering stores along with running stores. + * + * @return boolean staleStores + */ + public boolean staleStoresEnabled() { + return staleStores; + } + + @Override + public boolean equals(final Object obj) { + if (!(obj instanceof StoreQueryParams)) { + return false; + } + final StoreQueryParams storeQueryParams = (StoreQueryParams) obj; + return Objects.equals(storeQueryParams.partition, partition) + && Objects.equals(storeQueryParams.staleStores, staleStores) + && Objects.equals(storeQueryParams.storeName, storeName) + && Objects.equals(storeQueryParams.queryableStoreType, queryableStoreType); + } + + @Override + public String toString() { + return "StoreQueryParams {" + + "partition=" + partition + + ", staleStores=" + staleStores + + ", storeName=" + storeName + + ", queryableStoreType=" + queryableStoreType + + '}'; + } + + @Override + public int hashCode() { + return Objects.hash(partition, staleStores, storeName, queryableStoreType); + } +} \ No newline at end of file diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopologyBuilder.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopologyBuilder.java index a1db8d9b3c9fb..ccedea353b14d 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopologyBuilder.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopologyBuilder.java @@ -1499,7 +1499,7 @@ public void addPredecessor(final TopologyDescription.Node predecessor) { @Override public String toString() { final String topicsString = topics == null ? topicPattern.toString() : topics.toString(); - + return "Source: " + name + " (topics: " + topicsString + ")\n --> " + nodeNames(successors); } @@ -1701,7 +1701,7 @@ public int hashCode() { public static class TopicsInfo { final Set sinkTopics; - final Set sourceTopics; + public final Set sourceTopics; public final Map stateChangelogTopics; public final Map repartitionSourceTopics; diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/QueryableStoreProvider.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/QueryableStoreProvider.java index b25207edba813..4458b265f12db 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/QueryableStoreProvider.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/QueryableStoreProvider.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.streams.state.internals; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.processor.StateStore; import org.apache.kafka.streams.state.QueryableStoreType; @@ -41,29 +42,29 @@ public QueryableStoreProvider(final List storePr * Get a composite object wrapping the instances of the {@link StateStore} with the provided * storeName and {@link QueryableStoreType} * - * @param storeName name of the store - * @param queryableStoreType accept stores passing {@link QueryableStoreType#accepts(StateStore)} - * @param includeStaleStores if true, include standbys and recovering stores; - * if false, only include running actives. + * @param storeQueryParams if stateStoresEnabled is used i.e. staleStoresEnabled is true, include standbys and recovering stores; + * if stateStoresDisabled i.e. staleStoresEnabled is false, only include running actives; + * if partition is null then it fetches all local partitions on the instance; + * if partition is set then it fetches a specific partition. * @param The expected type of the returned store * @return A composite object that wraps the store instances. */ - public T getStore(final String storeName, - final QueryableStoreType queryableStoreType, - final boolean includeStaleStores) { + public T getStore(final StoreQueryParams storeQueryParams) { + final String storeName = storeQueryParams.storeName(); + final QueryableStoreType queryableStoreType = storeQueryParams.queryableStoreType(); final List globalStore = globalStoreProvider.stores(storeName, queryableStoreType); if (!globalStore.isEmpty()) { return queryableStoreType.create(globalStoreProvider, storeName); } final List allStores = new ArrayList<>(); for (final StreamThreadStateStoreProvider storeProvider : storeProviders) { - allStores.addAll(storeProvider.stores(storeName, queryableStoreType, includeStaleStores)); + allStores.addAll(storeProvider.stores(storeQueryParams)); } if (allStores.isEmpty()) { throw new InvalidStateStoreException("The state store, " + storeName + ", may have migrated to another instance."); } return queryableStoreType.create( - new WrappingStoreProvider(storeProviders, includeStaleStores), + new WrappingStoreProvider(storeProviders, storeQueryParams), storeName ); } diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProvider.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProvider.java index 70c87ae191bc1..eeb852d4a2bcb 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProvider.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProvider.java @@ -16,9 +16,11 @@ */ package org.apache.kafka.streams.state.internals; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.processor.StateStore; import org.apache.kafka.streams.processor.TaskId; +import org.apache.kafka.streams.processor.internals.InternalTopologyBuilder; import org.apache.kafka.streams.processor.internals.StreamThread; import org.apache.kafka.streams.processor.internals.Task; import org.apache.kafka.streams.state.QueryableStoreType; @@ -30,27 +32,36 @@ import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.Set; +import java.util.stream.Collectors; public class StreamThreadStateStoreProvider { private final StreamThread streamThread; + private final InternalTopologyBuilder internalTopologyBuilder; - public StreamThreadStateStoreProvider(final StreamThread streamThread) { + public StreamThreadStateStoreProvider(final StreamThread streamThread, + final InternalTopologyBuilder internalTopologyBuilder) { this.streamThread = streamThread; + this.internalTopologyBuilder = internalTopologyBuilder; } @SuppressWarnings("unchecked") - public List stores(final String storeName, - final QueryableStoreType queryableStoreType, - final boolean includeStaleStores) { + public List stores(final StoreQueryParams storeQueryParams) { + final String storeName = storeQueryParams.storeName(); + final QueryableStoreType queryableStoreType = storeQueryParams.queryableStoreType(); + final TaskId keyTaskId = createKeyTaskId(storeName, storeQueryParams.partition()); if (streamThread.state() == StreamThread.State.DEAD) { return Collections.emptyList(); } final StreamThread.State state = streamThread.state(); - if (includeStaleStores ? state.isAlive() : state == StreamThread.State.RUNNING) { - final Map tasks = includeStaleStores ? streamThread.allTasks() : streamThread.activeTasks(); + if (storeQueryParams.staleStoresEnabled() ? state.isAlive() : state == StreamThread.State.RUNNING) { + final Map tasks = storeQueryParams.staleStoresEnabled() ? streamThread.allTasks() : streamThread.activeTasks(); final List stores = new ArrayList<>(); for (final Task streamTask : tasks.values()) { + if (keyTaskId != null && !keyTaskId.equals(streamTask.id())) { + continue; + } final StateStore store = streamTask.getStore(storeName); if (store != null && queryableStoreType.accepts(store)) { if (!store.isOpen()) { @@ -72,8 +83,23 @@ public List stores(final String storeName, } else { throw new InvalidStateStoreException("Cannot get state store " + storeName + " because the stream thread is " + state + ", not RUNNING" + - (includeStaleStores ? " or REBALANCING" : "")); + (storeQueryParams.staleStoresEnabled() ? " or REBALANCING" : "")); } } + private TaskId createKeyTaskId(final String storeName, final Integer partition) { + if (partition == null) { + return null; + } + final List sourceTopics = internalTopologyBuilder.stateStoreNameToSourceTopics().get(storeName); + final Set sourceTopicsSet = sourceTopics.stream().collect(Collectors.toSet()); + final Map topicGroups = internalTopologyBuilder.topicGroups(); + for (final Map.Entry topicGroup : topicGroups.entrySet()) { + if (topicGroup.getValue().sourceTopics.containsAll(sourceTopicsSet)) { + return new TaskId(topicGroup.getKey(), partition.intValue()); + } + } + throw new InvalidStateStoreException("Cannot get state store " + storeName + " because the requested partition " + partition + "is" + + "not available on this instance"); + } } diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/WrappingStoreProvider.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/WrappingStoreProvider.java index 9231dfd7ed239..2caee8b5e2d9b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/WrappingStoreProvider.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/WrappingStoreProvider.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.streams.state.internals; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.state.QueryableStoreType; @@ -28,12 +29,17 @@ public class WrappingStoreProvider implements StateStoreProvider { private final List storeProviders; - private final boolean includeStaleStores; + private StoreQueryParams storeQueryParams; WrappingStoreProvider(final List storeProviders, - final boolean includeStaleStores) { + final StoreQueryParams storeQueryParams) { this.storeProviders = storeProviders; - this.includeStaleStores = includeStaleStores; + this.storeQueryParams = storeQueryParams; + } + + //visible for testing + public void setStoreQueryParams(final StoreQueryParams storeQueryParams) { + this.storeQueryParams = storeQueryParams; } @Override @@ -41,7 +47,7 @@ public List stores(final String storeName, final QueryableStoreType queryableStoreType) { final List allStores = new ArrayList<>(); for (final StreamThreadStateStoreProvider provider : storeProviders) { - final List stores = provider.stores(storeName, queryableStoreType, includeStaleStores); + final List stores = provider.stores(storeQueryParams); allStores.addAll(stores); } if (allStores.isEmpty()) { diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/EosIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/EosIntegrationTest.java index 3e5a2c3fea2f2..ad7ada9012efd 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/EosIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/EosIntegrationTest.java @@ -26,6 +26,7 @@ import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster; import org.apache.kafka.streams.integration.utils.IntegrationTestUtils; @@ -750,7 +751,7 @@ private void verifyStateStore(final KafkaStreams streams, final long maxWaitingTime = System.currentTimeMillis() + 300000L; while (System.currentTimeMillis() < maxWaitingTime) { try { - store = streams.store(storeName, QueryableStoreTypes.keyValueStore()); + store = streams.store(StoreQueryParams.fromNameAndType(storeName, QueryableStoreTypes.keyValueStore())); break; } catch (final InvalidStateStoreException okJustRetry) { try { diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/GlobalKTableEOSIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/GlobalKTableEOSIntegrationTest.java index b247ec2db43ef..d770f0135b879 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/GlobalKTableEOSIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/GlobalKTableEOSIntegrationTest.java @@ -27,6 +27,7 @@ import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster; import org.apache.kafka.streams.integration.utils.IntegrationTestUtils; @@ -138,7 +139,7 @@ public void shouldKStreamGlobalKTableLeftJoin() throws Exception { produceGlobalTableValues(); final ReadOnlyKeyValueStore replicatedStore = - kafkaStreams.store(globalStore, QueryableStoreTypes.keyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(globalStore, QueryableStoreTypes.keyValueStore())); TestUtils.waitForCondition( () -> "J".equals(replicatedStore.get(5L)), @@ -182,7 +183,7 @@ public void shouldKStreamGlobalKTableJoin() throws Exception { produceGlobalTableValues(); final ReadOnlyKeyValueStore replicatedStore = - kafkaStreams.store(globalStore, QueryableStoreTypes.keyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(globalStore, QueryableStoreTypes.keyValueStore())); TestUtils.waitForCondition( () -> "J".equals(replicatedStore.get(5L)), @@ -219,7 +220,7 @@ public void shouldRestoreTransactionalMessages() throws Exception { () -> { final ReadOnlyKeyValueStore store; try { - store = kafkaStreams.store(globalStore, QueryableStoreTypes.keyValueStore()); + store = kafkaStreams.store(StoreQueryParams.fromNameAndType(globalStore, QueryableStoreTypes.keyValueStore())); } catch (final InvalidStateStoreException ex) { return false; } @@ -253,7 +254,7 @@ public void shouldNotRestoreAbortedMessages() throws Exception { () -> { final ReadOnlyKeyValueStore store; try { - store = kafkaStreams.store(globalStore, QueryableStoreTypes.keyValueStore()); + store = kafkaStreams.store(StoreQueryParams.fromNameAndType(globalStore, QueryableStoreTypes.keyValueStore())); } catch (final InvalidStateStoreException ex) { return false; } diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/GlobalKTableIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/GlobalKTableIntegrationTest.java index 0a9148d61fb10..638134db9f87d 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/GlobalKTableIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/GlobalKTableIntegrationTest.java @@ -26,6 +26,7 @@ import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster; import org.apache.kafka.streams.integration.utils.IntegrationTestUtils; import org.apache.kafka.streams.kstream.Consumed; @@ -37,8 +38,8 @@ import org.apache.kafka.streams.state.KeyValueStore; import org.apache.kafka.streams.state.QueryableStoreTypes; import org.apache.kafka.streams.state.ReadOnlyKeyValueStore; -import org.apache.kafka.streams.state.Stores; import org.apache.kafka.streams.state.ValueAndTimestamp; +import org.apache.kafka.streams.state.Stores; import org.apache.kafka.test.IntegrationTest; import org.apache.kafka.test.MockProcessorSupplier; import org.apache.kafka.test.TestUtils; @@ -141,7 +142,7 @@ public void shouldKStreamGlobalKTableLeftJoin() throws Exception { produceGlobalTableValues(); final ReadOnlyKeyValueStore replicatedStore = - kafkaStreams.store(globalStore, QueryableStoreTypes.keyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(globalStore, QueryableStoreTypes.keyValueStore())); TestUtils.waitForCondition( () -> "J".equals(replicatedStore.get(5L)), @@ -149,7 +150,7 @@ public void shouldKStreamGlobalKTableLeftJoin() throws Exception { "waiting for data in replicated store"); final ReadOnlyKeyValueStore> replicatedStoreWithTimestamp = - kafkaStreams.store(globalStore, QueryableStoreTypes.timestampedKeyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(globalStore, QueryableStoreTypes.timestampedKeyValueStore())); assertThat(replicatedStoreWithTimestamp.get(5L), equalTo(ValueAndTimestamp.make("J", firstTimestamp + 4L))); firstTimestamp = mockTime.milliseconds(); @@ -208,7 +209,7 @@ public void shouldKStreamGlobalKTableJoin() throws Exception { produceGlobalTableValues(); final ReadOnlyKeyValueStore replicatedStore = - kafkaStreams.store(globalStore, QueryableStoreTypes.keyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(globalStore, QueryableStoreTypes.keyValueStore())); TestUtils.waitForCondition( () -> "J".equals(replicatedStore.get(5L)), @@ -216,7 +217,7 @@ public void shouldKStreamGlobalKTableJoin() throws Exception { "waiting for data in replicated store"); final ReadOnlyKeyValueStore> replicatedStoreWithTimestamp = - kafkaStreams.store(globalStore, QueryableStoreTypes.timestampedKeyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(globalStore, QueryableStoreTypes.timestampedKeyValueStore())); assertThat(replicatedStoreWithTimestamp.get(5L), equalTo(ValueAndTimestamp.make("J", firstTimestamp + 4L))); firstTimestamp = mockTime.milliseconds(); @@ -253,17 +254,17 @@ public void shouldRestoreGlobalInMemoryKTableOnRestart() throws Exception { produceInitialGlobalTableValues(); startStreams(); - ReadOnlyKeyValueStore store = kafkaStreams.store(globalStore, QueryableStoreTypes.keyValueStore()); + ReadOnlyKeyValueStore store = kafkaStreams.store(StoreQueryParams.fromNameAndType(globalStore, QueryableStoreTypes.keyValueStore())); assertThat(store.approximateNumEntries(), equalTo(4L)); ReadOnlyKeyValueStore> timestampedStore = - kafkaStreams.store(globalStore, QueryableStoreTypes.timestampedKeyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(globalStore, QueryableStoreTypes.timestampedKeyValueStore())); assertThat(timestampedStore.approximateNumEntries(), equalTo(4L)); kafkaStreams.close(); startStreams(); - store = kafkaStreams.store(globalStore, QueryableStoreTypes.keyValueStore()); + store = kafkaStreams.store(StoreQueryParams.fromNameAndType(globalStore, QueryableStoreTypes.keyValueStore())); assertThat(store.approximateNumEntries(), equalTo(4L)); - timestampedStore = kafkaStreams.store(globalStore, QueryableStoreTypes.timestampedKeyValueStore()); + timestampedStore = kafkaStreams.store(StoreQueryParams.fromNameAndType(globalStore, QueryableStoreTypes.timestampedKeyValueStore())); assertThat(timestampedStore.approximateNumEntries(), equalTo(4L)); } diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/KStreamAggregationIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/KStreamAggregationIntegrationTest.java index dcf72a501ec1a..a23e59a3af227 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/KStreamAggregationIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/KStreamAggregationIntegrationTest.java @@ -29,9 +29,10 @@ import org.apache.kafka.common.serialization.StringSerializer; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.KeyValue; -import org.apache.kafka.streams.KeyValueTimestamp; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.KeyValueTimestamp; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster; import org.apache.kafka.streams.integration.utils.IntegrationTestUtils; import org.apache.kafka.streams.kstream.Aggregator; @@ -673,7 +674,7 @@ public void close() {} // verify can query data via IQ final ReadOnlySessionStore sessionStore = - kafkaStreams.store(userSessionsStore, QueryableStoreTypes.sessionStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(userSessionsStore, QueryableStoreTypes.sessionStore())); final KeyValueIterator, String> bob = sessionStore.fetch("bob"); assertThat(bob.next(), equalTo(KeyValue.pair(new Windowed<>("bob", new SessionWindow(t1, t1)), "start"))); assertThat(bob.next(), equalTo(KeyValue.pair(new Windowed<>("bob", new SessionWindow(t3, t4)), "pause:resume"))); diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/OptimizedKTableIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/OptimizedKTableIntegrationTest.java index e2ef3f13f1531..acdf8019aded5 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/OptimizedKTableIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/OptimizedKTableIntegrationTest.java @@ -45,11 +45,12 @@ import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.KeyQueryMetadata; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster; import org.apache.kafka.streams.integration.utils.IntegrationTestUtils; import org.apache.kafka.streams.kstream.Consumed; import org.apache.kafka.streams.kstream.Materialized; -import org.apache.kafka.streams.KeyQueryMetadata; import org.apache.kafka.streams.processor.StateRestoreListener; import org.apache.kafka.streams.state.KeyValueStore; import org.apache.kafka.streams.state.QueryableStoreTypes; @@ -152,10 +153,10 @@ public void shouldApplyUpdatesToStandbyStore() throws Exception { assertThat(semaphore.tryAcquire(batch1NumMessages, 60, TimeUnit.SECONDS), is(equalTo(true))); final ReadOnlyKeyValueStore store1 = kafkaStreams1 - .store(TABLE_NAME, QueryableStoreTypes.keyValueStore()); + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, QueryableStoreTypes.keyValueStore())); final ReadOnlyKeyValueStore store2 = kafkaStreams2 - .store(TABLE_NAME, QueryableStoreTypes.keyValueStore()); + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, QueryableStoreTypes.keyValueStore())); final boolean kafkaStreams1WasFirstActive; final KeyQueryMetadata keyQueryMetadata = kafkaStreams1.queryMetadataForKey(TABLE_NAME, key, (topic, somekey, value, numPartitions) -> 0); diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/QueryableStateIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/QueryableStateIntegrationTest.java index 64e91fc614fa9..f7dc302878f7a 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/QueryableStateIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/QueryableStateIntegrationTest.java @@ -31,11 +31,12 @@ import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.KafkaStreams.State; import org.apache.kafka.streams.KafkaStreamsTest; -import org.apache.kafka.streams.KeyQueryMetadata; import org.apache.kafka.streams.KeyValue; -import org.apache.kafka.streams.LagInfo; +import org.apache.kafka.streams.KeyQueryMetadata; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.LagInfo; import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster; import org.apache.kafka.streams.integration.utils.IntegrationTestUtils; @@ -48,13 +49,14 @@ import org.apache.kafka.streams.kstream.Produced; import org.apache.kafka.streams.kstream.TimeWindows; import org.apache.kafka.streams.kstream.ValueMapper; -import org.apache.kafka.streams.state.KeyValueIterator; -import org.apache.kafka.streams.state.KeyValueStore; +import org.apache.kafka.streams.state.QueryableStoreType; import org.apache.kafka.streams.state.QueryableStoreTypes; import org.apache.kafka.streams.state.ReadOnlyKeyValueStore; -import org.apache.kafka.streams.state.ReadOnlyWindowStore; import org.apache.kafka.streams.state.StreamsMetadata; +import org.apache.kafka.streams.state.ReadOnlyWindowStore; import org.apache.kafka.streams.state.WindowStoreIterator; +import org.apache.kafka.streams.state.KeyValueIterator; +import org.apache.kafka.streams.state.KeyValueStore; import org.apache.kafka.test.IntegrationTest; import org.apache.kafka.test.MockMapper; import org.apache.kafka.test.NoRetryException; @@ -303,8 +305,9 @@ private void verifyAllKVKeys(final List streamsList, final int index = queryMetadata.getActiveHost().port(); final KafkaStreams streamsWithKey = pickInstanceByPort ? streamsList.get(index) : streams; + final QueryableStoreType> queryableStoreType = QueryableStoreTypes.keyValueStore(); final ReadOnlyKeyValueStore store = - streamsWithKey.store(storeName, QueryableStoreTypes.keyValueStore(), true); + streamsWithKey.store(StoreQueryParams.fromNameAndType(storeName, queryableStoreType).enableStaleStores()); if (store == null) { nullStoreKeys.add(key); continue; @@ -362,8 +365,9 @@ private void verifyAllWindowedKeys(final List streamsList, final int index = queryMetadata.getActiveHost().port(); final KafkaStreams streamsWithKey = pickInstanceByPort ? streamsList.get(index) : streams; + final QueryableStoreType> queryableWindowStoreType = QueryableStoreTypes.windowStore(); final ReadOnlyWindowStore store = - streamsWithKey.store(storeName, QueryableStoreTypes.windowStore(), true); + streamsWithKey.store(StoreQueryParams.fromNameAndType(storeName, queryableWindowStoreType).enableStaleStores()); if (store == null) { nullStoreKeys.add(key); continue; @@ -646,10 +650,10 @@ public void concurrentAccesses() throws Exception { waitUntilAtLeastNumRecordProcessed(outputTopicConcurrentWindowed, numberOfWordsPerIteration); final ReadOnlyKeyValueStore keyValueStore = - kafkaStreams.store(storeName + "-" + streamConcurrent, QueryableStoreTypes.keyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(storeName + "-" + streamConcurrent, QueryableStoreTypes.keyValueStore())); final ReadOnlyWindowStore windowStore = - kafkaStreams.store(windowStoreName + "-" + streamConcurrent, QueryableStoreTypes.windowStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(windowStoreName + "-" + streamConcurrent, QueryableStoreTypes.windowStore())); final Map expectedWindowState = new HashMap<>(); final Map expectedCount = new HashMap<>(); @@ -712,9 +716,9 @@ public void shouldBeAbleToQueryFilterState() throws Exception { waitUntilAtLeastNumRecordProcessed(outputTopic, 1); final ReadOnlyKeyValueStore - myFilterStore = kafkaStreams.store("queryFilter", QueryableStoreTypes.keyValueStore()); + myFilterStore = kafkaStreams.store(StoreQueryParams.fromNameAndType("queryFilter", QueryableStoreTypes.keyValueStore())); final ReadOnlyKeyValueStore - myFilterNotStore = kafkaStreams.store("queryFilterNot", QueryableStoreTypes.keyValueStore()); + myFilterNotStore = kafkaStreams.store(StoreQueryParams.fromNameAndType("queryFilterNot", QueryableStoreTypes.keyValueStore())); for (final KeyValue expectedEntry : expectedBatch1) { TestUtils.waitForCondition(() -> expectedEntry.value.equals(myFilterStore.get(expectedEntry.key)), @@ -778,7 +782,7 @@ public void shouldBeAbleToQueryMapValuesState() throws Exception { waitUntilAtLeastNumRecordProcessed(outputTopic, 5); final ReadOnlyKeyValueStore myMapStore = - kafkaStreams.store("queryMapValues", QueryableStoreTypes.keyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType("queryMapValues", QueryableStoreTypes.keyValueStore())); for (final KeyValue batchEntry : batch1) { assertEquals(Long.valueOf(batchEntry.value), myMapStore.get(batchEntry.key)); } @@ -826,8 +830,8 @@ public void shouldBeAbleToQueryMapValuesAfterFilterState() throws Exception { waitUntilAtLeastNumRecordProcessed(outputTopic, 1); final ReadOnlyKeyValueStore - myMapStore = kafkaStreams.store("queryMapValues", - QueryableStoreTypes.keyValueStore()); + myMapStore = kafkaStreams.store(StoreQueryParams.fromNameAndType("queryMapValues", + QueryableStoreTypes.keyValueStore())); for (final KeyValue expectedEntry : expectedBatch1) { assertEquals(expectedEntry.value, myMapStore.get(expectedEntry.key)); } @@ -887,10 +891,10 @@ private void verifyCanQueryState(final int cacheSizeBytes) throws Exception { waitUntilAtLeastNumRecordProcessed(outputTopic, 1); final ReadOnlyKeyValueStore - myCount = kafkaStreams.store(storeName, QueryableStoreTypes.keyValueStore()); + myCount = kafkaStreams.store(StoreQueryParams.fromNameAndType(storeName, QueryableStoreTypes.keyValueStore())); final ReadOnlyWindowStore windowStore = - kafkaStreams.store(windowStoreName, QueryableStoreTypes.windowStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(windowStoreName, QueryableStoreTypes.windowStore())); verifyCanGetByKey(keys, expectedCount, expectedCount, @@ -929,7 +933,7 @@ public void shouldNotMakeStoreAvailableUntilAllStoresAvailable() throws Exceptio "waiting for store " + storeName); final ReadOnlyKeyValueStore store = - kafkaStreams.store(storeName, QueryableStoreTypes.keyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(storeName, QueryableStoreTypes.keyValueStore())); TestUtils.waitForCondition( () -> Long.valueOf(8).equals(store.get("hello")), @@ -949,7 +953,7 @@ public void shouldNotMakeStoreAvailableUntilAllStoresAvailable() throws Exceptio try { assertEquals( Long.valueOf(8L), - kafkaStreams.store(storeName, QueryableStoreTypes.keyValueStore()).get("hello")); + kafkaStreams.store(StoreQueryParams.fromNameAndType(storeName, QueryableStoreTypes.keyValueStore())).get("hello")); return true; } catch (final InvalidStateStoreException ise) { return false; @@ -970,7 +974,7 @@ private class WaitForStore implements TestCondition { @Override public boolean conditionMet() { try { - kafkaStreams.store(storeName, QueryableStoreTypes.keyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(storeName, QueryableStoreTypes.keyValueStore())); return true; } catch (final InvalidStateStoreException ise) { return false; @@ -1028,7 +1032,7 @@ public void shouldAllowToQueryAfterThreadDied() throws Exception { "waiting for store " + storeName); final ReadOnlyKeyValueStore store = - kafkaStreams.store(storeName, QueryableStoreTypes.keyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(storeName, QueryableStoreTypes.keyValueStore())); TestUtils.waitForCondition( () -> "12".equals(store.get("a")) && "34".equals(store.get("b")), @@ -1055,7 +1059,7 @@ public void shouldAllowToQueryAfterThreadDied() throws Exception { "waiting for store " + storeName); final ReadOnlyKeyValueStore store2 = - kafkaStreams.store(storeName, QueryableStoreTypes.keyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(storeName, QueryableStoreTypes.keyValueStore())); try { TestUtils.waitForCondition( diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/StoreQueryIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/StoreQueryIntegrationTest.java new file mode 100644 index 0000000000000..c03998f543c12 --- /dev/null +++ b/streams/src/test/java/org/apache/kafka/streams/integration/StoreQueryIntegrationTest.java @@ -0,0 +1,344 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.streams.integration; + +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.common.serialization.IntegerSerializer; +import org.apache.kafka.common.serialization.Serdes; +import org.apache.kafka.common.utils.Bytes; +import org.apache.kafka.common.utils.MockTime; +import org.apache.kafka.streams.KafkaStreams; +import org.apache.kafka.streams.KeyQueryMetadata; +import org.apache.kafka.streams.StoreQueryParams; +import org.apache.kafka.streams.StreamsBuilder; +import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.KeyValue; +import org.apache.kafka.streams.errors.InvalidStateStoreException; +import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster; +import org.apache.kafka.streams.integration.utils.IntegrationTestUtils; +import org.apache.kafka.streams.kstream.Consumed; +import org.apache.kafka.streams.kstream.Materialized; +import org.apache.kafka.streams.state.KeyValueStore; +import org.apache.kafka.streams.state.QueryableStoreType; +import org.apache.kafka.streams.state.QueryableStoreTypes; +import org.apache.kafka.streams.state.ReadOnlyKeyValueStore; +import org.apache.kafka.test.IntegrationTest; +import org.apache.kafka.test.TestUtils; + +import org.junit.After; +import org.junit.Before; +import org.junit.Rule; +import org.junit.Test; +import org.junit.experimental.categories.Category; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.Properties; +import java.util.concurrent.Semaphore; +import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; +import java.util.stream.IntStream; + +import static org.apache.kafka.streams.integration.utils.IntegrationTestUtils.startApplicationAndWaitUntilRunning; +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.is; +import static org.hamcrest.Matchers.notNullValue; +import static org.hamcrest.Matchers.nullValue; + + +@Category({IntegrationTest.class}) +public class StoreQueryIntegrationTest { + + private static final int NUM_BROKERS = 1; + private static int port = 0; + private static final String INPUT_TOPIC_NAME = "input-topic"; + private static final String TABLE_NAME = "source-table"; + + @Rule + public final EmbeddedKafkaCluster cluster = new EmbeddedKafkaCluster(NUM_BROKERS); + + private final List streamsToCleanup = new ArrayList<>(); + private final MockTime mockTime = cluster.time; + + @Before + public void before() throws InterruptedException { + cluster.createTopic(INPUT_TOPIC_NAME, 2, 1); + } + + @After + public void after() { + for (final KafkaStreams kafkaStreams : streamsToCleanup) { + kafkaStreams.close(); + } + } + + @Test + public void shouldQueryAllActivePartitionStoresByDefault() throws Exception { + final int batch1NumMessages = 100; + final int key = 1; + final Semaphore semaphore = new Semaphore(0); + + final StreamsBuilder builder = new StreamsBuilder(); + builder.table(INPUT_TOPIC_NAME, Consumed.with(Serdes.Integer(), Serdes.Integer()), + Materialized.>as(TABLE_NAME) + .withCachingDisabled()) + .toStream() + .peek((k, v) -> semaphore.release()); + + final KafkaStreams kafkaStreams1 = createKafkaStreams(builder, streamsConfiguration()); + final KafkaStreams kafkaStreams2 = createKafkaStreams(builder, streamsConfiguration()); + final List kafkaStreamsList = Arrays.asList(kafkaStreams1, kafkaStreams2); + + startApplicationAndWaitUntilRunning(kafkaStreamsList, Duration.ofSeconds(60)); + + produceValueRange(key, 0, batch1NumMessages); + + // Assert that all messages in the first batch were processed in a timely manner + assertThat(semaphore.tryAcquire(batch1NumMessages, 60, TimeUnit.SECONDS), is(equalTo(true))); + final KeyQueryMetadata keyQueryMetadata = kafkaStreams1.queryMetadataForKey(TABLE_NAME, key, (topic, somekey, value, numPartitions) -> 0); + + final QueryableStoreType> queryableStoreType = QueryableStoreTypes.keyValueStore(); + final ReadOnlyKeyValueStore store1 = kafkaStreams1 + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, queryableStoreType)); + + final ReadOnlyKeyValueStore store2 = kafkaStreams2 + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, queryableStoreType)); + + final boolean kafkaStreams1IsActive; + if ((keyQueryMetadata.getActiveHost().port() % 2) == 1) { + kafkaStreams1IsActive = true; + } else { + kafkaStreams1IsActive = false; + } + + // Assert that only active is able to query for a key by default + assertThat(kafkaStreams1IsActive ? store1.get(key) : store2.get(key), is(notNullValue())); + assertThat(kafkaStreams1IsActive ? store2.get(key) : store1.get(key), is(nullValue())); + } + + @Test + public void shouldQuerySpecificActivePartitionStores() throws Exception { + final int batch1NumMessages = 100; + final int key = 1; + final Semaphore semaphore = new Semaphore(0); + + final StreamsBuilder builder = new StreamsBuilder(); + builder.table(INPUT_TOPIC_NAME, Consumed.with(Serdes.Integer(), Serdes.Integer()), + Materialized.>as(TABLE_NAME) + .withCachingDisabled()) + .toStream() + .peek((k, v) -> semaphore.release()); + + final KafkaStreams kafkaStreams1 = createKafkaStreams(builder, streamsConfiguration()); + final KafkaStreams kafkaStreams2 = createKafkaStreams(builder, streamsConfiguration()); + final List kafkaStreamsList = Arrays.asList(kafkaStreams1, kafkaStreams2); + + startApplicationAndWaitUntilRunning(kafkaStreamsList, Duration.ofSeconds(60)); + + produceValueRange(key, 0, batch1NumMessages); + + // Assert that all messages in the first batch were processed in a timely manner + assertThat(semaphore.tryAcquire(batch1NumMessages, 60, TimeUnit.SECONDS), is(equalTo(true))); + final KeyQueryMetadata keyQueryMetadata = kafkaStreams1.queryMetadataForKey(TABLE_NAME, key, (topic, somekey, value, numPartitions) -> 0); + + //key belongs to this partition + final int keyPartition = keyQueryMetadata.getPartition(); + + //key doesn't belongs to this partition + final int keyDontBelongPartition = (keyPartition == 0) ? 1 : 0; + final QueryableStoreType> queryableStoreType = QueryableStoreTypes.keyValueStore(); + ReadOnlyKeyValueStore store1 = null; + ReadOnlyKeyValueStore store2 = null; + try { + store1 = kafkaStreams1 + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, queryableStoreType).withPartition(keyPartition)); + } catch (final InvalidStateStoreException exception) { + //Only one among kafkaStreams1 and kafkaStreams2 will contain the specific active store requested. The other will throw exception + } + try { + store2 = kafkaStreams2 + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, queryableStoreType).withPartition(keyPartition)); + } catch (final InvalidStateStoreException exception) { + //Only one among kafkaStreams1 and kafkaStreams2 will contain the specific active store requested. The other will throw exception + } + final boolean kafkaStreams1IsActive; + if ((keyQueryMetadata.getActiveHost().port() % 2) == 1) { + kafkaStreams1IsActive = true; + assertThat(store1, is(notNullValue())); + assertThat(store2, is(nullValue())); + } else { + kafkaStreams1IsActive = false; + assertThat(store2, is(notNullValue())); + assertThat(store1, is(nullValue())); + } + + // Assert that only active for a specific requested partition serves key if stale stores and not enabled + assertThat(kafkaStreams1IsActive ? store1.get(key) : store2.get(key), is(notNullValue())); + + ReadOnlyKeyValueStore store3 = null; + ReadOnlyKeyValueStore store4 = null; + try { + store3 = kafkaStreams1 + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, queryableStoreType).withPartition(keyDontBelongPartition)); + } catch (final InvalidStateStoreException exception) { + //Only one among kafkaStreams1 and kafkaStreams2 will contain the specific active store requested. The other will throw exception + } + try { + store4 = kafkaStreams2 + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, queryableStoreType).withPartition(keyDontBelongPartition)); + } catch (final InvalidStateStoreException exception) { + //Only one among kafkaStreams1 and kafkaStreams2 will contain the specific active store requested. The other will throw exception + } + + // Assert that key is not served when wrong specific partition is requested + // If kafkaStreams1 is active for keyPartition, kafkaStreams2 would be active for keyDontBelongPartition + // So, in that case, store3 would be null and the store4 would not return the value for key as wrong partition was requested + assertThat(kafkaStreams1IsActive ? store4.get(key) : store3.get(key), is(nullValue())); + } + + @Test + public void shouldQueryAllStalePartitionStores() throws Exception { + final int batch1NumMessages = 100; + final int key = 1; + final Semaphore semaphore = new Semaphore(0); + + final StreamsBuilder builder = new StreamsBuilder(); + builder.table(INPUT_TOPIC_NAME, Consumed.with(Serdes.Integer(), Serdes.Integer()), + Materialized.>as(TABLE_NAME) + .withCachingDisabled()) + .toStream() + .peek((k, v) -> semaphore.release()); + + final KafkaStreams kafkaStreams1 = createKafkaStreams(builder, streamsConfiguration()); + final KafkaStreams kafkaStreams2 = createKafkaStreams(builder, streamsConfiguration()); + final List kafkaStreamsList = Arrays.asList(kafkaStreams1, kafkaStreams2); + + startApplicationAndWaitUntilRunning(kafkaStreamsList, Duration.ofSeconds(60)); + + produceValueRange(key, 0, batch1NumMessages); + + // Assert that all messages in the first batch were processed in a timely manner + assertThat(semaphore.tryAcquire(batch1NumMessages, 60, TimeUnit.SECONDS), is(equalTo(true))); + + final QueryableStoreType> queryableStoreType = QueryableStoreTypes.keyValueStore(); + final ReadOnlyKeyValueStore store1 = kafkaStreams1 + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, queryableStoreType).enableStaleStores()); + + final ReadOnlyKeyValueStore store2 = kafkaStreams2 + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, queryableStoreType).enableStaleStores()); + + // Assert that both active and standby are able to query for a key + assertThat(store1.get(key), is(notNullValue())); + assertThat(store2.get(key), is(notNullValue())); + } + + @Test + public void shouldQuerySpecificStalePartitionStores() throws Exception { + final int batch1NumMessages = 100; + final int key = 1; + final Semaphore semaphore = new Semaphore(0); + + final StreamsBuilder builder = new StreamsBuilder(); + builder.table(INPUT_TOPIC_NAME, Consumed.with(Serdes.Integer(), Serdes.Integer()), + Materialized.>as(TABLE_NAME) + .withCachingDisabled()) + .toStream() + .peek((k, v) -> semaphore.release()); + + final KafkaStreams kafkaStreams1 = createKafkaStreams(builder, streamsConfiguration()); + final KafkaStreams kafkaStreams2 = createKafkaStreams(builder, streamsConfiguration()); + final List kafkaStreamsList = Arrays.asList(kafkaStreams1, kafkaStreams2); + + startApplicationAndWaitUntilRunning(kafkaStreamsList, Duration.ofSeconds(60)); + + produceValueRange(key, 0, batch1NumMessages); + + // Assert that all messages in the first batch were processed in a timely manner + assertThat(semaphore.tryAcquire(batch1NumMessages, 60, TimeUnit.SECONDS), is(equalTo(true))); + final KeyQueryMetadata keyQueryMetadata = kafkaStreams1.queryMetadataForKey(TABLE_NAME, key, (topic, somekey, value, numPartitions) -> 0); + + //key belongs to this partition + final int keyPartition = keyQueryMetadata.getPartition(); + + //key doesn't belongs to this partition + final int keyDontBelongPartition = (keyPartition == 0) ? 1 : 0; + final QueryableStoreType> queryableStoreType = QueryableStoreTypes.keyValueStore(); + final ReadOnlyKeyValueStore store1 = kafkaStreams1 + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, queryableStoreType).enableStaleStores().withPartition(keyPartition)); + + final ReadOnlyKeyValueStore store2 = kafkaStreams2 + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, queryableStoreType).enableStaleStores().withPartition(keyPartition)); + + // Assert that both active and standby are able to query for a key + assertThat(store1.get(key), is(notNullValue())); + assertThat(store2.get(key), is(notNullValue())); + + + final ReadOnlyKeyValueStore store3 = kafkaStreams1 + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, queryableStoreType).enableStaleStores().withPartition(keyDontBelongPartition)); + + final ReadOnlyKeyValueStore store4 = kafkaStreams2 + .store(StoreQueryParams.fromNameAndType(TABLE_NAME, queryableStoreType).enableStaleStores().withPartition(keyDontBelongPartition)); + + // Assert that + assertThat(store3.get(key), is(nullValue())); + assertThat(store4.get(key), is(nullValue())); + } + + private KafkaStreams createKafkaStreams(final StreamsBuilder builder, final Properties config) { + final KafkaStreams streams = new KafkaStreams(builder.build(config), config); + streamsToCleanup.add(streams); + return streams; + } + + private void produceValueRange(final int key, final int start, final int endExclusive) throws Exception { + final Properties producerProps = new Properties(); + producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, cluster.bootstrapServers()); + producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class); + producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class); + + IntegrationTestUtils.produceKeyValuesSynchronously( + INPUT_TOPIC_NAME, + IntStream.range(start, endExclusive) + .mapToObj(i -> KeyValue.pair(key, i)) + .collect(Collectors.toList()), + producerProps, + mockTime); + } + + private Properties streamsConfiguration() { + final String applicationId = "streamsApp"; + final Properties config = new Properties(); + config.put(StreamsConfig.TOPOLOGY_OPTIMIZATION, StreamsConfig.OPTIMIZE); + config.put(StreamsConfig.APPLICATION_ID_CONFIG, applicationId); + config.put(StreamsConfig.APPLICATION_SERVER_CONFIG, "localhost:" + String.valueOf(++port)); + config.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, cluster.bootstrapServers()); + config.put(StreamsConfig.STATE_DIR_CONFIG, TestUtils.tempDirectory(applicationId).getPath()); + config.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.Integer().getClass()); + config.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.Integer().getClass()); + config.put(StreamsConfig.NUM_STANDBY_REPLICAS_CONFIG, 1); + config.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 100); + config.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 200); + config.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 1000); + config.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 100); + return config; + } +} diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/StoreUpgradeIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/StoreUpgradeIntegrationTest.java index 9a82ad9f89e81..3fa5dbcdb5eb2 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/StoreUpgradeIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/StoreUpgradeIntegrationTest.java @@ -23,6 +23,7 @@ import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster; import org.apache.kafka.streams.integration.utils.IntegrationTestUtils; import org.apache.kafka.streams.kstream.Windowed; @@ -333,7 +334,7 @@ private void processKeyValueAndVerifyPlainCount(final K key, () -> { try { final ReadOnlyKeyValueStore store = - kafkaStreams.store(STORE_NAME, QueryableStoreTypes.keyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(STORE_NAME, QueryableStoreTypes.keyValueStore())); try (final KeyValueIterator all = store.all()) { final List> storeContent = new LinkedList<>(); while (all.hasNext()) { @@ -358,7 +359,7 @@ private void verifyCountWithTimestamp(final K key, () -> { try { final ReadOnlyKeyValueStore> store = - kafkaStreams.store(STORE_NAME, QueryableStoreTypes.timestampedKeyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(STORE_NAME, QueryableStoreTypes.timestampedKeyValueStore())); final ValueAndTimestamp count = store.get(key); return count.value() == value && count.timestamp() == timestamp; } catch (final Exception swallow) { @@ -377,7 +378,7 @@ private void verifyCountWithSurrogateTimestamp(final K key, () -> { try { final ReadOnlyKeyValueStore> store = - kafkaStreams.store(STORE_NAME, QueryableStoreTypes.timestampedKeyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(STORE_NAME, QueryableStoreTypes.timestampedKeyValueStore())); final ValueAndTimestamp count = store.get(key); return count.value() == value && count.timestamp() == -1L; } catch (final Exception swallow) { @@ -407,7 +408,7 @@ private void processKeyValueAndVerifyCount(final K key, () -> { try { final ReadOnlyKeyValueStore> store = - kafkaStreams.store(STORE_NAME, QueryableStoreTypes.timestampedKeyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(STORE_NAME, QueryableStoreTypes.timestampedKeyValueStore())); try (final KeyValueIterator> all = store.all()) { final List>> storeContent = new LinkedList<>(); while (all.hasNext()) { @@ -442,7 +443,7 @@ private void processKeyValueAndVerifyCountWithTimestamp(final K key, () -> { try { final ReadOnlyKeyValueStore> store = - kafkaStreams.store(STORE_NAME, QueryableStoreTypes.timestampedKeyValueStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(STORE_NAME, QueryableStoreTypes.timestampedKeyValueStore())); try (final KeyValueIterator> all = store.all()) { final List>> storeContent = new LinkedList<>(); while (all.hasNext()) { @@ -808,7 +809,7 @@ private void processWindowedKeyValueAndVerifyPlainCount(final K key, () -> { try { final ReadOnlyWindowStore store = - kafkaStreams.store(STORE_NAME, QueryableStoreTypes.windowStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(STORE_NAME, QueryableStoreTypes.windowStore())); try (final KeyValueIterator, V> all = store.all()) { final List, V>> storeContent = new LinkedList<>(); while (all.hasNext()) { @@ -832,7 +833,7 @@ private void verifyWindowedCountWithSurrogateTimestamp(final Windowed key () -> { try { final ReadOnlyWindowStore> store = - kafkaStreams.store(STORE_NAME, QueryableStoreTypes.timestampedWindowStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(STORE_NAME, QueryableStoreTypes.timestampedWindowStore())); final ValueAndTimestamp count = store.fetch(key.key(), key.window().start()); return count.value() == value && count.timestamp() == -1L; } catch (final Exception swallow) { @@ -852,7 +853,7 @@ private void verifyWindowedCountWithTimestamp(final Windowed key, () -> { try { final ReadOnlyWindowStore> store = - kafkaStreams.store(STORE_NAME, QueryableStoreTypes.timestampedWindowStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(STORE_NAME, QueryableStoreTypes.timestampedWindowStore())); final ValueAndTimestamp count = store.fetch(key.key(), key.window().start()); return count.value() == value && count.timestamp() == timestamp; } catch (final Exception swallow) { @@ -882,7 +883,7 @@ private void processKeyValueAndVerifyWindowedCountWithTimestamp(final K k () -> { try { final ReadOnlyWindowStore> store = - kafkaStreams.store(STORE_NAME, QueryableStoreTypes.timestampedWindowStore()); + kafkaStreams.store(StoreQueryParams.fromNameAndType(STORE_NAME, QueryableStoreTypes.timestampedWindowStore())); try (final KeyValueIterator, ValueAndTimestamp> all = store.all()) { final List, ValueAndTimestamp>> storeContent = new LinkedList<>(); while (all.hasNext()) { diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/StreamStreamJoinIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/StreamStreamJoinIntegrationTest.java index 4094fde1cce2b..f69a7b0da2a18 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/StreamStreamJoinIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/StreamStreamJoinIntegrationTest.java @@ -19,6 +19,7 @@ import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.KafkaStreams.State; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.errors.InvalidStateStoreException; @@ -93,7 +94,7 @@ public void shouldNotAccessJoinStoresWhenGivingName() throws InterruptedExceptio kafkaStreams.start(); latch.await(); - assertThrows(InvalidStateStoreException.class, () -> kafkaStreams.store("join-store", QueryableStoreTypes.keyValueStore())); + assertThrows(InvalidStateStoreException.class, () -> kafkaStreams.store(StoreQueryParams.fromNameAndType("join-store", QueryableStoreTypes.keyValueStore()))); } } diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/CompositeReadOnlyKeyValueStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/CompositeReadOnlyKeyValueStoreTest.java index f8af40ba63d48..b16f3bbd569c3 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/CompositeReadOnlyKeyValueStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/CompositeReadOnlyKeyValueStoreTest.java @@ -18,13 +18,16 @@ import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.KeyValue; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.processor.internals.ProcessorStateManager; -import org.apache.kafka.streams.state.KeyValueIterator; import org.apache.kafka.streams.state.KeyValueStore; -import org.apache.kafka.streams.state.QueryableStoreTypes; -import org.apache.kafka.streams.state.StateSerdes; +import org.apache.kafka.streams.state.QueryableStoreType; import org.apache.kafka.streams.state.Stores; +import org.apache.kafka.streams.state.StateSerdes; +import org.apache.kafka.streams.state.KeyValueIterator; +import org.apache.kafka.streams.state.QueryableStoreTypes; +import org.apache.kafka.streams.state.ReadOnlyKeyValueStore; import org.apache.kafka.test.InternalMockProcessorContext; import org.apache.kafka.test.NoOpReadOnlyStore; import org.apache.kafka.test.MockRecordCollector; @@ -61,9 +64,9 @@ public void before() { stubProviderOne.addStore(storeName, stubOneUnderlying); otherUnderlyingStore = newStoreInstance(); stubProviderOne.addStore("other-store", otherUnderlyingStore); - + final QueryableStoreType> queryableStoreType = QueryableStoreTypes.keyValueStore(); theStore = new CompositeReadOnlyKeyValueStore<>( - new WrappingStoreProvider(asList(stubProviderOne, stubProviderTwo), false), + new WrappingStoreProvider(asList(stubProviderOne, stubProviderTwo), StoreQueryParams.fromNameAndType(storeName, QueryableStoreTypes.keyValueStore())), QueryableStoreTypes.keyValueStore(), storeName ); @@ -294,8 +297,9 @@ public long approximateNumEntries() { } private CompositeReadOnlyKeyValueStore rebalancing() { + final QueryableStoreType> queryableStoreType = QueryableStoreTypes.keyValueStore(); return new CompositeReadOnlyKeyValueStore<>( - new WrappingStoreProvider(singletonList(new StateStoreProviderStub(true)), false), + new WrappingStoreProvider(singletonList(new StateStoreProviderStub(true)), StoreQueryParams.fromNameAndType(storeName, QueryableStoreTypes.keyValueStore())), QueryableStoreTypes.keyValueStore(), storeName ); diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/CompositeReadOnlySessionStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/CompositeReadOnlySessionStoreTest.java index 419b7b339cb7e..5e309e7d62072 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/CompositeReadOnlySessionStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/CompositeReadOnlySessionStoreTest.java @@ -17,10 +17,13 @@ package org.apache.kafka.streams.state.internals; import org.apache.kafka.streams.KeyValue; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.kstream.Windowed; import org.apache.kafka.streams.kstream.internals.SessionWindow; +import org.apache.kafka.streams.state.ReadOnlySessionStore; import org.apache.kafka.streams.state.KeyValueIterator; +import org.apache.kafka.streams.state.QueryableStoreType; import org.apache.kafka.streams.state.QueryableStoreTypes; import org.apache.kafka.test.ReadOnlySessionStoreStub; import org.apache.kafka.test.StateStoreProviderStub; @@ -52,10 +55,10 @@ public class CompositeReadOnlySessionStoreTest { public void before() { stubProviderOne.addStore(storeName, underlyingSessionStore); stubProviderOne.addStore("other-session-store", otherUnderlyingStore); - + final QueryableStoreType> queryableStoreType = QueryableStoreTypes.sessionStore(); sessionStore = new CompositeReadOnlySessionStore<>( - new WrappingStoreProvider(Arrays.asList(stubProviderOne, stubProviderTwo), false), + new WrappingStoreProvider(Arrays.asList(stubProviderOne, stubProviderTwo), StoreQueryParams.fromNameAndType(storeName, queryableStoreType)), QueryableStoreTypes.sessionStore(), storeName); } @@ -107,9 +110,10 @@ public void shouldNotGetValueFromOtherStores() { @Test(expected = InvalidStateStoreException.class) public void shouldThrowInvalidStateStoreExceptionOnRebalance() { + final QueryableStoreType> queryableStoreType = QueryableStoreTypes.sessionStore(); final CompositeReadOnlySessionStore store = new CompositeReadOnlySessionStore<>( - new WrappingStoreProvider(singletonList(new StateStoreProviderStub(true)), false), + new WrappingStoreProvider(singletonList(new StateStoreProviderStub(true)), StoreQueryParams.fromNameAndType("whateva", queryableStoreType)), QueryableStoreTypes.sessionStore(), "whateva" ); diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/CompositeReadOnlyWindowStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/CompositeReadOnlyWindowStoreTest.java index 6495d8070916d..8aeb3a09b6348 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/CompositeReadOnlyWindowStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/CompositeReadOnlyWindowStoreTest.java @@ -17,6 +17,7 @@ package org.apache.kafka.streams.state.internals; import org.apache.kafka.streams.KeyValue; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.kstream.Windowed; import org.apache.kafka.streams.kstream.internals.TimeWindow; @@ -66,9 +67,8 @@ public void before() { otherUnderlyingStore = new ReadOnlyWindowStoreStub<>(WINDOW_SIZE); stubProviderOne.addStore("other-window-store", otherUnderlyingStore); - windowStore = new CompositeReadOnlyWindowStore<>( - new WrappingStoreProvider(asList(stubProviderOne, stubProviderTwo), false), + new WrappingStoreProvider(asList(stubProviderOne, stubProviderTwo), StoreQueryParams.fromNameAndType(storeName, QueryableStoreTypes.windowStore())), QueryableStoreTypes.windowStore(), storeName ); @@ -140,7 +140,7 @@ public void shouldThrowInvalidStateStoreExceptionIfFetchThrows() { underlyingWindowStore.setOpen(false); final CompositeReadOnlyWindowStore store = new CompositeReadOnlyWindowStore<>( - new WrappingStoreProvider(singletonList(stubProviderOne), false), + new WrappingStoreProvider(singletonList(stubProviderOne), StoreQueryParams.fromNameAndType("window-store", QueryableStoreTypes.windowStore())), QueryableStoreTypes.windowStore(), "window-store" ); diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/QueryableStoreProviderTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/QueryableStoreProviderTest.java index de202f0cec9ab..0da07148bd855 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/QueryableStoreProviderTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/QueryableStoreProviderTest.java @@ -17,6 +17,7 @@ package org.apache.kafka.streams.state.internals; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.processor.StateStore; import org.apache.kafka.streams.state.NoOpWindowStore; @@ -53,38 +54,38 @@ public void before() { @Test(expected = InvalidStateStoreException.class) public void shouldThrowExceptionIfKVStoreDoesntExist() { - storeProvider.getStore("not-a-store", QueryableStoreTypes.keyValueStore(), false); + storeProvider.getStore(StoreQueryParams.fromNameAndType("not-a-store", QueryableStoreTypes.keyValueStore())); } @Test(expected = InvalidStateStoreException.class) public void shouldThrowExceptionIfWindowStoreDoesntExist() { - storeProvider.getStore("not-a-store", QueryableStoreTypes.windowStore(), false); + storeProvider.getStore(StoreQueryParams.fromNameAndType("not-a-store", QueryableStoreTypes.windowStore())); } @Test public void shouldReturnKVStoreWhenItExists() { - assertNotNull(storeProvider.getStore(keyValueStore, QueryableStoreTypes.keyValueStore(), false)); + assertNotNull(storeProvider.getStore(StoreQueryParams.fromNameAndType(keyValueStore, QueryableStoreTypes.keyValueStore()))); } @Test public void shouldReturnWindowStoreWhenItExists() { - assertNotNull(storeProvider.getStore(windowStore, QueryableStoreTypes.windowStore(), false)); + assertNotNull(storeProvider.getStore(StoreQueryParams.fromNameAndType(windowStore, QueryableStoreTypes.windowStore()))); } @Test(expected = InvalidStateStoreException.class) public void shouldThrowExceptionWhenLookingForWindowStoreWithDifferentType() { - storeProvider.getStore(windowStore, QueryableStoreTypes.keyValueStore(), false); + storeProvider.getStore(StoreQueryParams.fromNameAndType(windowStore, QueryableStoreTypes.keyValueStore())); } @Test(expected = InvalidStateStoreException.class) public void shouldThrowExceptionWhenLookingForKVStoreWithDifferentType() { - storeProvider.getStore(keyValueStore, QueryableStoreTypes.windowStore(), false); + storeProvider.getStore(StoreQueryParams.fromNameAndType(keyValueStore, QueryableStoreTypes.windowStore())); } @Test public void shouldFindGlobalStores() { globalStateStores.put("global", new NoOpReadOnlyStore<>()); - assertNotNull(storeProvider.getStore("global", QueryableStoreTypes.keyValueStore(), false)); + assertNotNull(storeProvider.getStore(StoreQueryParams.fromNameAndType("global", QueryableStoreTypes.keyValueStore()))); } diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProviderTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProviderTest.java index 48b1625678908..b2f5f4e186b82 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProviderTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProviderTest.java @@ -24,14 +24,16 @@ import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.common.utils.MockTime; import org.apache.kafka.common.utils.Utils; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.TopologyWrapper; import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.processor.TaskId; -import org.apache.kafka.streams.processor.internals.MockStreamsMetrics; import org.apache.kafka.streams.processor.internals.ProcessorTopology; -import org.apache.kafka.streams.processor.internals.StateDirectory; +import org.apache.kafka.streams.processor.internals.MockStreamsMetrics; import org.apache.kafka.streams.processor.internals.StoreChangelogReader; +import org.apache.kafka.streams.processor.internals.InternalTopologyBuilder; +import org.apache.kafka.streams.processor.internals.StateDirectory; import org.apache.kafka.streams.processor.internals.StreamTask; import org.apache.kafka.streams.processor.internals.StreamThread; import org.apache.kafka.streams.state.QueryableStoreTypes; @@ -125,7 +127,8 @@ public void before() { configureRestoreConsumer(clientSupplier, "applicationId-kv-store-changelog"); configureRestoreConsumer(clientSupplier, "applicationId-window-store-changelog"); - final ProcessorTopology processorTopology = topology.getInternalBuilder(applicationId).build(); + final InternalTopologyBuilder internalTopologyBuilder = topology.getInternalBuilder(applicationId); + final ProcessorTopology processorTopology = internalTopologyBuilder.build(); tasks = new HashMap<>(); stateDirectory = new StateDirectory(streamsConfig, new MockTime(), true); @@ -147,7 +150,7 @@ public void before() { tasks.put(new TaskId(0, 1), taskTwo); threadMock = EasyMock.createNiceMock(StreamThread.class); - provider = new StreamThreadStateStoreProvider(threadMock); + provider = new StreamThreadStateStoreProvider(threadMock, internalTopologyBuilder); } @@ -160,7 +163,7 @@ public void cleanUp() throws IOException { public void shouldFindKeyValueStores() { mockThread(true); final List> kvStores = - provider.stores("kv-store", QueryableStoreTypes.keyValueStore(), false); + provider.stores(StoreQueryParams.fromNameAndType("kv-store", QueryableStoreTypes.keyValueStore())); assertEquals(2, kvStores.size()); for (final ReadOnlyKeyValueStore store: kvStores) { assertThat(store, instanceOf(ReadOnlyKeyValueStore.class)); @@ -172,7 +175,7 @@ public void shouldFindKeyValueStores() { public void shouldFindTimestampedKeyValueStores() { mockThread(true); final List>> tkvStores = - provider.stores("timestamped-kv-store", QueryableStoreTypes.timestampedKeyValueStore(), false); + provider.stores(StoreQueryParams.fromNameAndType("timestamped-kv-store", QueryableStoreTypes.timestampedKeyValueStore())); assertEquals(2, tkvStores.size()); for (final ReadOnlyKeyValueStore> store: tkvStores) { assertThat(store, instanceOf(ReadOnlyKeyValueStore.class)); @@ -184,7 +187,7 @@ public void shouldFindTimestampedKeyValueStores() { public void shouldNotFindKeyValueStoresAsTimestampedStore() { mockThread(true); final List>> tkvStores = - provider.stores("kv-store", QueryableStoreTypes.timestampedKeyValueStore(), false); + provider.stores(StoreQueryParams.fromNameAndType("kv-store", QueryableStoreTypes.timestampedKeyValueStore())); assertEquals(0, tkvStores.size()); } @@ -192,7 +195,7 @@ public void shouldNotFindKeyValueStoresAsTimestampedStore() { public void shouldFindTimestampedKeyValueStoresAsKeyValueStores() { mockThread(true); final List>> tkvStores = - provider.stores("timestamped-kv-store", QueryableStoreTypes.keyValueStore(), false); + provider.stores(StoreQueryParams.fromNameAndType("timestamped-kv-store", QueryableStoreTypes.keyValueStore())); assertEquals(2, tkvStores.size()); for (final ReadOnlyKeyValueStore> store: tkvStores) { assertThat(store, instanceOf(ReadOnlyKeyValueStore.class)); @@ -204,7 +207,7 @@ public void shouldFindTimestampedKeyValueStoresAsKeyValueStores() { public void shouldFindWindowStores() { mockThread(true); final List> windowStores = - provider.stores("window-store", QueryableStoreTypes.windowStore(), false); + provider.stores(StoreQueryParams.fromNameAndType("window-store", QueryableStoreTypes.windowStore())); assertEquals(2, windowStores.size()); for (final ReadOnlyWindowStore store: windowStores) { assertThat(store, instanceOf(ReadOnlyWindowStore.class)); @@ -216,7 +219,7 @@ public void shouldFindWindowStores() { public void shouldFindTimestampedWindowStores() { mockThread(true); final List>> windowStores = - provider.stores("timestamped-window-store", QueryableStoreTypes.timestampedWindowStore(), false); + provider.stores(StoreQueryParams.fromNameAndType("timestamped-window-store", QueryableStoreTypes.timestampedWindowStore())); assertEquals(2, windowStores.size()); for (final ReadOnlyWindowStore> store: windowStores) { assertThat(store, instanceOf(ReadOnlyWindowStore.class)); @@ -228,7 +231,7 @@ public void shouldFindTimestampedWindowStores() { public void shouldNotFindWindowStoresAsTimestampedStore() { mockThread(true); final List>> windowStores = - provider.stores("window-store", QueryableStoreTypes.timestampedWindowStore(), false); + provider.stores(StoreQueryParams.fromNameAndType("window-store", QueryableStoreTypes.timestampedWindowStore())); assertEquals(0, windowStores.size()); } @@ -236,7 +239,7 @@ public void shouldNotFindWindowStoresAsTimestampedStore() { public void shouldFindTimestampedWindowStoresAsWindowStore() { mockThread(true); final List>> windowStores = - provider.stores("timestamped-window-store", QueryableStoreTypes.windowStore(), false); + provider.stores(StoreQueryParams.fromNameAndType("timestamped-window-store", QueryableStoreTypes.windowStore())); assertEquals(2, windowStores.size()); for (final ReadOnlyWindowStore> store: windowStores) { assertThat(store, instanceOf(ReadOnlyWindowStore.class)); @@ -248,28 +251,28 @@ public void shouldFindTimestampedWindowStoresAsWindowStore() { public void shouldThrowInvalidStoreExceptionIfKVStoreClosed() { mockThread(true); taskOne.getStore("kv-store").close(); - provider.stores("kv-store", QueryableStoreTypes.keyValueStore(), false); + provider.stores(StoreQueryParams.fromNameAndType("kv-store", QueryableStoreTypes.keyValueStore())); } @Test(expected = InvalidStateStoreException.class) public void shouldThrowInvalidStoreExceptionIfTsKVStoreClosed() { mockThread(true); taskOne.getStore("timestamped-kv-store").close(); - provider.stores("timestamped-kv-store", QueryableStoreTypes.timestampedKeyValueStore(), false); + provider.stores(StoreQueryParams.fromNameAndType("timestamped-kv-store", QueryableStoreTypes.timestampedKeyValueStore())); } @Test(expected = InvalidStateStoreException.class) public void shouldThrowInvalidStoreExceptionIfWindowStoreClosed() { mockThread(true); taskOne.getStore("window-store").close(); - provider.stores("window-store", QueryableStoreTypes.windowStore(), false); + provider.stores(StoreQueryParams.fromNameAndType("window-store", QueryableStoreTypes.windowStore())); } @Test(expected = InvalidStateStoreException.class) public void shouldThrowInvalidStoreExceptionIfTsWindowStoreClosed() { mockThread(true); taskOne.getStore("timestamped-window-store").close(); - provider.stores("timestamped-window-store", QueryableStoreTypes.timestampedWindowStore(), false); + provider.stores(StoreQueryParams.fromNameAndType("timestamped-window-store", QueryableStoreTypes.timestampedWindowStore())); } @Test @@ -277,7 +280,7 @@ public void shouldReturnEmptyListIfNoStoresFoundWithName() { mockThread(true); assertEquals( Collections.emptyList(), - provider.stores("not-a-store", QueryableStoreTypes.keyValueStore(), false)); + provider.stores(StoreQueryParams.fromNameAndType("not-a-store", QueryableStoreTypes.keyValueStore()))); } @Test @@ -285,14 +288,14 @@ public void shouldReturnEmptyListIfStoreExistsButIsNotOfTypeValueStore() { mockThread(true); assertEquals( Collections.emptyList(), - provider.stores("window-store", QueryableStoreTypes.keyValueStore(), false) + provider.stores(StoreQueryParams.fromNameAndType("window-store", QueryableStoreTypes.keyValueStore())) ); } @Test(expected = InvalidStateStoreException.class) public void shouldThrowInvalidStoreExceptionIfNotAllStoresAvailable() { mockThread(false); - provider.stores("kv-store", QueryableStoreTypes.keyValueStore(), false); + provider.stores(StoreQueryParams.fromNameAndType("kv-store", QueryableStoreTypes.keyValueStore())); } private StreamTask createStreamsTask(final StreamsConfig streamsConfig, diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/WrappingStoreProviderTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/WrappingStoreProviderTest.java index 651d185eb7e34..7c66509c8592c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/WrappingStoreProviderTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/WrappingStoreProviderTest.java @@ -18,12 +18,13 @@ import org.apache.kafka.common.serialization.Serdes; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.state.NoOpWindowStore; import org.apache.kafka.streams.state.QueryableStoreTypes; import org.apache.kafka.streams.state.ReadOnlyKeyValueStore; -import org.apache.kafka.streams.state.ReadOnlyWindowStore; import org.apache.kafka.streams.state.Stores; +import org.apache.kafka.streams.state.ReadOnlyWindowStore; import org.apache.kafka.test.StateStoreProviderStub; import org.junit.Before; import org.junit.Test; @@ -54,10 +55,9 @@ public void before() { Serdes.serdeFrom(String.class)) .build()); stubProviderTwo.addStore("window", new NoOpWindowStore()); - wrappingStoreProvider = new WrappingStoreProvider( Arrays.asList(stubProviderOne, stubProviderTwo), - false + StoreQueryParams.fromNameAndType("kv", QueryableStoreTypes.keyValueStore()) ); } @@ -70,6 +70,7 @@ public void shouldFindKeyValueStores() { @Test public void shouldFindWindowStores() { + wrappingStoreProvider.setStoreQueryParams(StoreQueryParams.fromNameAndType("window", windowStore())); final List> windowStores = wrappingStoreProvider.stores("window", windowStore()); @@ -78,6 +79,7 @@ public void shouldFindWindowStores() { @Test(expected = InvalidStateStoreException.class) public void shouldThrowInvalidStoreExceptionIfNoStoreOfTypeFound() { - wrappingStoreProvider.stores("doesn't exist", QueryableStoreTypes.keyValueStore()); + wrappingStoreProvider.setStoreQueryParams(StoreQueryParams.fromNameAndType("doesn't exist", QueryableStoreTypes.keyValueStore())); + wrappingStoreProvider.stores("doesn't exist", QueryableStoreTypes.keyValueStore()); } } \ No newline at end of file diff --git a/streams/src/test/java/org/apache/kafka/test/StateStoreProviderStub.java b/streams/src/test/java/org/apache/kafka/test/StateStoreProviderStub.java index 09a6829386e3b..28835d21c18d8 100644 --- a/streams/src/test/java/org/apache/kafka/test/StateStoreProviderStub.java +++ b/streams/src/test/java/org/apache/kafka/test/StateStoreProviderStub.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.test; +import org.apache.kafka.streams.StoreQueryParams; import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.processor.StateStore; import org.apache.kafka.streams.state.QueryableStoreType; @@ -32,15 +33,15 @@ public class StateStoreProviderStub extends StreamThreadStateStoreProvider { private final boolean throwException; public StateStoreProviderStub(final boolean throwException) { - super(null); + super(null, null); this.throwException = throwException; } @SuppressWarnings("unchecked") @Override - public List stores(final String storeName, - final QueryableStoreType queryableStoreType, - final boolean includeStaleStores) { + public List stores(final StoreQueryParams storeQueryParams) { + final String storeName = storeQueryParams.storeName(); + final QueryableStoreType queryableStoreType = storeQueryParams.queryableStoreType(); if (throwException) { throw new InvalidStateStoreException("store is unavailable"); }