diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java index c7238a8c3..8cbc2e8ad 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java @@ -18,6 +18,8 @@ package org.springframework.cloud.stream.binder.kafka.streams; import java.util.Set; +import org.apache.kafka.streams.errors.StreamsUncaughtExceptionHandler; + import org.springframework.beans.factory.DisposableBean; import org.springframework.context.SmartLifecycle; import org.springframework.kafka.KafkaException; @@ -84,6 +86,12 @@ class StreamsBuilderFactoryManager implements SmartLifecycle { if (this.listener != null) { streamsBuilderFactoryBean.addListener(this.listener); } + // By default, we shutdown the client if there is an uncaught exception in the application. + // Users can override this by customizing SBFB. See this issue for more details: + // https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/1110 + streamsBuilderFactoryBean.setStreamsUncaughtExceptionHandler(exception -> + StreamsUncaughtExceptionHandler.StreamThreadExceptionResponse.SHUTDOWN_CLIENT); + // Starting the stream. streamsBuilderFactoryBean.start(); this.kafkaStreamsRegistry.registerKafkaStreams(streamsBuilderFactoryBean); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java index 6c3d70b5c..a0694b152 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java @@ -257,19 +257,6 @@ public class KafkaStreamsBinderHealthIndicatorTests { }); } - @Bean - public StreamsBuilderFactoryBeanConfigurer customizer() { - return factoryBean -> { - factoryBean.setKafkaStreamsCustomizer(new KafkaStreamsCustomizer() { - @Override - public void customize(KafkaStreams kafkaStreams) { - kafkaStreams.setUncaughtExceptionHandler(exception -> - StreamsUncaughtExceptionHandler.StreamThreadExceptionResponse.SHUTDOWN_CLIENT); - } - }); - }; - } - } @EnableBinding({ KafkaStreamsProcessor.class, KafkaStreamsProcessorX.class })