Using dot(.) character in function bindings.

In Kafka Streams functions, binding names need to use dot character instead of
underscores as the delimiter.

Resolves #755
This commit is contained in:
Soby Chacko
2019-09-27 15:39:03 -04:00
committed by Oleg Zhurakousky
parent daf4b47d1c
commit 021943ec41
8 changed files with 65 additions and 39 deletions

View File

@@ -267,7 +267,7 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro
}
else {
for (int i = 0; i < outputs; i++) {
outputBindingNames.add(String.format("%s_%s_%d", functionName, KafkaStreamsBindableProxyFactory.DEFAULT_OUTPUT_SUFFIX, i));
outputBindingNames.add(String.format("%s-%s-%d", functionName, KafkaStreamsBindableProxyFactory.DEFAULT_OUTPUT_SUFFIX, i));
}
}
return outputBindingNames;

View File

@@ -85,13 +85,16 @@ public class KafkaStreamsBindableProxyFactory extends AbstractBindableProxyFacto
private final String functionName;
private final boolean onlySingleFunction;
private BeanFactory beanFactory;
public KafkaStreamsBindableProxyFactory(ResolvableType type, String functionName) {
public KafkaStreamsBindableProxyFactory(ResolvableType type, String functionName, boolean onlySingleFunction) {
super(type.getType().getClass());
this.type = type;
this.functionName = functionName;
this.onlySingleFunction = onlySingleFunction;
}
@Override
@@ -140,7 +143,15 @@ public class KafkaStreamsBindableProxyFactory extends AbstractBindableProxyFacto
}
else {
outputBinding = String.format("%s_%s", this.functionName, DEFAULT_OUTPUT_SUFFIX);
int numberOfInputs = this.type.getRawClass() != null &&
(this.type.getRawClass().isAssignableFrom(BiFunction.class) ||
this.type.getRawClass().isAssignableFrom(BiConsumer.class)) ? 2 : getNumberOfInputs();
if (this.onlySingleFunction && numberOfInputs == 1) {
outputBinding = "output";
}
else {
outputBinding = String.format("%s-%s-0", this.functionName, DEFAULT_OUTPUT_SUFFIX);
}
}
Assert.isTrue(outputBinding != null, "output binding is not inferred.");
KafkaStreamsBindableProxyFactory.this.outputHolders.put(outputBinding,
@@ -177,13 +188,31 @@ public class KafkaStreamsBindableProxyFactory extends AbstractBindableProxyFacto
(this.type.getRawClass().isAssignableFrom(BiFunction.class) ||
this.type.getRawClass().isAssignableFrom(BiConsumer.class)) ? 2 : getNumberOfInputs();
if (numberOfInputs == 1) {
inputs.add(String.format("%s_%s", this.functionName, DEFAULT_INPUT_SUFFIX));
ResolvableType outboundArgument = this.type.getGeneric(1);
while (isAnotherFunctionOrConsumerFound(outboundArgument)) {
//The function is a curried function. We should introspect the partial function chain hierarchy.
outboundArgument = outboundArgument.getGeneric(1);
}
if (this.onlySingleFunction && (outboundArgument == null || outboundArgument.getRawClass() == null)) {
inputs.add("input");
}
else if (this.onlySingleFunction && outboundArgument.getRawClass() != null
&& (!outboundArgument.isArray() &&
outboundArgument.getRawClass().isAssignableFrom(KStream.class))) {
inputs.add("input");
}
else {
inputs.add(String.format("%s-%s-0", this.functionName, DEFAULT_INPUT_SUFFIX));
}
return inputs;
}
else {
int i = 0;
while (i < numberOfInputs) {
inputs.add(String.format("%s_%s_%d", this.functionName, DEFAULT_INPUT_SUFFIX, i++));
inputs.add(String.format("%s-%s-%d", this.functionName, DEFAULT_INPUT_SUFFIX, i++));
}
return inputs;
}

View File

@@ -66,6 +66,8 @@ public class KafkaStreamsFunctionAutoConfiguration {
.addGenericArgumentValue(kafkaStreamsFunctionBeanPostProcessor.getResolvableTypes().get(s));
rootBeanDefinition.getConstructorArgumentValues()
.addGenericArgumentValue(s);
rootBeanDefinition.getConstructorArgumentValues()
.addGenericArgumentValue(kafkaStreamsFunctionBeanPostProcessor.getResolvableTypes().size() == 1);
registry.registerBeanDefinition("kafkaStreamsBindableProxyFactory-" + s, rootBeanDefinition);
}
};

