From 38d6deb4d5dd7309a2baa0be8346be26e9ef619a Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 29 Jun 2018 16:45:46 -0400 Subject: [PATCH] Provide programmatic access to KafkaStreams object Providing access to the underlying StreamBuilderFactoryBean by making the bean name deterministic. Eariler, the binder was using UUID to make the stream builder factory bean names unique in the event of multiple StreamListeners. Switching to use the method name instead to keep the StreamBuilder factory beans unique while providing a deterministic way to giving it programmatic access. Polishing docs Fixes #396 --- .../src/main/asciidoc/kafka-streams.adoc | 14 ++++++++++++++ ...reamsStreamListenerSetupMethodOrchestrator.java | 9 +++------ ...afkaStreamsBinderWordCountIntegrationTests.java | 10 ++++++++++ 3 files changed, 27 insertions(+), 6 deletions(-) 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 74666b8b1..3b0752c1f 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 @@ -652,3 +652,17 @@ else { //query from the remote host } ---- + +== Accessing the underlying KafkaStreams object + +`StreamBuilderFactoryBean` from spring-kafka that is responsible for constructing the `KafkaStreams` object can be accessed programmatically. +Each `StreamBuilderFactoryBean` is registered as `stream-builder` and appended with the `StreamListener` method name. +If your `StreamListener` method is named as `process` for example, the stream builder bean is named as `stream-builder-process`. +Since this is a factory bean, it should be accessed by prepending an ampersand (`&`) when accessing it programmatically. +Following is an example and it assumes the `StreamListener` method is named as `process` + +[source] +---- +StreamsBuilderFactoryBean streamsBuilderFactoryBean = context.getBean("&stream-builder-process", StreamsBuilderFactoryBean.class); + KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); +---- \ No newline at end of file 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 0a5a4ed9d..eceaed414 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 @@ -21,7 +21,6 @@ import java.util.Arrays; import java.util.Collection; import java.util.HashMap; import java.util.Map; -import java.util.UUID; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -35,7 +34,6 @@ import org.apache.kafka.streams.kstream.KStream; 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; @@ -386,12 +384,11 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene ConfigurableListableBeanFactory beanFactory = this.applicationContext.getBeanFactory(); StreamsBuilderFactoryBean streamsBuilder = new StreamsBuilderFactoryBean(); streamsBuilder.setAutoStartup(false); - String uuid = UUID.randomUUID().toString(); BeanDefinition streamsBuilderBeanDefinition = BeanDefinitionBuilder.genericBeanDefinition((Class) streamsBuilder.getClass(), () -> streamsBuilder) .getRawBeanDefinition(); - ((BeanDefinitionRegistry) beanFactory).registerBeanDefinition("stream-builder-" + uuid, streamsBuilderBeanDefinition); - StreamsBuilderFactoryBean streamsBuilderX = applicationContext.getBean("&stream-builder-" + uuid, StreamsBuilderFactoryBean.class); + ((BeanDefinitionRegistry) beanFactory).registerBeanDefinition("stream-builder-" + method.getName(), streamsBuilderBeanDefinition); + StreamsBuilderFactoryBean streamsBuilderX = applicationContext.getBean("&stream-builder-" + method.getName(), StreamsBuilderFactoryBean.class); String group = bindingProperties.getGroup(); if (!StringUtils.hasText(group)) { group = binderConfigurationProperties.getApplicationId(); @@ -421,7 +418,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene BeanDefinition streamsConfigBeanDefinition = BeanDefinitionBuilder.genericBeanDefinition((Class) streamsConfig.getClass(), () -> streamsConfig) .getRawBeanDefinition(); - ((BeanDefinitionRegistry) beanFactory).registerBeanDefinition("streamsConfig-" + uuid, streamsConfigBeanDefinition); + ((BeanDefinitionRegistry) beanFactory).registerBeanDefinition("streamsConfig-" + method.getName(), streamsConfigBeanDefinition); streamsBuilder.setStreamsConfig(streamsConfig); methodStreamsBuilderFactoryBeanMap.put(method, streamsBuilderX); diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderWordCountIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderWordCountIntegrationTests.java index 41f0374c4..001cd3f9f 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderWordCountIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderWordCountIntegrationTests.java @@ -24,11 +24,14 @@ import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; 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.kstream.TimeWindows; +import org.apache.kafka.streams.state.QueryableStoreTypes; +import org.apache.kafka.streams.state.ReadOnlyWindowStore; import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.ClassRule; @@ -48,6 +51,7 @@ import org.springframework.context.ConfigurableApplicationContext; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.StreamsBuilderFactoryBean; import org.springframework.kafka.test.rule.KafkaEmbedded; import org.springframework.kafka.test.utils.KafkaTestUtils; import org.springframework.messaging.handler.annotation.SendTo; @@ -101,6 +105,12 @@ public class KafkaStreamsBinderWordCountIntegrationTests { "--spring.cloud.stream.kafka.streams.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString()); try { receiveAndValidate(context); + //Assertions on StreamBuilderFactoryBean + StreamsBuilderFactoryBean streamsBuilderFactoryBean = context.getBean("&stream-builder-process", StreamsBuilderFactoryBean.class); + KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); + ReadOnlyWindowStore store = kafkaStreams.store("foo-WordCounts", QueryableStoreTypes.windowStore()); + assertThat(store).isNotNull(); + } finally { context.close(); }