From 89eb9a31d40130bb86b9c2828c4d103bc7c4d6ec Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Mon, 11 Jul 2022 14:33:18 -0400 Subject: [PATCH] Code cleanup - Kafka Streams binder --- .../AbstractKafkaStreamsBinderProcessor.java | 25 ++++++++----------- .../EncodingDecodingBindAdviceHandler.java | 7 +++--- .../streams/InteractiveQueryService.java | 2 +- .../binder/kafka/streams/KStreamBinder.java | 4 +-- .../KafkaStreamsBinderHealthIndicator.java | 6 ++--- ...StreamsBinderSupportAutoConfiguration.java | 5 ++-- .../KafkaStreamsFunctionProcessor.java | 5 +--- .../kafka/streams/SerdeResolverUtils.java | 4 +-- .../function/FunctionDetectorCondition.java | 4 +-- ...KafkaStreamsFunctionBeanPostProcessor.java | 12 ++++----- 10 files changed, 33 insertions(+), 41 deletions(-) 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 50613fa79..d1ff1bf94 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 @@ -129,16 +129,10 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application .getStartOffset(); Topology.AutoOffsetReset autoOffsetReset = null; if (startOffset != null) { - switch (startOffset) { - case earliest: - autoOffsetReset = Topology.AutoOffsetReset.EARLIEST; - break; - case latest: - autoOffsetReset = Topology.AutoOffsetReset.LATEST; - break; - default: - break; - } + autoOffsetReset = switch (startOffset) { + case earliest -> Topology.AutoOffsetReset.EARLIEST; + case latest -> Topology.AutoOffsetReset.LATEST; + }; } if (extendedConsumerProperties.isResetOffsets()) { AbstractKafkaStreamsBinderProcessor.LOG.warn("Detected resetOffsets configured on binding " @@ -221,11 +215,11 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application final MutablePropertySources propertySources = environment.getPropertySources(); - if (!StringUtils.isEmpty(bindingProperties.getBinder())) { + if (StringUtils.hasText(bindingProperties.getBinder())) { final KafkaStreamsBinderConfigurationProperties multiBinderKafkaStreamsBinderConfigurationProperties = applicationContext.getBean(bindingProperties.getBinder() + "-KafkaStreamsBinderConfigurationProperties", KafkaStreamsBinderConfigurationProperties.class); String connectionString = multiBinderKafkaStreamsBinderConfigurationProperties.getKafkaConnectionString(); - if (StringUtils.isEmpty(connectionString)) { + if (!StringUtils.hasText(connectionString)) { connectionString = (String) propertySources.get(bindingProperties.getBinder() + "-kafkaStreamsBinderEnv").getProperty("spring.cloud.stream.kafka.binder.brokers"); } @@ -587,7 +581,10 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application // Processor to retrieve the header value. stream.process(() -> eventTypeProcessor(kafkaStreamsConsumerProperties, matched, topicObject, headersObject)); // Branching based on event type match. - final KStream[] branch = stream.branch((key, value) -> matched.getAndSet(false)); + final Map> stringKStreamMap = stream.split() + .branch((key, value) -> matched.getAndSet(false)) + .noDefaultBranch(); + final KStream[] branch = stringKStreamMap.values().toArray(new KStream[0]); // Deserialize if we have a branch from above. final KStream deserializedKStream = branch[0].mapValues(value -> valueSerde.deserializer().deserialize( topicObject.get(), headersObject.get(), ((Bytes) value).get())); @@ -600,7 +597,7 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application private Consumed getConsumed(KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties, Serde keySerde, Serde valueSerde, Topology.AutoOffsetReset autoOffsetReset) { TimestampExtractor timestampExtractor = null; - if (!StringUtils.isEmpty(kafkaStreamsConsumerProperties.getTimestampExtractorBeanName())) { + if (StringUtils.hasText(kafkaStreamsConsumerProperties.getTimestampExtractorBeanName())) { timestampExtractor = applicationContext.getBean(kafkaStreamsConsumerProperties.getTimestampExtractorBeanName(), TimestampExtractor.class); } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/EncodingDecodingBindAdviceHandler.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/EncodingDecodingBindAdviceHandler.java index 91e50adee..7ef9d37bc 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/EncodingDecodingBindAdviceHandler.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/EncodingDecodingBindAdviceHandler.java @@ -1,5 +1,5 @@ /* - * Copyright 2019-2019 the original author or authors. + * Copyright 2019-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. @@ -37,7 +37,7 @@ public class EncodingDecodingBindAdviceHandler implements ConfigurationPropertie private boolean decodingSettingProvided; public boolean isDecodingSettingProvided() { - return decodingSettingProvided; + return this.decodingSettingProvided; } public boolean isEncodingSettingProvided() { @@ -46,7 +46,7 @@ public class EncodingDecodingBindAdviceHandler implements ConfigurationPropertie @Override public BindHandler apply(BindHandler bindHandler) { - BindHandler handler = new AbstractBindHandler(bindHandler) { + return new AbstractBindHandler(bindHandler) { @Override public Bindable onStart(ConfigurationPropertyName name, Bindable target, BindContext context) { @@ -67,6 +67,5 @@ public class EncodingDecodingBindAdviceHandler implements ConfigurationPropertie return bindHandler.onStart(name, target, context); } }; - return handler; } } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/InteractiveQueryService.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/InteractiveQueryService.java index 5cd035f13..9dce4fe9e 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/InteractiveQueryService.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/InteractiveQueryService.java @@ -169,7 +169,7 @@ public class InteractiveQueryService { String applicationServer = configuration.get("application.server"); String[] splits = StringUtils.split(applicationServer, ":"); - return new HostInfo(splits[0], Integer.valueOf(splits[1])); + return new HostInfo(splits[0], Integer.parseInt(splits[1])); } return null; } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java index 0e81494f1..e36221964 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java @@ -1,5 +1,5 @@ /* - * Copyright 2017-2021 the original author or authors. + * Copyright 2017-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. @@ -244,7 +244,7 @@ class KStreamBinder extends produced.withName(properties.getProducedAs()); } StreamPartitioner streamPartitioner = null; - if (!StringUtils.isEmpty(properties.getStreamPartitionerBeanName())) { + if (StringUtils.hasText(properties.getStreamPartitionerBeanName())) { streamPartitioner = getApplicationContext().getBean(properties.getStreamPartitionerBeanName(), StreamPartitioner.class); } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java index 16a5a534d..93eb00cf1 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java @@ -1,5 +1,5 @@ /* - * Copyright 2019-2021 the original author or authors. + * Copyright 2019-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. @@ -163,9 +163,9 @@ public class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator i if (isRunningResult) { final Set threadMetadata = kafkaStreams.metadataForLocalThreads(); - final Map threadDetails = new HashMap(); + final Map threadDetails = new HashMap<>(); for (ThreadMetadata metadata : threadMetadata) { - final Map threadDetail = new HashMap(); + final Map threadDetail = new HashMap<>(); threadDetail.put("threadName", metadata.threadName()); threadDetail.put("threadState", metadata.threadState()); threadDetail.put("adminClientId", metadata.adminClientId()); 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 f0482902f..195a8c401 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 @@ -1,5 +1,5 @@ /* - * Copyright 2017-2021 the original author or authors. + * Copyright 2017-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. @@ -222,8 +222,7 @@ public class KafkaStreamsBinderSupportAutoConfiguration { kafkaConnectionString); } } - else if (bootstrapServerConfig instanceof List) { - List bootStrapCollection = (List) bootstrapServerConfig; + else if (bootstrapServerConfig instanceof List bootStrapCollection) { if (bootStrapCollection.size() == 1 && bootStrapCollection.get(0).equals("localhost:9092")) { properties.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaConnectionString); diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java index ff0286c1a..b07a8449d 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java @@ -1,5 +1,5 @@ /* - * Copyright 2019-2021 the original author or authors. + * Copyright 2019-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. @@ -560,9 +560,6 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro throw new IllegalStateException(ex); } } - else { - //throw new IllegalStateException(StreamListenerErrorMessages.INVALID_DECLARATIVE_METHOD_PARAMETERS); - } } return arguments; } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtils.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtils.java index f03368746..c8ceb40e2 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtils.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtils.java @@ -215,8 +215,8 @@ abstract class SerdeResolverUtils { * Private internal class used strictly to 'remember' a score for a serde and use it for sorting later. */ private static class SerdeWithSpecificityScore implements Comparable { - private Integer score; - private String serdeBeanName; + private final Integer score; + private final String serdeBeanName; SerdeWithSpecificityScore(Integer score, String serdeBeanName) { this.score = Objects.requireNonNull(score); diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/FunctionDetectorCondition.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/FunctionDetectorCondition.java index 1f5890f83..7c8f1bf59 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/FunctionDetectorCondition.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/FunctionDetectorCondition.java @@ -1,5 +1,5 @@ /* - * Copyright 2019-2021 the original author or authors. + * Copyright 2019-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. @@ -96,7 +96,7 @@ public class FunctionDetectorCondition extends SpringBootCondition { Method[] methods = classObj.getMethods(); Optional kafkaStreamMethod = Arrays.stream(methods).filter(m -> m.getName().equals(key)).findFirst(); // check if the bean name is overridden. - if (!kafkaStreamMethod.isPresent()) { + if (kafkaStreamMethod.isEmpty()) { final BeanDefinition beanDefinition = context.getBeanFactory().getBeanDefinition(key); final String factoryMethodName = beanDefinition.getFactoryMethodName(); kafkaStreamMethod = Arrays.stream(methods).filter(m -> m.getName().equals(factoryMethodName)).findFirst(); 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 153f19d42..829f8a655 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 @@ -1,5 +1,5 @@ /* - * Copyright 2019-2021 the original author or authors. + * Copyright 2019-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. @@ -65,13 +65,13 @@ public class KafkaStreamsFunctionBeanPostProcessor implements InitializingBean, private ConfigurableListableBeanFactory beanFactory; private boolean onlySingleFunction; - private Map resolvableTypeMap = new TreeMap<>(); - private Map methods = new TreeMap<>(); + private final Map resolvableTypeMap = new TreeMap<>(); + private final Map methods = new TreeMap<>(); private final StreamFunctionProperties streamFunctionProperties; - private Map kafkaStreamsOnlyResolvableTypes = new HashMap<>(); - private Map kafakStreamsOnlyMethods = new HashMap<>(); + private final Map kafkaStreamsOnlyResolvableTypes = new HashMap<>(); + private final Map kafakStreamsOnlyMethods = new HashMap<>(); public KafkaStreamsFunctionBeanPostProcessor(StreamFunctionProperties streamFunctionProperties) { this.streamFunctionProperties = streamFunctionProperties; @@ -177,7 +177,7 @@ public class KafkaStreamsFunctionBeanPostProcessor implements InitializingBean, try { Method[] methods = classObj.getMethods(); Optional functionalBeanMethods = Arrays.stream(methods).filter(m -> m.getName().equals(key)).findFirst(); - if (!functionalBeanMethods.isPresent()) { + if (functionalBeanMethods.isEmpty()) { final BeanDefinition beanDefinition = this.beanFactory.getBeanDefinition(key); final String factoryMethodName = beanDefinition.getFactoryMethodName(); functionalBeanMethods = Arrays.stream(methods).filter(m -> m.getName().equals(factoryMethodName)).findFirst();