Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@
*
* @author Soby Chacko
* @author Steven Gantz
* @sinc 4.1.0
* @since 4.1.0
*/
public class DltAwareProcessor<KIn, VIn, KOut, VOut> extends RecordRecoverableProcessor<KIn, VIn, KOut, VOut> {

Expand All @@ -51,11 +51,6 @@ public class DltAwareProcessor<KIn, VIn, KOut, VOut> extends RecordRecoverablePr
*/
private final DltPublishingContext dltPublishingContext;

/**
* A {@link BiConsumer} that does the recovery of a failed record.
*/
private BiConsumer<Record<KIn, VIn>, Exception> processorRecordRecoverer;

/**
*
* @param delegateFunction {@link Function} to process the data
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -82,8 +80,6 @@ public class KafkaStreamsBinderMetrics {

private MeterBinder meterBinder;

private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();

private volatile Set<MetricName> currentMeters = new HashSet<>();

private static final ReentrantLock metricsLock = new ReentrantLock();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

/**
Expand Down Expand Up @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
/**
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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[]{};
Expand Down
Loading