【发布时间】:2019-02-05 03:52:54
【问题描述】:
我最近开始研究适用于 Kafka 的 Spring Cloud Stream,并且一直在努力使 TestBinder 与 Kstreams 一起工作。这是一个已知的限制,还是我只是忽略了一些东西?
这很好用:
字符串处理器:
@StreamListener(TopicBinding.INPUT)
@SendTo(TopicBinding.OUTPUT)
public String process(String message) {
return message + " world";
}
字符串测试:
@Test
@SuppressWarnings("unchecked")
public void testString() {
Message<String> message = new GenericMessage<>("Hello");
topicBinding.input().send(message);
Message<String> received = (Message<String>) messageCollector.forChannel(topicBinding.output()).poll();
assertThat(received.getPayload(), equalTo("Hello world"));
}
但是当我尝试在我的流程中使用 KStream 时,我无法让 TestBinder 工作。
K流处理器:
@SendTo(TopicBinding.OUTPUT)
public KStream<String, String> process(
@Input(TopicBinding.INPUT) KStream<String, String> events) {
return events.mapValues((value) -> value + " world");
}
KStream 测试:
@Test
@SuppressWarnings("unchecked")
public void testKstream() {
Message<String> message = MessageBuilder
.withPayload("Hello")
.setHeader(KafkaHeaders.TOPIC, "event.sirism.dev".getBytes())
.setHeader(KafkaHeaders.MESSAGE_KEY, "Test".getBytes())
.build();
topicBinding.input().send(message);
Message<String> received = (Message<String>)
messageCollector.forChannel(topicBinding.output()).poll();
assertThat(received.getPayload(), equalTo("Hello world"));
}
您可能已经注意到,我从 Kstream 处理器中省略了 @StreamListener,但没有它,testbinder 似乎无法找到处理程序。 (但有了它,启动应用时就不行了)
这是一个已知的错误/限制,还是我只是在这里做一些愚蠢的事情?
【问题讨论】:
标签: spring-boot apache-kafka spring-cloud-stream