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 83dd8c9ca09e5..57d16eefdb955 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 @@ -59,7 +59,14 @@ public List stores(final StoreQueryParameters storeQueryParams) { final Map tasks = storeQueryParams.staleStoresEnabled() ? streamThread.allTasks() : streamThread.activeTaskMap(); final List stores = new ArrayList<>(); if (keyTaskId != null) { - final T store = validateAndListStores(tasks.get(keyTaskId).getStore(storeName), queryableStoreType, storeName, keyTaskId); + final Task task = tasks.get(keyTaskId); + if (task == null) { + throw new InvalidStateStoreException( + String.format("The specified partition %d for store %s does not exist.", + storeQueryParams.partition(), + storeName)); + } + final T store = validateAndListStores(task.getStore(storeName), queryableStoreType, storeName, keyTaskId); if (store != null) { return Collections.singletonList(store); } 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 aa14a2e3372b1..911d2e35afa3b 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 @@ -72,7 +72,9 @@ import static org.hamcrest.MatcherAssert.assertThat; import static org.hamcrest.Matchers.instanceOf; import static org.hamcrest.Matchers.not; +import static org.hamcrest.core.IsEqual.equalTo; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertThrows; public class StreamThreadStateStoreProviderTest { @@ -290,6 +292,48 @@ public void shouldReturnEmptyListIfNoStoresFoundWithName() { provider.stores(StoreQueryParameters.fromNameAndType("not-a-store", QueryableStoreTypes.keyValueStore()))); } + @Test + public void shouldReturnSingleStoreForPartition() { + mockThread(true); + { + final List> kvStores = + provider.stores( + StoreQueryParameters + .fromNameAndType("kv-store", QueryableStoreTypes.keyValueStore()) + .withPartition(0)); + assertEquals(1, kvStores.size()); + for (final ReadOnlyKeyValueStore store : kvStores) { + assertThat(store, instanceOf(ReadOnlyKeyValueStore.class)); + assertThat(store, not(instanceOf(TimestampedKeyValueStore.class))); + } + } + { + final List> kvStores = + provider.stores( + StoreQueryParameters + .fromNameAndType("kv-store", QueryableStoreTypes.keyValueStore()) + .withPartition(1)); + assertEquals(1, kvStores.size()); + for (final ReadOnlyKeyValueStore store : kvStores) { + assertThat(store, instanceOf(ReadOnlyKeyValueStore.class)); + assertThat(store, not(instanceOf(TimestampedKeyValueStore.class))); + } + } + } + + @Test + public void shouldThrowForInvalidPartitions() { + mockThread(true); + final InvalidStateStoreException thrown = assertThrows( + InvalidStateStoreException.class, + () -> provider.stores( + StoreQueryParameters + .fromNameAndType("kv-store", QueryableStoreTypes.keyValueStore()) + .withPartition(2)) + ); + assertThat(thrown.getMessage(), equalTo("The specified partition 2 for store kv-store does not exist.")); + } + @Test public void shouldReturnEmptyListIfStoreExistsButIsNotOfTypeValueStore() { mockThread(true);