From 33e58a0baaa099371c55643bb45502b14e4dc826 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 10 Jan 2023 11:09:25 -0500 Subject: [PATCH] Minor code cleanup/deprecation removal etc. --- .../properties/KafkaBinderConfigurationProperties.java | 2 +- .../kafka/streams/AbstractKafkaStreamsBinderProcessor.java | 4 ++-- .../streams/KafkaStreamsBinderSupportAutoConfiguration.java | 4 ++-- .../binder/kafka/streams/KafkaStreamsBinderUtils.java | 6 +++--- .../streams/KafkaStreamsMessageConversionDelegate.java | 2 +- .../streams/endpoint/KafkaStreamsTopologyEndpoint.java | 3 +-- .../rabbit/provisioning/RabbitExchangeQueueProvisioner.java | 2 +- 7 files changed, 11 insertions(+), 12 deletions(-) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java index 91d36ae69..96dad5fca 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java @@ -272,7 +272,7 @@ public class KafkaBinderConfigurationProperties { private String toConnectionString(String[] hosts, String defaultPort) { String[] fullyFormattedHosts = new String[hosts.length]; for (int i = 0; i < hosts.length; i++) { - if (hosts[i].contains(":") || StringUtils.isEmpty(defaultPort)) { + if (hosts[i].contains(":") || !StringUtils.hasText(defaultPort)) { fullyFormattedHosts[i] = hosts[i]; } else { diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java index d1ff1bf94..5facdb478 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java @@ -206,7 +206,7 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application } String applicationId = functionConfig.getApplicationId(); - if (!StringUtils.isEmpty(applicationId)) { + if (StringUtils.hasText(applicationId)) { streamConfiguration.put(StreamsConfig.APPLICATION_ID_CONFIG, applicationId); } } @@ -614,7 +614,7 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application private Processor eventTypeProcessor(KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties, AtomicBoolean matched, AtomicReference topicObject, AtomicReference headersObject) { - return new Processor() { + return new Processor<>() { org.apache.kafka.streams.processor.api.ProcessorContext context; diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java index 7690b8963..6bc779c37 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java @@ -195,11 +195,11 @@ public class KafkaStreamsBinderSupportAutoConfiguration { //Making sure that the application indeed set a property. String kafkaStreamsBinderBroker = environment.getProperty("spring.cloud.stream.kafka.streams.binder.brokers"); - if (StringUtils.isEmpty(kafkaStreamsBinderBroker)) { + if (!StringUtils.hasText(kafkaStreamsBinderBroker)) { //Kafka Streams binder specific property for brokers is not set by the application. //See if there is one configured at the kafka binder level. String kafkaBinderBroker = environment.getProperty("spring.cloud.stream.kafka.binder.brokers"); - if (!StringUtils.isEmpty(kafkaBinderBroker)) { + if (StringUtils.hasText(kafkaBinderBroker)) { kafkaConnectionString = kafkaBinderBroker; configProperties.setBrokers(kafkaConnectionString); } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java index 947325812..1299b8b45 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java @@ -128,12 +128,12 @@ final class KafkaStreamsBinderUtils { (cr, e) -> new TopicPartition(dlqDestinationResolvers.values().iterator().next().apply(cr, e), partitionFunction.apply(group, cr, e)); - DeadLetterPublishingRecoverer kafkaStreamsBinderDlqRecoverer = !dlqDestinationResolvers.isEmpty() || !StringUtils - .isEmpty(extendedConsumerProperties.getExtension().getDlqName()) + DeadLetterPublishingRecoverer kafkaStreamsBinderDlqRecoverer = !dlqDestinationResolvers.isEmpty() || StringUtils + .hasText(extendedConsumerProperties.getExtension().getDlqName()) ? new DeadLetterPublishingRecoverer(kafkaTemplate, destinationResolver) : null; for (String inputTopic : inputTopics) { - if (StringUtils.isEmpty( + if (!StringUtils.hasText( extendedConsumerProperties.getExtension().getDlqName()) && dlqDestinationResolvers.isEmpty()) { destinationResolver = (cr, e) -> new TopicPartition("error." + inputTopic + "." + group, partitionFunction.apply(group, cr, e)); diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java index 9d8b1610f..d805d9c14 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java @@ -97,7 +97,7 @@ public class KafkaStreamsMessageConversionDelegate { Message message = v instanceof Message ? (Message) v : MessageBuilder.withPayload(v).build(); Map headers = new HashMap<>(message.getHeaders()); - if (!StringUtils.isEmpty(contentType)) { + if (StringUtils.hasText(contentType)) { headers.put(MessageHeaders.CONTENT_TYPE, contentType); } MessageHeaders messageHeaders = new MessageHeaders(headers); diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/endpoint/KafkaStreamsTopologyEndpoint.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/endpoint/KafkaStreamsTopologyEndpoint.java index 598a0165f..e3e470bdb 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/endpoint/KafkaStreamsTopologyEndpoint.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/endpoint/KafkaStreamsTopologyEndpoint.java @@ -49,7 +49,6 @@ public class KafkaStreamsTopologyEndpoint { @ReadOperation public List kafkaStreamsTopologies() { final List streamsBuilderFactoryBeans = this.kafkaStreamsRegistry.streamsBuilderFactoryBeans(); - final StringBuilder topologyDescription = new StringBuilder(); final List descs = new ArrayList<>(); streamsBuilderFactoryBeans.stream() .forEach(streamsBuilderFactoryBean -> @@ -59,7 +58,7 @@ public class KafkaStreamsTopologyEndpoint { @ReadOperation public String kafkaStreamsTopology(@Selector String applicationId) { - if (!StringUtils.isEmpty(applicationId)) { + if (StringUtils.hasText(applicationId)) { final StreamsBuilderFactoryBean streamsBuilderFactoryBean = this.kafkaStreamsRegistry.streamsBuilderFactoryBean(applicationId); if (streamsBuilderFactoryBean != null) { return streamsBuilderFactoryBean.getTopology().describe().toString(); diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java index 6c1d7bc2e..984d853d2 100644 --- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java +++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java @@ -167,7 +167,7 @@ public class RabbitExchangeQueueProvisioner producerProperties.getExtension(), false)); declareQueue(queue.getName(), queue); String prefix = producerProperties.getExtension().getPrefix(); - String destination = StringUtils.isEmpty(prefix) ? exchangeName + String destination = !StringUtils.hasText(prefix) ? exchangeName : exchangeName.substring(prefix.length()); String[] routingKeys = bindingRoutingKeys(producerProperties.getExtension()); if (ObjectUtils.isEmpty(routingKeys)) {