GH-2783 Ensure proper cashing of StreamBridge function

Resolves #2783
This commit is contained in:
Oleg Zhurakousky
2023-08-09 15:58:47 +02:00
parent b0fb7740b7
commit b4e976f371
2 changed files with 44 additions and 4 deletions

View File

@@ -179,6 +179,36 @@ public class StreamBridgeTests {
}
}
@SuppressWarnings("rawtypes")
@Test
void test_2785() throws Exception {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
EmptyConfiguration.class)).web(WebApplicationType.NONE).run(
"--spring.cloud.stream.source=outputA;outputB",
"--spring.cloud.stream.bindings.outputA-out-0.producer.use-native-encoding=false",
"--spring.cloud.stream.bindings.outputB-out-0.producer.partition-count=1",
"--spring.cloud.stream.bindings.outputC-out-0.producer.use-native-encoding=true",
"--spring.cloud.stream.bindings.outputD-out-0.content-type=text/html",
"--spring.jmx.enabled=false")) {
StreamBridge streamBridge = context.getBean(StreamBridge.class);
Field field = ReflectionUtils.findField(StreamBridge.class, "streamBridgeFunctionCache");
field.setAccessible(true);
Map functionCache = (Map) field.get(streamBridge);
streamBridge.send("outputA-out-0", MessageBuilder.withPayload("A").build());
assertThat(functionCache.size()).isEqualTo(1);
streamBridge.send("foo", MessageBuilder.withPayload("A").build());
assertThat(functionCache.size()).isEqualTo(1);
streamBridge.send("outputB-out-0", MessageBuilder.withPayload("A").build());
assertThat(functionCache.size()).isEqualTo(1);
streamBridge.send("outputC-out-0", MessageBuilder.withPayload("A").build());
assertThat(functionCache.size()).isEqualTo(2);
streamBridge.send("outputD-out-0", MessageBuilder.withPayload("A").build());
assertThat(functionCache.size()).isEqualTo(3);
}
}
/*
* This test verifies that when a partition key expression is set, then scst_partition is always set, even in
* concurrent scenarios.

View File

@@ -96,7 +96,7 @@ public final class StreamBridge implements StreamOperations, SmartInitializingSi
private final BindingService bindingService;
private final Map<String, FunctionInvocationWrapper> streamBridgeFunctionCache;
private final Map<Integer, FunctionInvocationWrapper> streamBridgeFunctionCache;
private final FunctionInvocationHelper<?> functionInvocationHelper;
@@ -187,13 +187,23 @@ public final class StreamBridge implements StreamOperations, SmartInitializingSi
return messageChannel.send(resultMessage);
}
private int hashProducerProperties(ProducerProperties producerProperties, String outputContentType) {
int hash = outputContentType.hashCode()
+ Boolean.hashCode(producerProperties.isUseNativeEncoding())
+ Boolean.hashCode(producerProperties.isPartitioned())
+ producerProperties.getPartitionCount();
return hash;
}
private synchronized FunctionInvocationWrapper getStreamBridgeFunction(String outputContentType, ProducerProperties producerProperties) {
if (StringUtils.hasText(outputContentType) && this.streamBridgeFunctionCache.containsKey(outputContentType)) {
return this.streamBridgeFunctionCache.get(outputContentType);
int streamBridgeFunctionKey = this.hashProducerProperties(producerProperties, outputContentType);
if (this.streamBridgeFunctionCache.containsKey(streamBridgeFunctionKey)) {
return this.streamBridgeFunctionCache.get(streamBridgeFunctionKey);
}
else {
FunctionInvocationWrapper functionToInvoke = this.functionCatalog.lookup(STREAM_BRIDGE_FUNC_NAME, outputContentType.toString());
this.streamBridgeFunctionCache.put(outputContentType, functionToInvoke);
this.streamBridgeFunctionCache.put(streamBridgeFunctionKey, functionToInvoke);
functionToInvoke.setSkipOutputConversion(producerProperties.isUseNativeEncoding());
return functionToInvoke;
}