diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderEnvironmentPostProcessor.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderEnvironmentPostProcessor.java new file mode 100644 index 000000000..18a319188 --- /dev/null +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderEnvironmentPostProcessor.java @@ -0,0 +1,46 @@ +/* + * Copyright 2022-2022 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 + * + * https://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.HashMap; +import java.util.Map; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.env.EnvironmentPostProcessor; +import org.springframework.core.env.ConfigurableEnvironment; +import org.springframework.core.env.MapPropertySource; + +/** + * {@link EnvironmentPostProcessor} to exclude the SendToDlqAndContinue BiFunction + * in core Spring Cloud Function by adding it to the ineligible function definitions. + * + * @author Soby Chacko + */ +public class KafkaStreamsBinderEnvironmentPostProcessor implements EnvironmentPostProcessor { + + @Override + public void postProcessEnvironment(ConfigurableEnvironment environment, + SpringApplication application) { + Map kafkaStreamsBinderIneligibleDefns = new HashMap<>(); + kafkaStreamsBinderIneligibleDefns.put("spring.cloud.function.ineligible-definitions", + "sendToDlqAndContinue"); + + environment.getPropertySources().addLast(new MapPropertySource( + "KAFKA_STREAMS_BINDER_INELIGIBLE_DEFINITIONS", kafkaStreamsBinderIneligibleDefns)); + } + +} diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/resources/META-INF/spring.factories b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/resources/META-INF/spring.factories index 00fdadf9e..fa6f0369e 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/resources/META-INF/spring.factories +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/resources/META-INF/spring.factories @@ -3,3 +3,5 @@ org.springframework.boot.autoconfigure.EnableAutoConfiguration=\ org.springframework.cloud.stream.binder.kafka.streams.KafkaStreamsBinderSupportAutoConfiguration,\ org.springframework.cloud.stream.binder.kafka.streams.function.KafkaStreamsFunctionAutoConfiguration,\ org.springframework.cloud.stream.binder.kafka.streams.endpoint.KafkaStreamsTopologyEndpointAutoConfiguration +org.springframework.boot.env.EnvironmentPostProcessor=\ +org.springframework.cloud.stream.binder.kafka.streams.KafkaStreamsBinderEnvironmentPostProcessor diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/ExtendedBindingHandlerMappingsProviderAutoConfigurationTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/ExtendedBindingHandlerMappingsProviderAutoConfigurationTests.java index 43f21e243..fdfbd2953 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/ExtendedBindingHandlerMappingsProviderAutoConfigurationTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/ExtendedBindingHandlerMappingsProviderAutoConfigurationTests.java @@ -32,6 +32,7 @@ class ExtendedBindingHandlerMappingsProviderAutoConfigurationTests { private final ApplicationContextRunner contextRunner = new ApplicationContextRunner() .withUserConfiguration(KafkaStreamsTestApp.class) .withPropertyValues( + "spring.cloud.function.ineligible-definitions: sendToDlqAndContinue", "spring.cloud.stream.kafka.streams.default.consumer.application-id: testApp123", "spring.cloud.stream.kafka.streams.default.consumer.consumed-as: default-consumer", "spring.cloud.stream.kafka.streams.default.consumer.materialized-as: default-materializer", diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtilsTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtilsTests.java index c57bd9123..8355b0515 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtilsTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtilsTests.java @@ -74,6 +74,7 @@ class SerdeResolverUtilsTests { void returnsSerdeBeanForMatchingType() { this.contextRunner .withConfiguration(AutoConfigurations.of(SerdeResolverSimpleTestApp.class)) + .withPropertyValues("spring.cloud.function.ineligible-definitions: sendToDlqAndContinue") .run((context) -> { ResolvableType fooType = ResolvableType.forClass(Foo.class); assertThat(SerdeResolverUtils.resolveForType(context, fooType, fallback)).isInstanceOf(FooSerde.class); @@ -214,7 +215,9 @@ class SerdeResolverUtilsTests { ResolvableType geWildcard = ResolvableType.forType(new ParameterizedTypeReference>() { }); ResolvableType geRaw = ResolvableType.forRawClass(GenericEvent.class); - new ApplicationContextRunner().withUserConfiguration(SerdeResolverSimpleTestApp.class).run((context) -> { + new ApplicationContextRunner().withUserConfiguration(SerdeResolverSimpleTestApp.class) + .withPropertyValues("spring.cloud.function.ineligible-definitions: sendToDlqAndContinue") + .run((context) -> { assertThat(SerdeResolverUtils.beanNamesForMatchingSerdes(context, geDate)) .containsExactly( @@ -284,7 +287,9 @@ class SerdeResolverUtilsTests { ResolvableType geWildcard = ResolvableType.forType(new ParameterizedTypeReference>() { }); ResolvableType geRaw = ResolvableType.forRawClass(GenericEvent.class); - new ApplicationContextRunner().withUserConfiguration(SerdeResolverComplexTestApp.class).run((context) -> { + new ApplicationContextRunner().withUserConfiguration(SerdeResolverComplexTestApp.class) + .withPropertyValues("spring.cloud.function.ineligible-definitions: sendToDlqAndContinue") + .run((context) -> { assertThat(SerdeResolverUtils.beanNamesForMatchingSerdes(context, geFooDate)) .containsExactly(