View File

@@ -80,8 +80,6 @@ public class KafkaStreamsBinderWordCountFunctionTests {
app.setWebApplicationType(WebApplicationType.NONE);
try (ConfigurableApplicationContext context = app.run(
"--spring.cloud.stream.function.inputBindings.process=input",
"--spring.cloud.stream.function.outputBindings.process=output",
"--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.destination=words",
@@ -103,8 +101,6 @@ public class KafkaStreamsBinderWordCountFunctionTests {
app.setWebApplicationType(WebApplicationType.NONE);
try (ConfigurableApplicationContext context = app.run(
"--spring.cloud.stream.function.inputBindings.process=input",
"--spring.cloud.stream.function.outputBindings.process=output",
"--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.destination=words-1",

View File

@@ -58,7 +58,7 @@ public class KafkaStreamsFunctionStateStoreTests {
try (ConfigurableApplicationContext context = app.run("--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.process_in.destination=words",
"--spring.cloud.stream.bindings.input.destination=words",
"--spring.cloud.stream.kafka.streams.default.consumer.application-id=testKafkaStreamsFuncionWithMultipleStateStores",
"--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000",
"--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" +

View File

@@ -81,11 +81,11 @@ public class MultipleFunctionsInSameAppTests {
try (ConfigurableApplicationContext ignored = app.run(
"--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.process_in.destination=purchases",
"--spring.cloud.stream.bindings.process_out_0.destination=coffee",
"--spring.cloud.stream.bindings.process_out_1.destination=electronics",
"--spring.cloud.stream.bindings.analyze_in_0.destination=coffee",
"--spring.cloud.stream.bindings.analyze_in_1.destination=electronics",
"--spring.cloud.stream.bindings.process-in-0.destination=purchases",
"--spring.cloud.stream.bindings.process-out-0.destination=coffee",
"--spring.cloud.stream.bindings.process-out-1.destination=electronics",
"--spring.cloud.stream.bindings.analyze-in-0.destination=coffee",
"--spring.cloud.stream.bindings.analyze-in-1.destination=electronics",
"--spring.cloud.stream.kafka.streams.binder.functions.analyze.applicationId=analyze-id-0",
"--spring.cloud.stream.kafka.streams.binder.functions.process.applicationId=process-id-0",
"--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000",

View File

@@ -59,8 +59,8 @@ public class SerdesProvidedAsBeansTests {
try (ConfigurableApplicationContext context = app.run(
"--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.process_in.destination=purchases",
"--spring.cloud.stream.bindings.process_out.destination=coffee",
"--spring.cloud.stream.bindings.input.destination=purchases",
"--spring.cloud.stream.bindings.output.destination=coffee",
"--spring.cloud.stream.kafka.streams.binder.functions.process.applicationId=process-id-0",
"--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000",
"--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" +
@@ -77,16 +77,16 @@ public class SerdesProvidedAsBeansTests {
final BindingServiceProperties bindingServiceProperties = context.getBean(BindingServiceProperties.class);
final KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties = context.getBean(KafkaStreamsExtendedBindingProperties.class);
final ConsumerProperties consumerProperties = bindingServiceProperties.getBindingProperties("process_in").getConsumer();
final KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties = kafkaStreamsExtendedBindingProperties.getExtendedConsumerProperties("process_in");
kafkaStreamsExtendedBindingProperties.getExtendedConsumerProperties("process_in");
final ConsumerProperties consumerProperties = bindingServiceProperties.getBindingProperties("input").getConsumer();
final KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties = kafkaStreamsExtendedBindingProperties.getExtendedConsumerProperties("input");
kafkaStreamsExtendedBindingProperties.getExtendedConsumerProperties("input");
final Serde<?> inboundValueSerde = keyValueSerdeResolver.getInboundValueSerde(consumerProperties, kafkaStreamsConsumerProperties, resolvableType.getGeneric(0));
Assert.isTrue(inboundValueSerde instanceof FooSerde, "Inbound Value Serde is not matched");
final ProducerProperties producerProperties = bindingServiceProperties.getBindingProperties("process_out").getProducer();
final KafkaStreamsProducerProperties kafkaStreamsProducerProperties = kafkaStreamsExtendedBindingProperties.getExtendedProducerProperties("process_out");
kafkaStreamsExtendedBindingProperties.getExtendedProducerProperties("process_out");
final ProducerProperties producerProperties = bindingServiceProperties.getBindingProperties("output").getProducer();
final KafkaStreamsProducerProperties kafkaStreamsProducerProperties = kafkaStreamsExtendedBindingProperties.getExtendedProducerProperties("output");
kafkaStreamsExtendedBindingProperties.getExtendedProducerProperties("output");
final Serde<?> outboundValueSerde = keyValueSerdeResolver.getOutboundValueSerde(producerProperties, kafkaStreamsProducerProperties, resolvableType.getGeneric(1));
Assert.isTrue(outboundValueSerde instanceof FooSerde, "Outbound Value Serde is not matched");

View File

@@ -120,14 +120,14 @@ public class StreamToTableJoinFunctionTests {
try (ConfigurableApplicationContext ignored = app.run("--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.process_in_0.destination=user-clicks-1",
"--spring.cloud.stream.bindings.process_in_1.destination=user-regions-1",
"--spring.cloud.stream.bindings.process-in-0.destination=user-clicks-1",
"--spring.cloud.stream.bindings.process-in-1.destination=user-regions-1",
"--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" +
"=org.apache.kafka.common.serialization.Serdes$StringSerde",
"--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" +
"=org.apache.kafka.common.serialization.Serdes$StringSerde",
"--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=10000",
"--spring.cloud.stream.kafka.streams.bindings.process_in_0.consumer.applicationId" +
"--spring.cloud.stream.kafka.streams.bindings.process-in-0.consumer.applicationId" +
"=testStreamToTableBiConsumer",
"--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) {
@@ -177,15 +177,15 @@ public class StreamToTableJoinFunctionTests {
private void runTest(SpringApplication app, Consumer<String, Long> consumer) {
try (ConfigurableApplicationContext ignored = app.run("--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.process_in_0.destination=user-clicks-1",
"--spring.cloud.stream.bindings.process_in_1.destination=user-regions-1",
"--spring.cloud.stream.bindings.process_out.destination=output-topic-1",
"--spring.cloud.stream.bindings.process-in-0.destination=user-clicks-1",
"--spring.cloud.stream.bindings.process-in-1.destination=user-regions-1",
"--spring.cloud.stream.bindings.process-out-0.destination=output-topic-1",
"--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" +
"=org.apache.kafka.common.serialization.Serdes$StringSerde",
"--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" +
"=org.apache.kafka.common.serialization.Serdes$StringSerde",
"--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=10000",
"--spring.cloud.stream.kafka.streams.bindings.process_in_0.consumer.applicationId" +
"--spring.cloud.stream.kafka.streams.bindings.process-in-0.consumer.applicationId" +
"=StreamToTableJoinFunctionTests-abc",
"--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) {
@@ -307,21 +307,20 @@ public class StreamToTableJoinFunctionTests {
try (ConfigurableApplicationContext context = app.run("--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.function.inputBindings.process=input-1,input-2",
"--spring.cloud.stream.bindings.input-1.destination=user-clicks-2",
"--spring.cloud.stream.bindings.input-2.destination=user-regions-2",
"--spring.cloud.stream.bindings.process_out.destination=output-topic-2",
"--spring.cloud.stream.bindings.input-1.consumer.useNativeDecoding=true",
"--spring.cloud.stream.bindings.input-2.consumer.useNativeDecoding=true",
"--spring.cloud.stream.bindings.output.producer.useNativeEncoding=true",
"--spring.cloud.stream.bindings.process-in-0.destination=user-clicks-2",
"--spring.cloud.stream.bindings.process-in-1.destination=user-regions-2",
"--spring.cloud.stream.bindings.process-out-0.destination=output-topic-2",
"--spring.cloud.stream.bindings.process-in-0.consumer.useNativeDecoding=true",
"--spring.cloud.stream.bindings.process-in-1.consumer.useNativeDecoding=true",
"--spring.cloud.stream.bindings.process-out-0.producer.useNativeEncoding=true",
"--spring.cloud.stream.kafka.streams.binder.configuration.auto.offset.reset=latest",
"--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.startOffset=earliest",
"--spring.cloud.stream.kafka.streams.bindings.process-in-0.consumer.startOffset=earliest",
"--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" +
"=org.apache.kafka.common.serialization.Serdes$StringSerde",
"--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" +
"=org.apache.kafka.common.serialization.Serdes$StringSerde",
"--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=10000",
"--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.application-id" +
"--spring.cloud.stream.kafka.streams.bindings.process-in-0.consumer.application-id" +
"=StreamToTableJoinFunctionTests-foobar",
"--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(),
"--spring.cloud.stream.kafka.streams.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString())) {