Вам нужно добавить kstream binder в pom;стартер только добавляет связыватель канала сообщений.
РЕДАКТИРОВАТЬ
Я только что скопировал подобный код в приложение без проблем.
@SpringBootApplication
@EnableBinding(So50693858Application.PersonBinding.class)
public class So50693858Application {
public static void main(String[] args) {
SpringApplication.run(So50693858Application.class, args);
}
@StreamListener
public void process(@Input(PersonBinding.PERSON_IN) KStream<String, String> events) {
events.foreach(((key, value) -> System.out.println("Key: " + key + "; Value: " + value)));
}
interface PersonBinding {
String PERSON_IN = "pin";
String PERSON_OUT = "pout";
@Output(PERSON_OUT)
MessageChannel personOut();
@Input(PERSON_IN)
KStream<String, String> personIn();
}
}
иотправил сообщение от производителя консоли на pout
и
Key: null; Value: foo
Однако не ясно, почему у вас есть привязки ввода и вывода к одному и тому же назначению (не то, что могло бы вызвать проблему, которую вы видите).
РЕДАКТИРОВАТЬ
Это также работает (с вашими свойствами):
@SpringBootApplication
@EnableBinding(So50693858Application.PersonBinding.class)
public class So50693858Application {
public static void main(String[] args) {
SpringApplication.run(So50693858Application.class, args);
}
@Bean
public ApplicationRunner runner(MessageChannel pout) {
return args -> {
pout.send(new GenericMessage<>("foo".getBytes(),
Collections.singletonMap(KafkaHeaders.MESSAGE_KEY, "bar".getBytes())));
pout.send(new GenericMessage<>("baz".getBytes(),
Collections.singletonMap(KafkaHeaders.MESSAGE_KEY, "qux".getBytes())));
};
}
@StreamListener
public void process(@Input(PersonBinding.PERSON_IN) KStream<String, String> events) {
events.foreach(((key, value) -> System.out.println("Key: " + key + "; Value: " + value)));
}
interface PersonBinding {
String PERSON_IN = "pin";
String PERSON_OUT = "pout";
@Output(PERSON_OUT)
MessageChannel personOut();
@Input(PERSON_IN)
KStream<String, String> personIn();
}
}
и
Key: bar; Value: foo
Key: qux; Value: baz
EDIT3
Pom для версии 2.0.x:
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.example</groupId>
<artifactId>so50693858</artifactId>
<version>0.0.1-SNAPSHOT</version>
<packaging>jar</packaging>
<name>so50693858</name>
<description>Demo project for Spring Boot</description>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>2.0.2.RELEASE</version>
<relativePath/> <!-- lookup parent from repository -->
</parent>
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
<java.version>1.8</java.version>
<spring-cloud.version>Finchley.RC2</spring-cloud.version>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-stream</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-stream-binder-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-stream-test-support</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-stream-binder-kafka-streams</artifactId>
</dependency>
</dependencies>
<dependencyManagement>
<dependencies>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-dependencies</artifactId>
<version>${spring-cloud.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
<repositories>
<repository>
<id>spring-milestones</id>
<name>Spring Milestones</name>
<url>https://repo.spring.io/milestone</url>
<snapshots>
<enabled>false</enabled>
</snapshots>
</repository>
</repositories>
</project>
config:
# Default Configuration
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
# 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