diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltAwareProcessor.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltAwareProcessor.java index 385d25d15b..801b60eb20 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltAwareProcessor.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltAwareProcessor.java @@ -35,7 +35,7 @@ * * @author Soby Chacko * @author Steven Gantz - * @sinc 4.1.0 + * @since 4.1.0 */ public class DltAwareProcessor extends RecordRecoverableProcessor { @@ -51,11 +51,6 @@ public class DltAwareProcessor extends RecordRecoverablePr */ private final DltPublishingContext dltPublishingContext; - /** - * A {@link BiConsumer} that does the recovery of a failed record. - */ - private BiConsumer, Exception> processorRecordRecoverer; - /** * * @param delegateFunction {@link Function} to process the data diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderMetrics.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderMetrics.java index 5c5d49bed9..8ba8ca95d3 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderMetrics.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderMetrics.java @@ -22,8 +22,6 @@ import java.util.Map; import java.util.Objects; import java.util.Set; -import java.util.concurrent.Executors; -import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.locks.ReentrantLock; import java.util.function.ToDoubleFunction; @@ -82,8 +80,6 @@ public class KafkaStreamsBinderMetrics { private MeterBinder meterBinder; - private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); - private volatile Set currentMeters = new HashSet<>(); private static final ReentrantLock metricsLock = new ReentrantLock(); diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java index 134e8154da..3565dbc246 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java @@ -282,11 +282,6 @@ private void discoverOnlyKafkaStreamsResolvableTypes(String key, ResolvableType kafkaStreamsOnlyResolvableTypes.put(key, resolvableType); } - private void discoverOnlyKafkaStreamsResolvableTypesAndMethods(String key, ResolvableType resolvableType, Method method) { - kafkaStreamsOnlyResolvableTypes.put(key, resolvableType); - kafakStreamsOnlyMethods.put(key, method); - } - private void addResolvableTypeInfo(String key, Method method) { if (kafakStreamsOnlyMethods.size() == 1) { this.methods.put(key, method); diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index 4e96e6a55b..da75021fd9 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -16,7 +16,6 @@ package org.springframework.cloud.stream.binder.kafka; -import java.io.IOException; import java.lang.reflect.Field; import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; diff --git a/core/spring-cloud-stream-test-binder/src/main/java/org/springframework/cloud/stream/binder/test/TestChannelBinder.java b/core/spring-cloud-stream-test-binder/src/main/java/org/springframework/cloud/stream/binder/test/TestChannelBinder.java index a0ff8b21d6..135c6895c7 100644 --- a/core/spring-cloud-stream-test-binder/src/main/java/org/springframework/cloud/stream/binder/test/TestChannelBinder.java +++ b/core/spring-cloud-stream-test-binder/src/main/java/org/springframework/cloud/stream/binder/test/TestChannelBinder.java @@ -18,7 +18,6 @@ import java.util.function.Consumer; -import org.reactivestreams.Subscription; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cloud.stream.binder.AbstractMessageChannelBinder; @@ -29,8 +28,6 @@ import org.springframework.cloud.stream.binder.test.TestChannelBinderProvisioner.SpringIntegrationProducerDestination; import org.springframework.cloud.stream.provisioning.ConsumerDestination; import org.springframework.cloud.stream.provisioning.ProducerDestination; -import org.springframework.context.ApplicationEvent; -import org.springframework.context.ApplicationListener; import org.springframework.core.retry.RetryException; import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.acks.AcknowledgmentCallback; diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java index 3a14d737c2..495c38f1fc 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java @@ -16,7 +16,6 @@ package org.springframework.cloud.stream.binder; -import java.io.IOException; import java.util.LinkedHashMap; import java.util.Map; import java.util.concurrent.atomic.AtomicReference; diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/PartitionHandler.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/PartitionHandler.java index 683eb22f5f..c420998a11 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/PartitionHandler.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/PartitionHandler.java @@ -16,17 +16,13 @@ package org.springframework.cloud.stream.binder; -import java.lang.reflect.Field; import java.util.Map; -import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; -import org.springframework.context.expression.BeanFactoryResolver; import org.springframework.expression.EvaluationContext; import org.springframework.messaging.Message; import org.springframework.util.Assert; import org.springframework.util.CollectionUtils; -import org.springframework.util.ReflectionUtils; import org.springframework.util.StringUtils; /** @@ -186,18 +182,6 @@ private PartitionSelectorStrategy getPartitionSelectorStrategy( return partitionSelector; } - private static BeanFactory extractBeanFactoryFromEvaluationContext(EvaluationContext evaluationContext) { - try { - Field field = ReflectionUtils.findField(BeanFactoryResolver.class, "beanFactory"); - field.setAccessible(true); - return (BeanFactory) field.get(evaluationContext); - } - catch (Exception e) { - throw new RuntimeException("Failed to extract beanFactory from EvaluationContext. Please use different constructor" - + " which allows you to pass the instance of the beanFactory."); - } - } - /** * Default partition strategy; only works on keys with "real" hash codes, such as * String. Caller now always applies modulo so no need to do so here. diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java index 33fb152f2a..8ca93ce0ee 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java @@ -203,13 +203,6 @@ private FunctionInvocationWrapper retrieveFunction(String functionName) { return function; } - private void skipOutputConversionIfNecessary(String functionName) { - FunctionInvocationWrapper function = retrieveFunction(functionName); - if (function != null) { - function.setSkipOutputConversion(true); - } - } - private boolean isNativeEncodingNotSet(ProducerProperties producerProperties, ConsumerProperties consumerProperties, boolean input) { if (input) { diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/RoutingFunctionEnvironmentPostProcessor.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/RoutingFunctionEnvironmentPostProcessor.java index 7e37deb04c..a02e17e7d0 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/RoutingFunctionEnvironmentPostProcessor.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/RoutingFunctionEnvironmentPostProcessor.java @@ -21,7 +21,6 @@ import org.springframework.boot.SpringApplication; import org.springframework.cloud.function.context.config.RoutingFunction; import org.springframework.core.env.ConfigurableEnvironment; -import org.springframework.core.env.StandardEnvironment; import org.springframework.util.StringUtils; /** * diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java index 3ca1fa8742..7d6b7627e0 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java @@ -43,8 +43,6 @@ import org.springframework.cloud.function.context.catalog.SimpleFunctionRegistry.FunctionInvocationWrapper; import org.springframework.cloud.function.context.catalog.SimpleFunctionRegistry.PassThruFunction; import org.springframework.cloud.function.core.FunctionInvocationHelper; -import org.springframework.cloud.stream.binder.Binder; -import org.springframework.cloud.stream.binder.BinderFactory; import org.springframework.cloud.stream.binder.BinderWrapper; import org.springframework.cloud.stream.binder.ProducerProperties; import org.springframework.cloud.stream.binding.BindingService; @@ -326,14 +324,6 @@ private void addPartitioningInterceptorIfNeedBe(ProducerProperties producerPrope } } - private String resolveBinderTargetType(String channelName, String binderName, Class bindableType, BinderFactory binderFactory) { - String binderConfigurationName = binderName != null ? binderName : this.bindingServiceProperties - .getBinder(channelName); - Binder binder = binderFactory.getBinder(binderConfigurationName, bindableType); - String targetProtocol = binder.getClass().getSimpleName().startsWith("Rabbit") ? "amqp" : "kafka"; - return targetProtocol; - } - private void addGlobalChannelInterceptorProcessor(AbstractMessageChannel messageChannel, String destinationName) { final GlobalChannelInterceptorProcessor globalChannelInterceptorProcessor = this.applicationContext.getBean(GlobalChannelInterceptorProcessor.class); diff --git a/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/DefaultPollableMessageSourceTests.java b/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/DefaultPollableMessageSourceTests.java index fe855c982f..8ef5255115 100644 --- a/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/DefaultPollableMessageSourceTests.java +++ b/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/DefaultPollableMessageSourceTests.java @@ -23,7 +23,6 @@ import org.junit.jupiter.api.Test; import org.springframework.integration.channel.DirectChannel; -import org.springframework.integration.core.MessageSource; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; diff --git a/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroSchemaRegistryClientMessageConverter.java b/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroSchemaRegistryClientMessageConverter.java index dcdf2eeb2d..f69acbcecd 100644 --- a/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroSchemaRegistryClientMessageConverter.java +++ b/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroSchemaRegistryClientMessageConverter.java @@ -113,9 +113,6 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag */ public static final MimeType DEFAULT_AVRO_MIME_TYPE = new MimeType("application", "*+" + AVRO_FORMAT); - private static final AvroSchemaServiceManager defaultAvroSchemaServiceManager = - new AvroSchemaServiceManagerImpl(); - private final CacheManager cacheManager; protected Resource[] schemaImports = new Resource[]{};