From 3e4af104b85ab56f2000bb1798291dbcbe8b5c00 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 23 Aug 2019 17:28:49 -0400 Subject: [PATCH] Introducing retry for finding state stores. During application startup, there are cases when state stores might take longer time to initialize. The call to find the state store in the binder (InteractiveQueryService) will fail in those situations. Providing a basic retry mechanism using which applicaitons can opt-in for this retry. Adding properties for retrying state store at the binder level. Adding docs. Polishing. Resolves #706 --- docs/src/main/asciidoc/kafka-streams.adoc | 10 +++++ .../streams/InteractiveQueryService.java | 42 ++++++++++++++----- ...aStreamsBinderConfigurationProperties.java | 33 +++++++++++++++ ...reamsInteractiveQueryIntegrationTests.java | 32 +++++++++++++- 4 files changed, 105 insertions(+), 12 deletions(-) rename spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/{integration => }/KafkaStreamsInteractiveQueryIntegrationTests.java (84%) diff --git a/docs/src/main/asciidoc/kafka-streams.adoc b/docs/src/main/asciidoc/kafka-streams.adoc index b00293705..701f029f7 100644 --- a/docs/src/main/asciidoc/kafka-streams.adoc +++ b/docs/src/main/asciidoc/kafka-streams.adoc @@ -952,6 +952,16 @@ ReadOnlyKeyValueStore keyValueStore = interactiveQueryService.getQueryableStoreType("my-store", QueryableStoreTypes.keyValueStore()); ---- +During the startup, the above method call to retrieve the startup might fail. +For e.g it might still be in the middle of initializing the state store. +In such cases, it will be useful to retry this operation. +Kafka Streams binder provides a simple retry mechanism to accommodate this. + +Following are the two properties that you can use to control this retrying. + +* spring.cloud.stream.binder.kafka.streams.stateStoreRetry.maxAttempts - Default is `1` . +* spring.cloud.stream.binder.kafka.streams.stateStoreRetry.backOffInterval - Default is `1000` milliseconds. + If there are multiple instances of the kafka streams application running, then before you can query them interactively, you need to identify which application instance hosts the key. `InteractiveQueryService` API provides methods for identifying the host information. diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/InteractiveQueryService.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/InteractiveQueryService.java index 055370e08..3dc7e201a 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/InteractiveQueryService.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/InteractiveQueryService.java @@ -19,6 +19,8 @@ package org.springframework.cloud.stream.binder.kafka.streams; import java.util.Map; import java.util.Optional; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; import org.apache.kafka.common.serialization.Serializer; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.errors.InvalidStateStoreException; @@ -27,6 +29,10 @@ import org.apache.kafka.streams.state.QueryableStoreType; import org.apache.kafka.streams.state.StreamsMetadata; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; +import org.springframework.retry.RetryPolicy; +import org.springframework.retry.backoff.FixedBackOffPolicy; +import org.springframework.retry.policy.SimpleRetryPolicy; +import org.springframework.retry.support.RetryTemplate; import org.springframework.util.StringUtils; /** @@ -41,6 +47,8 @@ import org.springframework.util.StringUtils; */ public class InteractiveQueryService { + private static final Log LOG = LogFactory.getLog(InteractiveQueryService.class); + private final KafkaStreamsRegistry kafkaStreamsRegistry; private final KafkaStreamsBinderConfigurationProperties binderConfigurationProperties; @@ -64,18 +72,32 @@ public class InteractiveQueryService { * @return queryable store. */ public T getQueryableStore(String storeName, QueryableStoreType storeType) { - for (KafkaStreams kafkaStream : this.kafkaStreamsRegistry.getKafkaStreams()) { - try { - T store = kafkaStream.store(storeName, storeType); - if (store != null) { - return store; + + RetryTemplate retryTemplate = new RetryTemplate(); + + KafkaStreamsBinderConfigurationProperties.StateStoreRetry stateStoreRetry = this.binderConfigurationProperties.getStateStoreRetry(); + RetryPolicy retryPolicy = new SimpleRetryPolicy(stateStoreRetry.getMaxAttempts()); + FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy(); + backOffPolicy.setBackOffPeriod(stateStoreRetry.getBackoffPeriod()); + + retryTemplate.setBackOffPolicy(backOffPolicy); + retryTemplate.setRetryPolicy(retryPolicy); + + return retryTemplate.execute(context -> { + T store; + for (KafkaStreams kafkaStream : InteractiveQueryService.this.kafkaStreamsRegistry.getKafkaStreams()) { + try { + store = kafkaStream.store(storeName, storeType); + if (store != null) { + return store; + } + } + catch (InvalidStateStoreException e) { + LOG.warn("Error when retrieving state store: " + storeName, e); } } - catch (InvalidStateStoreException ignored) { - // pass through - } - } - return null; + throw new IllegalStateException("Error when retrieving state store: " + storeName); + }); } /** diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsBinderConfigurationProperties.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsBinderConfigurationProperties.java index 747157a54..0f2f577d7 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsBinderConfigurationProperties.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsBinderConfigurationProperties.java @@ -54,6 +54,16 @@ public class KafkaStreamsBinderConfigurationProperties private String applicationId; + private StateStoreRetry stateStoreRetry = new StateStoreRetry(); + + public StateStoreRetry getStateStoreRetry() { + return stateStoreRetry; + } + + public void setStateStoreRetry(StateStoreRetry stateStoreRetry) { + this.stateStoreRetry = stateStoreRetry; + } + public String getApplicationId() { return this.applicationId; } @@ -79,4 +89,27 @@ public class KafkaStreamsBinderConfigurationProperties this.serdeError = serdeError; } + public static class StateStoreRetry { + + private int maxAttempts = 1; + + private long backoffPeriod = 1000; + + public int getMaxAttempts() { + return maxAttempts; + } + + public void setMaxAttempts(int maxAttempts) { + this.maxAttempts = maxAttempts; + } + + public long getBackoffPeriod() { + return backoffPeriod; + } + + public void setBackoffPeriod(long backoffPeriod) { + this.backoffPeriod = backoffPeriod; + } + } + } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsInteractiveQueryIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsInteractiveQueryIntegrationTests.java similarity index 84% rename from spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsInteractiveQueryIntegrationTests.java rename to spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsInteractiveQueryIntegrationTests.java index 171eb7f1c..c0d435ae2 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsInteractiveQueryIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsInteractiveQueryIntegrationTests.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.stream.binder.kafka.streams.integration; +package org.springframework.cloud.stream.binder.kafka.streams; import java.util.Map; @@ -23,25 +23,29 @@ import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.serialization.IntegerSerializer; import org.apache.kafka.common.serialization.Serdes; +import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Materialized; import org.apache.kafka.streams.kstream.Serialized; import org.apache.kafka.streams.state.HostInfo; +import org.apache.kafka.streams.state.QueryableStoreType; import org.apache.kafka.streams.state.QueryableStoreTypes; import org.apache.kafka.streams.state.ReadOnlyKeyValueStore; import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.ClassRule; import org.junit.Test; +import org.mockito.Mockito; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.cloud.stream.binder.kafka.streams.InteractiveQueryService; import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; +import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; @@ -54,6 +58,7 @@ import org.springframework.kafka.test.utils.KafkaTestUtils; import org.springframework.messaging.handler.annotation.SendTo; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.internal.verification.VerificationModeFactory.times; /** * @author Soby Chacko @@ -86,6 +91,29 @@ public class KafkaStreamsInteractiveQueryIntegrationTests { consumer.close(); } + @Test + public void testStateStoreRetrievalRetry() { + + KafkaStreams mock = Mockito.mock(KafkaStreams.class); + KafkaStreamsRegistry kafkaStreamsRegistry = new KafkaStreamsRegistry(); + kafkaStreamsRegistry.registerKafkaStreams(mock); + KafkaStreamsBinderConfigurationProperties binderConfigurationProperties = + new KafkaStreamsBinderConfigurationProperties(new KafkaProperties()); + binderConfigurationProperties.getStateStoreRetry().setMaxAttempts(3); + InteractiveQueryService interactiveQueryService = new InteractiveQueryService(kafkaStreamsRegistry, + binderConfigurationProperties); + + QueryableStoreType> storeType = QueryableStoreTypes.keyValueStore(); + try { + interactiveQueryService.getQueryableStore("foo", storeType); + } + catch (Exception ignored) { + + } + + Mockito.verify(mock, times(3)).store("foo", storeType); + } + @Test public void testKstreamBinderWithPojoInputAndStringOuput() throws Exception { SpringApplication app = new SpringApplication(ProductCountApplication.class);