diff --git a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/kafka-streams.adoc b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/kafka-streams.adoc index b0b8e284d..42f9f23c6 100644 --- a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/kafka-streams.adoc +++ b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/kafka-streams.adoc @@ -577,6 +577,38 @@ public KStream process(KStream input) { } ---- +== State Store + +State store is created automatically by Kafka Stream when Streas DSL is used. When use processor API, in case you want to +create and register a state store manually, you can use `KafkaStreamsStateStore` annotation. You can specify store name, +type, whether to enable log, whether disable cache, etc, and those parameters will be injected into KStream building +process in Kafka Streams binder to create and register the store to your KStream. After that, you can access the same way +how you access in normal Kafka Streams code. + +Creation code: +[source] +---- +@KafkaStreamsStateStore(name="mystate", type= KafkaStreamsStateStoreProperties.StoreType.WINDOW, lengthMs=300000) +public void process(KStream input) { + ... +} +---- + +Access code: +[source] +---- +Processor() { + + WindowStore state; + + @Override + public void init(ProcessorContext processorContext) { + state = (WindowStore)processorContext.getStateStore("mystate"); + } + ... +} +---- + == Interactive Queries As part of the public Kafka Streams binder API, we expose a class called `QueryableStoreRegistry`. You can access this diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java index 128e0d43e..7caa21a12 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java @@ -35,6 +35,9 @@ import org.apache.kafka.streams.kstream.KTable; import org.apache.kafka.streams.kstream.Materialized; import org.apache.kafka.streams.state.KeyValueStore; +import org.apache.kafka.streams.state.StoreBuilder; +import org.apache.kafka.streams.state.Stores; + import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanInitializationException; import org.springframework.beans.factory.config.BeanDefinition; @@ -44,9 +47,11 @@ import org.springframework.beans.factory.support.BeanDefinitionRegistry; import org.springframework.cloud.stream.annotation.Input; import org.springframework.cloud.stream.annotation.StreamListener; import org.springframework.cloud.stream.binder.ConsumerProperties; +import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsStateStore; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsConsumerProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsExtendedBindingProperties; +import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsStateStoreProperties; import org.springframework.cloud.stream.binding.StreamListenerErrorMessages; import org.springframework.cloud.stream.binding.StreamListenerParameterAdapter; import org.springframework.cloud.stream.binding.StreamListenerResultAdapter; @@ -79,6 +84,7 @@ import org.springframework.util.StringUtils; * 3. Each StreamListener method that it orchestrates gets its own {@link StreamsBuilderFactoryBean} and {@link StreamsConfig} * * @author Soby Chacko + * @author Lei Chen */ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListenerSetupMethodOrchestrator, ApplicationContextAware { @@ -230,10 +236,12 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene StreamsBuilderFactoryBean streamsBuilderFactoryBean = methodStreamsBuilderFactoryBeanMap.get(method); StreamsBuilder streamsBuilder = streamsBuilderFactoryBean.getObject(); KafkaStreamsConsumerProperties extendedConsumerProperties = kafkaStreamsExtendedBindingProperties.getExtendedConsumerProperties(inboundName); + //get state store spec + KafkaStreamsStateStoreProperties spec = buildStateStoreSpec(method); Serde keySerde = this.keyValueSerdeResolver.getInboundKeySerde(extendedConsumerProperties); Serde valueSerde = this.keyValueSerdeResolver.getInboundValueSerde(bindingProperties.getConsumer(), extendedConsumerProperties); if (parameterType.isAssignableFrom(KStream.class)) { - KStream stream = getkStream(inboundName, bindingProperties, streamsBuilder, keySerde, valueSerde); + KStream stream = getkStream(inboundName, spec, bindingProperties, streamsBuilder, keySerde, valueSerde); KStreamBoundElementFactory.KStreamWrapper kStreamWrapper = (KStreamBoundElementFactory.KStreamWrapper) targetBean; //wrap the proxy created during the initial target type binding with real object (KStream) kStreamWrapper.wrap((KStream) stream); @@ -288,8 +296,51 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene .withValueSerde(v)); } - private KStream getkStream(String inboundName, BindingProperties bindingProperties, StreamsBuilder streamsBuilder, + private StoreBuilder buildStateStore(KafkaStreamsStateStoreProperties spec) { + try { + Serde keySerde = this.keyValueSerdeResolver.getStateStoreKeySerde(spec.getKeySerdeString()); + Serde valueSerde = this.keyValueSerdeResolver.getStateStoreValueSerde(spec.getValueSerdeString()); + StoreBuilder builder; + switch (spec.getType()) { + case KEYVALUE: + builder = Stores.keyValueStoreBuilder(Stores.persistentKeyValueStore(spec.getName()), keySerde, valueSerde); + break; + case WINDOW: + builder = Stores.windowStoreBuilder(Stores.persistentWindowStore(spec.getName(), spec.getRetention(), 3, spec.getLength(), false), + keySerde, + valueSerde); + break; + case SESSION: + builder = Stores.sessionStoreBuilder(Stores.persistentSessionStore(spec.getName(), spec.getRetention()), keySerde, valueSerde); + break; + default: + throw new UnsupportedOperationException("state store type (" + spec.getType() + ") is not supported!"); + } + if (spec.isCacheEnabled()) { + builder = builder.withCachingEnabled(); + } + if (spec.isLoggingDisabled()) { + builder = builder.withLoggingDisabled(); + } + + return builder; + + }catch (Exception e) { + LOG.error("failed to build state store exception : " + e); + throw e; + } + } + + + private KStream getkStream(String inboundName, KafkaStreamsStateStoreProperties storeSpec, BindingProperties bindingProperties, StreamsBuilder streamsBuilder, Serde keySerde, Serde valueSerde) { + if (storeSpec != null) { + StoreBuilder storeBuilder = buildStateStore(storeSpec); + streamsBuilder.addStateStore(storeBuilder); + if (LOG.isInfoEnabled()) { + LOG.info("state store " + storeBuilder.name() + " added to topology"); + } + } KStream stream = streamsBuilder.stream(bindingServiceProperties.getBindingDestination(inboundName), Consumed.with(keySerde, valueSerde)); final boolean nativeDecoding = bindingServiceProperties.getConsumerProperties(inboundName).isUseNativeDecoding(); @@ -431,4 +482,24 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene } return null; } + + @SuppressWarnings({"unchecked"}) + private static KafkaStreamsStateStoreProperties buildStateStoreSpec(Method method) { + KafkaStreamsStateStore spec = AnnotationUtils.findAnnotation(method, KafkaStreamsStateStore.class); + if (spec != null) { + Assert.isTrue(!ObjectUtils.isEmpty(spec.name()), "name cannot be empty"); + Assert.isTrue(spec.name().length() >= 1, "name cannot be empty."); + KafkaStreamsStateStoreProperties props = new KafkaStreamsStateStoreProperties(); + props.setName(spec.name()); + props.setType(spec.type()); + props.setLength(spec.lengthMs()); + props.setKeySerdeString(spec.keySerde()); + props.setRetention(spec.retentionMs()); + props.setValueSerdeString(spec.valueSerde()); + props.setCacheEnabled(spec.cache()); + props.setLoggingDisabled(!spec.logging()); + return props; + } + return null; + } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KeyValueSerdeResolver.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KeyValueSerdeResolver.java index 6461fc99e..4fd5d4cf7 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KeyValueSerdeResolver.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KeyValueSerdeResolver.java @@ -38,12 +38,16 @@ import org.springframework.util.StringUtils; * If native decoding is disabled, then the binder will do the deserialization on value and ignore any Serde set for value * and rely on the contentType provided. Keys are always deserialized at the broker. * + * * Same rules apply on the outbound. If native encoding is enabled, then value serialization is done at the broker using * any binder level Serde for value, if not using common Serde, if not, then byte[]. * If native encoding is disabled, then the binder will do serialization using a contentType. Keys are always serialized * by the broker. * + * For state store, use serdes class specified in {@link KafkaStreamsStateStore} to create Serde accordingly. + * * @author Soby Chacko + * @author Lei Chen */ class KeyValueSerdeResolver { @@ -130,6 +134,31 @@ class KeyValueSerdeResolver { return valueSerde; } + /** + * Provide the {@link Serde} for state store + * + * @param keySerdeString serde class used for key + * @return {@link Serde} for the state store key. + */ + public Serde getStateStoreKeySerde(String keySerdeString) { + return getKeySerde(keySerdeString); + } + + /** + * Provide the {@link Serde} for state store value + * + * @param valueSerdeString serde class used for value + * @return {@link Serde} for the state store value. + */ + public Serde getStateStoreValueSerde(String valueSerdeString) { + try { + return getValueSerde(valueSerdeString); + } + catch (ClassNotFoundException e) { + throw new IllegalStateException("Serde class not found: ", e); + } + } + private Serde getKeySerde(String keySerdeString) { Serde keySerde; try { diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/annotations/KafkaStreamsStateStore.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/annotations/KafkaStreamsStateStore.java new file mode 100644 index 000000000..96b06569b --- /dev/null +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/annotations/KafkaStreamsStateStore.java @@ -0,0 +1,105 @@ +/* + * Copyright 2018 the original author or authors. + * + * Licensed 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.springframework.cloud.stream.binder.kafka.streams.annotations; + + +import java.lang.annotation.ElementType; +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; +import java.lang.annotation.Target; + +import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsStateStoreProperties;; + + +/** + * Interface for Kafka Stream state store. + * + * This interface can be used to inject a state store specification into KStream building process so + * that the desired store can be built by StreamBuilder and added to topology for later use by processors. + * This is particularly useful when need to combine stream DSL with low level processor APIs. In those cases, + * if a writable state store is desired in processors, it needs to be created using this annotation. + * Here is the example. + * + *
+ *     @StreamListener("input")
+ *     @KafkaStreamsStateStore(name="mystate", type= KafkaStreamsStateStoreProperties.StoreType.WINDOW, size=300000)
+ *	   public void process(KStream input) {
+ *         ......
+ *     }
+ *
+ * + * With that, you should be able to read/write this state store in your processor/transformer code. + * + *
+ * 		new Processor() {
+ * 			WindowStore state;
+ * 			@Override
+ *			public void init(ProcessorContext processorContext) {
+ *			state = (WindowStore)processorContext.getStateStore("mystate");
+ *				......
+ *			}
+ *		}
+ *
+ * + * @author Lei Chen + */ + +@Target({ ElementType.TYPE, ElementType.METHOD, ElementType.ANNOTATION_TYPE }) +@Retention(RetentionPolicy.RUNTIME) + +public @interface KafkaStreamsStateStore { + + /** + * @return name of state store. + */ + String name() default ""; + + /** + * @return {@link KafkaStreamsStateStoreProperties.StoreType} of state store. + */ + KafkaStreamsStateStoreProperties.StoreType type() default KafkaStreamsStateStoreProperties.StoreType.KEYVALUE; + + /** + * @return key serde of state store. + */ + String keySerde() default "org.apache.kafka.common.serialization.Serdes$StringSerde"; + + /** + * @return value serde of state store. + */ + String valueSerde() default "org.apache.kafka.common.serialization.Serdes$StringSerde"; + + /** + * @return length in milli-second of window(for windowed store). + */ + long lengthMs() default 0; + + /** + * @return the maximum period of time in milli-second to keep each window in this store(for windowed store). + */ + long retentionMs() default 0; + + /** + * @return whether caching should be enabled on the created store. + */ + boolean cache() default false; + + /** + * @return whether logging should be enabled on the created store. + */ + boolean logging() default true; +} diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsStateStoreProperties.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsStateStoreProperties.java new file mode 100644 index 000000000..a2d9f67bf --- /dev/null +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsStateStoreProperties.java @@ -0,0 +1,151 @@ +/* + * Copyright 2018 the original author or authors. + * + * Licensed 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.springframework.cloud.stream.binder.kafka.streams.properties; + + +/** + * @author Lei Chen + */ +public class KafkaStreamsStateStoreProperties { + + public enum StoreType { + KEYVALUE("keyvalue"), + WINDOW("window"), + SESSION("session") + ; + + private final String type; + + /** + * @param type + */ + StoreType(final String type) { + this.type = type; + } + + @Override + public String toString() { + return type; + } + } + + + /** + * name for this state store + */ + private String name; + + /** + * type for this state store + */ + private StoreType type; + + /** + * Size/length of this state store in ms. Only applicable for window store. + */ + private long length; + + /** + * Retention period for this state store in ms. + */ + private long retention; + + /** + * Key serde class specified per state store. + */ + private String keySerdeString; + + /** + * Value serde class specified per state store. + */ + private String valueSerdeString; + + /** + * Whether enable cache in this state store. + */ + private boolean cacheEnabled; + + /** + * Whether enable logging in this state store. + */ + private boolean loggingDisabled; + + + public String getName() { + return name; + } + + public void setName(String name) { + this.name = name; + } + + public StoreType getType() { + return type; + } + + public void setType(StoreType type) { + this.type = type; + } + + public long getLength() { + return length; + } + + public void setLength(long length) { + this.length = length; + } + + public long getRetention() { + return retention; + } + + public void setRetention(long retention) { + this.retention = retention; + } + + public String getKeySerdeString() { + return keySerdeString; + } + + public void setKeySerdeString(String keySerdeString) { + this.keySerdeString = keySerdeString; + } + + public String getValueSerdeString() { + return valueSerdeString; + } + + public void setValueSerdeString(String valueSerdeString) { + this.valueSerdeString = valueSerdeString; + } + + public boolean isCacheEnabled() { + return cacheEnabled; + } + + public void setCacheEnabled(boolean cacheEnabled) { + this.cacheEnabled = cacheEnabled; + } + + public boolean isLoggingDisabled() { + return loggingDisabled; + } + + public void setLoggingDisabled(boolean loggingDisabled) { + this.loggingDisabled = loggingDisabled; + } +} diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStateStoreIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStateStoreIntegrationTests.java new file mode 100644 index 000000000..e01b503ba --- /dev/null +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStateStoreIntegrationTests.java @@ -0,0 +1,152 @@ +/* + * Copyright 2018 the original author or authors. + * + * Licensed 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.springframework.cloud.stream.binder.kafka.streams; + + +import java.util.Map; + +import org.apache.kafka.streams.kstream.KStream; +import org.apache.kafka.streams.processor.Processor; +import org.apache.kafka.streams.processor.ProcessorContext; +import org.apache.kafka.streams.state.WindowStore; +import org.junit.ClassRule; +import org.junit.Test; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.WebApplicationType; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.annotation.Input; +import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsStateStore; +import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsStateStoreProperties; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.kafka.core.DefaultKafkaProducerFactory; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.test.rule.KafkaEmbedded; +import org.springframework.kafka.test.utils.KafkaTestUtils; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * @author Lei Chen + * @author Soby Chacko + */ +public class KafkaStreamsStateStoreIntegrationTests { + + @ClassRule + public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, "counts-id"); + + @Test + public void testKstreamStateStore() throws Exception { + SpringApplication app = new SpringApplication(ProductCountApplication.class); + app.setWebApplicationType(WebApplicationType.NONE); + ConfigurableApplicationContext context = app.run("--server.port=0", + "--spring.jmx.enabled=false", + "--spring.cloud.stream.bindings.input.destination=foobar", + "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", + "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.bindings.input.consumer.headerMode=raw", + "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), + "--spring.cloud.stream.kafka.streams.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString()); + try { + Thread.sleep(2000); + receiveAndValidateFoo(context); + } catch (Exception e) { + throw e; + } finally { + context.close(); + } + } + + private void receiveAndValidateFoo(ConfigurableApplicationContext context) throws Exception { + Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>(senderProps); + KafkaTemplate template = new KafkaTemplate<>(pf, true); + template.setDefaultTopic("foobar"); + template.sendDefault("{\"id\":\"123\"}"); + Thread.sleep(1000); + + //assertions + ProductCountApplication productCount = context.getBean(ProductCountApplication.class); + WindowStore state = productCount.state; + assertThat(state != null).isTrue(); + assertThat(state.name()).isEqualTo("mystate"); + assertThat(state.persistent()).isTrue(); + assertThat(productCount.processed).isTrue(); + } + + @EnableBinding(KafkaStreamsProcessorX.class) + @EnableAutoConfiguration + public static class ProductCountApplication { + + WindowStore state; + boolean processed; + + @StreamListener("input") + @KafkaStreamsStateStore(name = "mystate", type = KafkaStreamsStateStoreProperties.StoreType.WINDOW, lengthMs = 300000) + @SuppressWarnings({"deprecation", "unchecked"}) + public void process(KStream input) { + + input + .process(() -> new Processor() { + + @Override + public void init(ProcessorContext processorContext) { + state = (WindowStore) processorContext.getStateStore("mystate"); + } + + @Override + public void process(Object s, Product product) { + processed = true; + } + + @Override + public void punctuate(long l) { + + } + + @Override + public void close() { + if (state != null) { + state.close(); + } + } + }, "mystate"); + } + } + + public static class Product { + + Integer id; + + public Integer getId() { + return id; + } + + public void setId(Integer id) { + this.id = id; + } + } + + interface KafkaStreamsProcessorX { + + @Input("input") + KStream input(); + } +}