From ccc26f88632333b40b5e68243a67c208727ec2e8 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 | 17 ++++------------- 2 files changed, 6 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 6764184de..9ab2b4aaf 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 @@ -134,10 +134,10 @@ public class FunctionConfiguration { @SuppressWarnings("rawtypes") @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 947258d0a..09f0b4992 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 @@ -84,8 +84,6 @@ public final class StreamBridge implements SmartInitializingSingleton { private final FunctionCatalog functionCatalog; - private final FunctionRegistry functionRegistry; - private final NewDestinationBindingCallback destinationBindingCallback; private BindingServiceProperties bindingServiceProperties; @@ -103,17 +101,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; @@ -260,14 +255,10 @@ public final class StreamBridge implements SmartInitializingSingleton { FunctionRegistration> fr = new FunctionRegistration<>(new PassThruFunction(), STREAM_BRIDGE_FUNC_NAME); fr.getProperties().put("singleton", "false"); + Type functionType = ResolvableType.forClassWithGenerics(Function.class, Object.class, Object.class).getType(); - this.functionRegistry.register(fr.type(functionType)); + ((FunctionRegistry) this.functionCatalog).register(fr.type(functionType)); 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; }