I am trying build out a simple streams app based on Kafka Streams using this example.
Word Count
However when I am starting the app, I get the below error: Can someone please point out on what I am missing out here? Here is the code, config & error
@SpringBootApplication
@Slf4j
@EnableScheduling
@EnableBinding(PersonBinding.class)
public class DemoApplication {
public static void main(String[] args) {
SpringApplication.run(DemoApplication.class, args);
}
@Component
public static class PersonSource {
private final MessageChannel personOut;
@Autowired
PersonSource(PersonBinding personBinding) {
this.personOut = personBinding.personOut();
}
@Scheduled(fixedDelay = 5000L)
public void run() {
Message<Person> message = MessageBuilder
.withPayload(new Person("John", "Doe", Instant.now()))
.build();
try {
personOut.send(message);
log.info("Published message: {}", message);
} catch (Exception e) {
e.printStackTrace();
throw e;
}
}
}
@Component
public static class PersonProcessor {
@StreamListener
public void process(@Input(PersonBinding.PERSON_IN) KStream<String, Person> events) {
events.foreach(((key, value) -> System.out.println("Key: " + key + "; Value: " + value)));
}
}
}
@Data
@AllArgsConstructor
class Person {
String firstName;
String lastName;
Instant createdOn;
}
interface PersonBinding {
String PERSON_IN = "pin";
String PERSON_OUT = "pout";
@Output(PERSON_OUT)
MessageChannel personOut();
@Input(PERSON_IN)
KStream<String, Person> personIn();
}
Dependency Management (Spring Boot 1.5.13.RELEASE)
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-stream-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-stream-binder-kstream</artifactId>
</dependency>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams</artifactId>
<version>1.1.0</version>
</dependency>
Configuration
# Default Configuration
spring.cloud.stream.kstream.binder.configuration.default.key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde
spring.cloud.stream.kstream.binder.configuration.default.value.serde=org.apache.kafka.common.serialization.Serdes$StringSerde
# Out Bindings Configuration
spring.cloud.stream.bindings.pout.destination=pout
spring.cloud.stream.bindings.pout.producer.header-mode=raw
# In Bindings Configuration
spring.cloud.stream.bindings.pin.destination=pout
spring.cloud.stream.bindings.pin.consumer.header-mode=raw
Error
Field configurationProperties in org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfiguration required a single bean, but 2 were found:
- spring.cloud.stream.kafka.binder-org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties: a programmatically registered singleton - binderConfigurationProperties: defined by method 'binderConfigurationProperties' in class path resource [org/springframework/cloud/stream/binder/kstream/config/KStreamBinderSupportAutoConfiguration.class]
Action:
Consider marking one of the beans as @Primary, updating the consumer to accept multiple beans, or using @Qualifier to identify the bean that should be consumed
** EDIT 1 **
Uploaded code to Github https://github.com/tapitoe/demo-spring-cloud-streams/tree/master/src
I had a similar problem, It got fixed after I added:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
If you love us? You can donate to us via Paypal or buy me a coffee so we can maintain and grow! Thank you!
Donate Us With