I'm new to kafka streams and faced a problem with Kafka Streams Store. It looks like store is not accesible after rebalancing. Here's my sample code:
@Bean
public KafkaStreams kafkaStreams(StreamsBuilder streamsBuilder, KafkaStreamsConfiguration kafkaStreamsConfiguration) {
KStream<String, String> inputTopic = streamsBuilder.stream("test-topic");
inputTopic.filter(
(key, value) -> value.contains("test")
).toTable(Materialized.<String, String, KeyValueStore<Bytes, byte[]>> as("test-store"));
KafkaStreams streams = new KafkaStreams(streamsBuilder.build(), kafkaStreamsConfiguration.asProperties());
streams.cleanUp();
streams.start();
return streams;
}
@Component
@AllArgsConstructor
public class TestKafkaStore {
private final KafkaStreams kafkaStreams;
public void readKafka() {
final ReadOnlyKeyValueStore<String, String> keyValueStore =
kafkaStreams.store(StoreQueryParameters.fromNameAndType("test-store", QueryableStoreTypes.keyValueStore()));
}
}
I'm able to access the store and read values from keyValueStore, but after rebalance
2021-12-09 00:16:52,968 INFO org.apache.kafka.streams.KafkaStreams [test-stream-a8da1f15-1742-4256-8515-f36de8eae941-StreamThread-1] [] stream-client [test-stream-a8da1f15-1742-4256-8515-f36de8eae941] State transition from RUNNING to REBALANCING
2021-12-09 00:16:53,032 INFO org.apache.kafka.streams.KafkaStreams [test-stream-621d13d3-78cf-413e-b6cc-3e46f24bf731-StreamThread-1] [] stream-client [test-stream-621d13d3-78cf-413e-b6cc-3e46f24bf731] State transition from REBALANCING to RUNNING
when I'm trying to run readKafka() again i'm getting error:
org.apache.kafka.streams.errors.InvalidStateStoreException: The state store, test-store, may have migrated to another instance.
at org.apache.kafka.streams.state.internals.QueryableStoreProvider.getStore(QueryableStoreProvider.java:76)
at org.apache.kafka.streams.KafkaStreams.store(KafkaStreams.java:1202)