diff --git a/pom.xml b/pom.xml
index e065e0835..af349050b 100644
--- a/pom.xml
+++ b/pom.xml
@@ -121,6 +121,14 @@
+
+
+ org.junit.vintage
+ junit-vintage-engine
+ test
+
+
+
diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderAutoConfigurationPropertiesTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderAutoConfigurationPropertiesTest.java
index 84ee1c71d..e7023239d 100644
--- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderAutoConfigurationPropertiesTest.java
+++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderAutoConfigurationPropertiesTest.java
@@ -65,10 +65,10 @@ public class KafkaBinderAutoConfigurationPropertiesTest {
new KafkaProducerProperties());
Method getProducerFactoryMethod = KafkaMessageChannelBinder.class
.getDeclaredMethod("getProducerFactory", String.class,
- ExtendedProducerProperties.class);
+ ExtendedProducerProperties.class, String.class);
getProducerFactoryMethod.setAccessible(true);
DefaultKafkaProducerFactory producerFactory = (DefaultKafkaProducerFactory) getProducerFactoryMethod
- .invoke(this.kafkaMessageChannelBinder, "foo", producerProperties);
+ .invoke(this.kafkaMessageChannelBinder, "foo", producerProperties, "foo.producer");
Field producerFactoryConfigField = ReflectionUtils
.findField(DefaultKafkaProducerFactory.class, "configs", Map.class);
ReflectionUtils.makeAccessible(producerFactoryConfigField);
@@ -88,12 +88,12 @@ public class KafkaBinderAutoConfigurationPropertiesTest {
.containsAll(bootstrapServers))).isTrue();
Method createKafkaConsumerFactoryMethod = KafkaMessageChannelBinder.class
.getDeclaredMethod("createKafkaConsumerFactory", boolean.class,
- String.class, ExtendedConsumerProperties.class);
+ String.class, ExtendedConsumerProperties.class, String.class);
createKafkaConsumerFactoryMethod.setAccessible(true);
ExtendedConsumerProperties consumerProperties = new ExtendedConsumerProperties<>(
new KafkaConsumerProperties());
DefaultKafkaConsumerFactory consumerFactory = (DefaultKafkaConsumerFactory) createKafkaConsumerFactoryMethod
- .invoke(this.kafkaMessageChannelBinder, true, "test", consumerProperties);
+ .invoke(this.kafkaMessageChannelBinder, true, "test", consumerProperties, "test.consumer");
Field consumerFactoryConfigField = ReflectionUtils
.findField(DefaultKafkaConsumerFactory.class, "configs", Map.class);
ReflectionUtils.makeAccessible(consumerFactoryConfigField);
diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java
index a6e65027e..d7d8a8a83 100644
--- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java
+++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java
@@ -70,10 +70,10 @@ public class KafkaBinderConfigurationPropertiesTest {
kafkaProducerProperties);
Method getProducerFactoryMethod = KafkaMessageChannelBinder.class
.getDeclaredMethod("getProducerFactory", String.class,
- ExtendedProducerProperties.class);
+ ExtendedProducerProperties.class, String.class);
getProducerFactoryMethod.setAccessible(true);
DefaultKafkaProducerFactory producerFactory = (DefaultKafkaProducerFactory) getProducerFactoryMethod
- .invoke(this.kafkaMessageChannelBinder, "bar", producerProperties);
+ .invoke(this.kafkaMessageChannelBinder, "bar", producerProperties, "bar.producer");
Field producerFactoryConfigField = ReflectionUtils
.findField(DefaultKafkaProducerFactory.class, "configs", Map.class);
ReflectionUtils.makeAccessible(producerFactoryConfigField);
@@ -100,12 +100,12 @@ public class KafkaBinderConfigurationPropertiesTest {
.contains("10.98.09.199:9082"))).isTrue();
Method createKafkaConsumerFactoryMethod = KafkaMessageChannelBinder.class
.getDeclaredMethod("createKafkaConsumerFactory", boolean.class,
- String.class, ExtendedConsumerProperties.class);
+ String.class, ExtendedConsumerProperties.class, String.class);
createKafkaConsumerFactoryMethod.setAccessible(true);
ExtendedConsumerProperties consumerProperties = new ExtendedConsumerProperties<>(
new KafkaConsumerProperties());
DefaultKafkaConsumerFactory consumerFactory = (DefaultKafkaConsumerFactory) createKafkaConsumerFactoryMethod
- .invoke(this.kafkaMessageChannelBinder, true, "test", consumerProperties);
+ .invoke(this.kafkaMessageChannelBinder, true, "test", consumerProperties, "test.consumer");
Field consumerFactoryConfigField = ReflectionUtils
.findField(DefaultKafkaConsumerFactory.class, "configs", Map.class);
ReflectionUtils.makeAccessible(consumerFactoryConfigField);
diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderUnitTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderUnitTests.java
index 0eb4873a1..ef7a71511 100644
--- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderUnitTests.java
+++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderUnitTests.java
@@ -86,17 +86,17 @@ public class KafkaBinderUnitTests {
consumerProps);
Method method = KafkaMessageChannelBinder.class.getDeclaredMethod(
"createKafkaConsumerFactory", boolean.class, String.class,
- ExtendedConsumerProperties.class);
+ ExtendedConsumerProperties.class, String.class);
method.setAccessible(true);
// test default for anon
- Object factory = method.invoke(binder, true, "foo-1", ecp);
+ Object factory = method.invoke(binder, true, "foo-1", ecp, "foo.consumer");
Map, ?> configs = TestUtils.getPropertyValue(factory, "configs", Map.class);
assertThat(configs.get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG))
.isEqualTo("latest");
// test default for named
- factory = method.invoke(binder, false, "foo-2", ecp);
+ factory = method.invoke(binder, false, "foo-2", ecp, "foo.consumer");
configs = TestUtils.getPropertyValue(factory, "configs", Map.class);
assertThat(configs.get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG))
.isEqualTo("earliest");
@@ -104,7 +104,7 @@ public class KafkaBinderUnitTests {
// binder level setting
binderConfigurationProperties.setConfiguration(Collections
.singletonMap(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest"));
- factory = method.invoke(binder, false, "foo-3", ecp);
+ factory = method.invoke(binder, false, "foo-3", ecp, "foo.consumer");
configs = TestUtils.getPropertyValue(factory, "configs", Map.class);
assertThat(configs.get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG))
.isEqualTo("latest");
@@ -112,7 +112,7 @@ public class KafkaBinderUnitTests {
// consumer level setting
consumerProps.setConfiguration(Collections
.singletonMap(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"));
- factory = method.invoke(binder, false, "foo-4", ecp);
+ factory = method.invoke(binder, false, "foo-4", ecp, "foo.consumer");
configs = TestUtils.getPropertyValue(factory, "configs", Map.class);
assertThat(configs.get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG))
.isEqualTo("earliest");