Merge pull request #845 from sobychacko/gh-844

Kafka streams concurrency with multiple bindings
This commit is contained in:
Oleg Zhurakousky
2020-02-12 16:38:53 +01:00
committed by GitHub
2 changed files with 16 additions and 1 deletions

View File

@@ -265,7 +265,12 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application
int concurrency = this.bindingServiceProperties.getConsumerProperties(inboundName)
.getConcurrency();
// override concurrency if set at the individual binding level.
if (concurrency > 1) {
// Concurrency will be mapped to num.stream.threads. Since this is going into a global config,
// we are explicitly assigning concurrency left at default of 1 to num.stream.threads. Otherwise,
// a potential previous value might still be used in the case of multiple processors or a processor
// with multiple input bindings with various concurrency values.
// See this GH issue: https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/844
if (concurrency >= 1) {
streamConfigGlobalProperties.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG,
concurrency);
}

View File

@@ -26,6 +26,7 @@ import java.util.function.Function;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.KStream;
import org.junit.AfterClass;
import org.junit.BeforeClass;
@@ -92,6 +93,8 @@ public class MultipleFunctionsInSameAppTests {
"--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",
"--spring.cloud.stream.bindings.process-in-0.consumer.concurrency=2",
"--spring.cloud.stream.bindings.analyze-in-0.consumer.concurrency=1",
"--spring.cloud.stream.kafka.streams.binder.functions.process.configuration.client.id=process-client",
"--spring.cloud.stream.kafka.streams.binder.functions.analyze.configuration.client.id=analyze-client",
"--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) {
@@ -108,6 +111,13 @@ public class MultipleFunctionsInSameAppTests {
assertThat(processStreamsConfiguration.getProperty("client.id")).isEqualTo("process-client");
assertThat(analyzeStreamsConfiguration.getProperty("client.id")).isEqualTo("analyze-client");
Integer concurrency = (Integer) processStreamsBuilderFactoryBean.getStreamsConfiguration()
.get(StreamsConfig.NUM_STREAM_THREADS_CONFIG);
assertThat(concurrency).isEqualTo(2);
concurrency = (Integer) analyzeStreamsBuilderFactoryBean.getStreamsConfiguration()
.get(StreamsConfig.NUM_STREAM_THREADS_CONFIG);
assertThat(concurrency).isEqualTo(1);
}
}