diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/StreamsBuilderFactoryBean.java b/spring-kafka/src/main/java/org/springframework/kafka/core/StreamsBuilderFactoryBean.java index e1a45492..d5d9b352 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/StreamsBuilderFactoryBean.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/StreamsBuilderFactoryBean.java @@ -1,5 +1,5 @@ /* - * Copyright 2017 the original author or authors. + * Copyright 2017-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. @@ -36,6 +36,7 @@ import org.springframework.util.Assert; * * @author Artem Bilan * @author Ivan Ursul + * @author Soby Chacko * * @since 1.1.4 */ @@ -43,17 +44,17 @@ public class StreamsBuilderFactoryBean extends AbstractFactoryBean props = new HashMap<>(); + props.put(StreamsConfig.APPLICATION_ID_CONFIG, APPLICATION_ID); + props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, this.brokerAddresses); + StreamsConfig streamsConfig = new StreamsConfig(props); + streamsBuilderFactoryBean.setStreamsConfig(streamsConfig); + + assertThat(streamsBuilderFactoryBean.isRunning()).isFalse(); + streamsBuilderFactoryBean.start(); + assertThat(streamsBuilderFactoryBean.isRunning()).isTrue(); + } + + @Configuration + @EnableKafka + @EnableKafkaStreams + public static class KafkaStreamsConfiguration { + + @Bean(name = KafkaStreamsDefaultConfiguration.DEFAULT_STREAMS_BUILDER_BEAN_NAME) + public StreamsBuilderFactoryBean defaultKafkaStreamsBuilder() { + StreamsBuilderFactoryBean streamsBuilderFactoryBean = new StreamsBuilderFactoryBean(); + streamsBuilderFactoryBean.setAutoStartup(false); + return streamsBuilderFactoryBean; + } + } +}