From 3fc30978526664f875e4f94a0d6997cb75b89412 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 29 Aug 2022 16:51:42 +0200 Subject: [PATCH] Small cleanup in StreamBridge Remove redundant injection of FunctionRegistry since it is the same as FunctionCatalog --- .../stream/function/FunctionConfiguration.java | 4 ++-- .../cloud/stream/function/StreamBridge.java | 16 +++------------- 2 files changed, 5 insertions(+), 15 deletions(-) diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index 267fe9bd8..7f05cb277 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -137,10 +137,10 @@ public class FunctionConfiguration { private final static String SOURCE_PROPERY = "spring.cloud.stream.source"; @Bean - public StreamBridge streamBridgeUtils(FunctionCatalog functionCatalog, FunctionRegistry functionRegistry, + public StreamBridge streamBridgeUtils(FunctionCatalog functionCatalog, BindingServiceProperties bindingServiceProperties, ConfigurableApplicationContext applicationContext, @Nullable NewDestinationBindingCallback callback) { - return new StreamBridge(functionCatalog, functionRegistry, bindingServiceProperties, applicationContext, callback); + return new StreamBridge(functionCatalog, bindingServiceProperties, applicationContext, callback); } @Bean 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 d3025a919..0d8f8e464 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 @@ -82,8 +82,6 @@ public final class StreamBridge implements SmartInitializingSingleton { private final FunctionCatalog functionCatalog; - private final FunctionRegistry functionRegistry; - private final NewDestinationBindingCallback destinationBindingCallback; private BindingServiceProperties bindingServiceProperties; @@ -101,17 +99,14 @@ public final class StreamBridge implements SmartInitializingSingleton { /** * * @param functionCatalog instance of {@link FunctionCatalog} - * @param functionRegistry instance of {@link FunctionRegistry} * @param bindingServiceProperties instance of {@link BindingServiceProperties} * @param applicationContext instance of {@link ConfigurableApplicationContext} */ @SuppressWarnings("serial") - StreamBridge(FunctionCatalog functionCatalog, FunctionRegistry functionRegistry, - BindingServiceProperties bindingServiceProperties, ConfigurableApplicationContext applicationContext, - @Nullable NewDestinationBindingCallback destinationBindingCallback) { + StreamBridge(FunctionCatalog functionCatalog, BindingServiceProperties bindingServiceProperties, + ConfigurableApplicationContext applicationContext, @Nullable NewDestinationBindingCallback destinationBindingCallback) { this.bindingService = applicationContext.getBean(BindingService.class); this.functionCatalog = functionCatalog; - this.functionRegistry = functionRegistry; this.applicationContext = applicationContext; this.bindingServiceProperties = bindingServiceProperties; this.destinationBindingCallback = destinationBindingCallback; @@ -257,13 +252,8 @@ public final class StreamBridge implements SmartInitializingSingleton { } FunctionRegistration> fr = new FunctionRegistration<>(v -> v, STREAM_BRIDGE_FUNC_NAME); fr.getProperties().put("singleton", "false"); - this.functionRegistry.register(fr.type(FunctionType.from(Object.class).to(Object.class).message())); + ((FunctionRegistry) this.functionCatalog).register(fr.type(FunctionType.from(Object.class).to(Object.class).message())); Map channels = applicationContext.getBeansOfType(DirectWithAttributesChannel.class); -// for (Entry channelEntry : channels.entrySet()) { -// if (channelEntry.getValue().getAttribute("type").equals("output")) { -// this.channelCache.put(channelEntry.getKey(), channelEntry.getValue()); -// } -// } this.initialized = true; }