From 895570fb48269b4cce689d629395d6fc55c426ed Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 29 Sep 2023 15:33:09 -0400 Subject: [PATCH] Code cleanup: AbstractKafkaStreamsBinderProcessor - Remove deprecation warnings - Other code cleanup --- .../AbstractKafkaStreamsBinderProcessor.java | 21 ++++++++++++------- 1 file changed, 13 insertions(+), 8 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 f1fc8cb19..85145e2ec 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 @@ -19,6 +19,7 @@ package org.springframework.cloud.stream.binder.kafka.streams; import java.util.Arrays; import java.util.HashMap; import java.util.Map; +import java.util.Objects; import java.util.Optional; import java.util.concurrent.atomic.AtomicReference; import java.util.regex.Pattern; @@ -79,6 +80,7 @@ import org.springframework.kafka.config.StreamsBuilderFactoryBean; import org.springframework.kafka.config.StreamsBuilderFactoryBeanConfigurer; import org.springframework.kafka.core.CleanupConfig; import org.springframework.kafka.streams.RecoveringDeserializationExceptionHandler; +import org.springframework.lang.NonNull; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.support.MessageBuilder; import org.springframework.util.Assert; @@ -120,7 +122,7 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application } @Override - public final void setApplicationContext(ApplicationContext applicationContext) + public final void setApplicationContext(@NonNull ApplicationContext applicationContext) throws BeansException { this.applicationContext = (ConfigurableApplicationContext) applicationContext; } @@ -221,7 +223,7 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application applicationContext.getBean(bindingProperties.getBinder() + "-KafkaStreamsBinderConfigurationProperties", KafkaStreamsBinderConfigurationProperties.class); String connectionString = multiBinderKafkaStreamsBinderConfigurationProperties.getKafkaConnectionString(); if (!StringUtils.hasText(connectionString)) { - connectionString = (String) propertySources.get(bindingProperties.getBinder() + "-kafkaStreamsBinderEnv").getProperty("spring.cloud.stream.kafka.binder.brokers"); + connectionString = (String) Objects.requireNonNull(propertySources.get(bindingProperties.getBinder() + "-kafkaStreamsBinderEnv")).getProperty("spring.cloud.stream.kafka.binder.brokers"); } streamConfiguration.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, connectionString); @@ -450,15 +452,18 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application AtomicReference topicObject = new AtomicReference<>(); AtomicReference headersObject = new AtomicReference<>(); // Processor to retrieve the header value. - stream.process(() -> eventTypeProcessor(kafkaStreamsConsumerProperties, matchedRecordThreadLocal, topicObject, headersObject)); + stream.process(() -> eventTypeProcessor(kafkaStreamsConsumerProperties, topicObject, headersObject)); // Branching based on event type match. - final KStream[] branch = stream.branch((key, value) -> { + final Map> stringKStreamMap = stream.split() + .branch((key, value) -> { if (matchedRecordThreadLocal.get()) { matchedRecordThreadLocal.set(false); return true; } return false; - }); + }) + .noDefaultBranch(); + final KStream[] branch = stringKStreamMap.values().toArray(new KStream[0]); // Deserialize if we have a branch from above. final KStream deserializedKStream = !kafkaStreamsConsumerProperties.isUseConfiguredSerdeWhenRoutingEvents() ? branch[0].mapValues(value -> valueSerde.deserializer().deserialize( @@ -591,7 +596,7 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application final KStream stream = kTable.toStream(); // Processor to retrieve the header value. - stream.process(() -> eventTypeProcessor(kafkaStreamsConsumerProperties, matchedRecordThreadLocal, topicObject, headersObject)); + stream.process(() -> eventTypeProcessor(kafkaStreamsConsumerProperties, topicObject, headersObject)); // Branching based on event type match. final Map> stringKStreamMap = stream.split() .branch((key, value) -> { @@ -632,7 +637,7 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application } private Processor eventTypeProcessor(KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties, - ThreadLocal matchedValHolder, AtomicReference topicObject, AtomicReference headersObject) { + AtomicReference topicObject, AtomicReference headersObject) { return new Processor<>() { org.apache.kafka.streams.processor.api.ProcessorContext context; @@ -658,7 +663,7 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application final String[] eventTypesFromBinding = StringUtils.commaDelimitedListToStringArray(kafkaStreamsConsumerProperties.getEventTypes()); for (String eventTypeFromBinding : eventTypesFromBinding) { if (eventTypeFromHeader.equals(eventTypeFromBinding)) { - matchedValHolder.set(true); + AbstractKafkaStreamsBinderProcessor.matchedRecordThreadLocal.set(true); break; } }