【发布时间】:2017-11-25 21:04:08
【问题描述】:
我正在尝试创建一个基于事件的系统,用于使用 Apache Kafka 作为消息传递系统和 Spring Cloud Stream Kafka 的服务之间的通信。
我已经编写了如下的接收器类方法,
@StreamListener(target = Sink.INPUT, condition = "headers['eventType']=='EmployeeCreatedEvent'")
public void handleEmployeeCreatedEvent(@Payload String payload) {
logger.info("Received EmployeeCreatedEvent: " + payload);
}
此方法专门用于捕获与 EmployeeCreatedEvent 相关的消息或事件。
@StreamListener(target = Sink.INPUT, condition = "headers['eventType']=='EmployeeTransferredEvent'")
public void handleEmployeeTransferredEvent(@Payload String payload) {
logger.info("Received EmployeeTransferredEvent: " + payload);
}
此方法专门用于捕获与 EmployeeTransferredEvent 相关的消息或事件。
@StreamListener(target = Sink.INPUT)
public void handleDefaultEvent(@Payload String payload) {
logger.info("Received payload: " + payload);
}
这是默认方法。
当我运行应用程序时,我看不到使用条件属性注释的方法被调用。我只看到调用了 handleDefaultEvent 方法。
我正在使用下面的 CustomMessageSource 类从发送/源应用程序向此接收器应用程序发送消息,如下所示,
@Component
@EnableBinding(Source.class)
public class CustomMessageSource {
@Autowired
private Source source;
public void sendMessage(String payload,String eventType) {
Message<String> myMessage = MessageBuilder.withPayload(payload)
.setHeader("eventType", eventType)
.build();
source.output().send(myMessage);
}
}
我正在源应用程序中从我的控制器调用该方法,如下所示,
customMessageSource.sendMessage("Hello","EmployeeCreatedEvent");
customMessageSource 实例如下自动装配,
@Autowired
CustomMessageSource customMessageSource;
基本上,我想过滤接收器/接收器应用程序收到的消息并相应地处理它们。
为此,我使用@StreamListener 注解和条件属性来模拟处理不同事件的行为。
我正在使用 Spring Cloud Stream Chelsea.SR2 版本。
谁能帮我解决这个问题。
【问题讨论】:
标签: events apache-kafka spring-integration spring-cloud-stream