From 7d7d7aa0217fbb4a66205bd7437add48d40b1a3f Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Thu, 16 Jun 2022 18:12:11 -0400 Subject: [PATCH] SendToDlqAndContinue needs to be excluded In Kafka Streams binder, we need to exclude the SendToDlqAndContinue BiFunction by adding it to the ineligible function definitions. See this commit for more context: https://github.com/spring-cloud/spring-cloud-function/commit/8ec15fa3ca7fe0d0c58ca6f6c564073c85f7388a --- ...StreamsBinderEnvironmentPostProcessor.java | 46 +++++++++++++++++++ .../main/resources/META-INF/spring.factories | 2 + ...appingsProviderAutoConfigurationTests.java | 1 + .../streams/SerdeResolverUtilsTests.java | 9 +++- 4 files changed, 56 insertions(+), 2 deletions(-) create mode 100644 binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderEnvironmentPostProcessor.java 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(