Change name of property endpoint -> defaultRoute
This commit is contained in:
@@ -1,4 +1,4 @@
|
|||||||
spring.cloud.stream.bindings.input.destination: foos
|
spring.cloud.stream.bindings.input.destination: foos
|
||||||
spring.cloud.stream.bindings.output.destination: bars
|
spring.cloud.stream.bindings.output.destination: bars
|
||||||
spring.cloud.function.stream.endpoint: uppercase
|
spring.cloud.function.stream.default-route: uppercase
|
||||||
management.security.enabled: false
|
management.security.enabled: false
|
||||||
@@ -1,2 +1,2 @@
|
|||||||
spring.cloud.function.stream.endpoint: uppercase
|
spring.cloud.function.stream.default-route: uppercase
|
||||||
spring.cloud.function.scan.packages: com.example.functions
|
spring.cloud.function.scan.packages: com.example.functions
|
||||||
@@ -60,7 +60,7 @@ public class StreamConfiguration {
|
|||||||
FunctionInspector functionInspector,
|
FunctionInspector functionInspector,
|
||||||
@Lazy CompositeMessageConverterFactory compositeMessageConverterFactory) {
|
@Lazy CompositeMessageConverterFactory compositeMessageConverterFactory) {
|
||||||
return new StreamListeningFunctionInvoker(registry, functionInspector,
|
return new StreamListeningFunctionInvoker(registry, functionInspector,
|
||||||
compositeMessageConverterFactory, properties.getEndpoint());
|
compositeMessageConverterFactory, properties.getDefaultRoute());
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -24,7 +24,11 @@ import org.springframework.boot.context.properties.ConfigurationProperties;
|
|||||||
@ConfigurationProperties(prefix = "spring.cloud.function.stream")
|
@ConfigurationProperties(prefix = "spring.cloud.function.stream")
|
||||||
public class StreamConfigurationProperties {
|
public class StreamConfigurationProperties {
|
||||||
|
|
||||||
private String endpoint;
|
/**
|
||||||
|
* The default route for a message if more than one is available and no explicit route
|
||||||
|
* key is provided.
|
||||||
|
*/
|
||||||
|
private String defaultRoute;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Interval to be used for the Duration (in milliseconds) of a non-Flux producing
|
* Interval to be used for the Duration (in milliseconds) of a non-Flux producing
|
||||||
@@ -34,12 +38,12 @@ public class StreamConfigurationProperties {
|
|||||||
|
|
||||||
public static final String ROUTE_KEY = "stream_routekey";
|
public static final String ROUTE_KEY = "stream_routekey";
|
||||||
|
|
||||||
public String getEndpoint() {
|
public String getDefaultRoute() {
|
||||||
return endpoint;
|
return defaultRoute;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void setEndpoint(String endpoint) {
|
public void setDefaultRoute(String defaultRoute) {
|
||||||
this.endpoint = endpoint;
|
this.defaultRoute = defaultRoute;
|
||||||
}
|
}
|
||||||
|
|
||||||
public long getInterval() {
|
public long getInterval() {
|
||||||
|
|||||||
@@ -55,7 +55,7 @@ public class StreamListeningFunctionInvoker implements SmartInitializingSingleto
|
|||||||
|
|
||||||
private MessageConverter converter;
|
private MessageConverter converter;
|
||||||
|
|
||||||
private final String defaultEndpoint;
|
private final String defaultRoute;
|
||||||
|
|
||||||
private final Map<String, FluxMessageProcessor> processors = new HashMap<>();
|
private final Map<String, FluxMessageProcessor> processors = new HashMap<>();
|
||||||
|
|
||||||
@@ -67,11 +67,11 @@ public class StreamListeningFunctionInvoker implements SmartInitializingSingleto
|
|||||||
|
|
||||||
public StreamListeningFunctionInvoker(FunctionCatalog functionCatalog,
|
public StreamListeningFunctionInvoker(FunctionCatalog functionCatalog,
|
||||||
FunctionInspector functionInspector,
|
FunctionInspector functionInspector,
|
||||||
CompositeMessageConverterFactory converterFactory, String defaultEndpoint) {
|
CompositeMessageConverterFactory converterFactory, String defaultRoute) {
|
||||||
this.functionCatalog = functionCatalog;
|
this.functionCatalog = functionCatalog;
|
||||||
this.functionInspector = functionInspector;
|
this.functionInspector = functionInspector;
|
||||||
this.converterFactory = converterFactory;
|
this.converterFactory = converterFactory;
|
||||||
this.defaultEndpoint = defaultEndpoint;
|
this.defaultRoute = defaultRoute;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
@@ -137,8 +137,8 @@ public class StreamListeningFunctionInvoker implements SmartInitializingSingleto
|
|||||||
.get(StreamConfigurationProperties.ROUTE_KEY);
|
.get(StreamConfigurationProperties.ROUTE_KEY);
|
||||||
name = stash(key);
|
name = stash(key);
|
||||||
}
|
}
|
||||||
if (name==null && defaultEndpoint != null) {
|
if (name==null && defaultRoute != null) {
|
||||||
name = stash(defaultEndpoint);
|
name = stash(defaultRoute);
|
||||||
}
|
}
|
||||||
if (name == null) {
|
if (name == null) {
|
||||||
Set<String> names = new LinkedHashSet<>(functionCatalog.getFunctionNames());
|
Set<String> names = new LinkedHashSet<>(functionCatalog.getFunctionNames());
|
||||||
|
|||||||
@@ -44,7 +44,7 @@ import static org.assertj.core.api.Assertions.assertThat;
|
|||||||
*/
|
*/
|
||||||
@RunWith(SpringRunner.class)
|
@RunWith(SpringRunner.class)
|
||||||
@SpringBootTest(classes = PojoStreamingExplicitEndpointTests.StreamingFunctionApplication.class, properties = {
|
@SpringBootTest(classes = PojoStreamingExplicitEndpointTests.StreamingFunctionApplication.class, properties = {
|
||||||
"spring.cloud.function.stream.endpoint=uppercase",
|
"spring.cloud.function.stream.default-route=uppercase",
|
||||||
"logging.level.org.springframework.integration=DEBUG", "debug=TRUE" })
|
"logging.level.org.springframework.integration=DEBUG", "debug=TRUE" })
|
||||||
public class PojoStreamingExplicitEndpointTests {
|
public class PojoStreamingExplicitEndpointTests {
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user