【发布时间】:2019-01-18 23:00:11
【问题描述】:
我正在使用 Spring Cloud Stream 版本 2.1.0.RELEASE 将消息(在本例中为 Kafka)发送到基于接收到的输入动态定义的通道。问题是只有所有其他消息最终都在正确的频道中,另一半最终在默认频道中。
我使用this 示例作为起点。
我将要发送到的通道放入特定的消息标头,然后使用 HeaderValueRouter 检查相同的标头值以查看要输出到哪个通道。
我正在按如下方式配置我的应用程序:
@EnableBinding(CloudStreamConfig.DynamicSource.class)
public class CloudStreamConfig {
@Autowired
private BinderAwareChannelResolver resolver;
public static final String CHANNEL_HEADER = "channelHeader";
public static final String OUTPUT_CHANNEL = "outputChannel";
private final String defaultChannel = "defaultChannel";
@ServiceActivator(inputChannel = OUTPUT_CHANNEL)
@Bean
public HeaderValueRouter router() {
HeaderValueRouter router = new HeaderValueRouter(CHANNEL_HEADER);
router.setDefaultOutputChannelName(defaultChannel);
router.setChannelResolver(resolver);
return router;
}
public interface DynamicSource {
@Output(OUTPUT_CHANNEL)
MessageChannel output();
}
}
在我的控制器中,我接收一个对象以及一个定义将其发送到哪个通道的参数,然后将其发送到 MessageChannel。相关代码如下:
@Autowired
@Qualifier(CloudStreamConfig.OUTPUT_CHANNEL)
public MessageChannel localChannel;
...
@GetMapping(path = "/error/{channel}")
@ResponseStatus(HttpStatus.OK)
public void error(@PathVariable String channel) {
// build my object
Message message = MessageBuilder.createMessage(myObject,
new MessageHeaders(Collections.singletonMap(CloudStreamConfig.CHANNEL_HEADER, channel)));
localChannel.send(message);
}
如果我向 /error/someChannel 发送 10 条消息,我希望在 someChannel 中看到 10 条消息。但是,我在someChannel 中看到一半消息,在defaultChannel 中看到另一半消息。我在消息中放置了一个调试计数器变量,它将第一条消息发送到正确的通道,然后每隔一条消息发送到正确的通道,而其他消息都发送到默认通道。
这是什么原因造成的,我该如何解决?我在滥用我的DynamicSource 课程吗?我以为它会与任何同名的自动连线MessageChannel 相关联(而且看起来确实如此),但我想知道我是否缺少某些东西。还是与BinderAwareChannelResolver 发生了意外交互? (老实说,我不知道这是做什么的,我只是因为样本做了才包含它)
【问题讨论】:
标签: java spring spring-integration spring-cloud-stream