From 4505874ddce464e0a95d7e33f4a0f9eebc0e85a3 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Thu, 1 Feb 2018 18:31:00 -0500 Subject: [PATCH] GH-548: StreamsBuilderFactoryBean enhancements Fixes https://github.com/spring-projects/spring-kafka/issues/548 * Make `StreamsConfig` customizable in the `StreamsBuilderFactoryBean` * Set the phase on `StreamsBuilderFactoryBean` to the `Integer.MAX_VALUE - 1000` * Adding tests * Addressing PR review comments * Addressing PR review comments --- .../kafka/core/StreamsBuilderFactoryBean.java | 33 ++++++- .../StreamsBuilderFactoryLateConfigTests.java | 95 +++++++++++++++++++ 2 files changed, 124 insertions(+), 4 deletions(-) create mode 100644 spring-kafka/src/test/java/org/springframework/kafka/core/StreamsBuilderFactoryLateConfigTests.java 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; + } + } +}