【问题标题】:Spring cloud stream: dynamic output channel strange behaviorSpring Cloud Stream:动态输出通道的奇怪行为
【发布时间】: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


    【解决方案1】:

    输出通道上有两个订阅者 - 通道绑定(在活页夹中)和您的路由器。

    对于DirectChannels,默认的调度算法是轮询,因此您将消息交替发送到路由器和直接发送到活页夹。

    您需要一个不同的DirectChannel @Bean 用于服务激活器,以便所有消息都到达那里,然后在路由后到达活页夹。

    请参阅该示例中的sourceChannel

    【讨论】:

    • 我对你的第一句话感到困惑,如果我只创建路由器,为什么输出频道上还有另一个订阅者?路由器是不是在输出通道之前的一个层次,决定输出到哪个通道?
    • 另外,当我直接按照示例进行操作并有一个名为“sourceChannel”的 DirectChannel 和一个名为“sourceChannel”的 ServiceActivator 时,我得到了循环依赖,这就是我尝试上述尝试的原因跨度>
    • 这很简单——你的接口的@EnableBinding 创建一个消费者,@ServiceActivator 创建第二个消费者。您是否尝试过实际运行示例?如果这得到循环依赖,请在该 repo 中打开一个问题。如果您仍然无法弄清楚有什么不同,请将您的项目发布到某个地方。
    • 我的错误,按原样运行项目并且不会创建循环依赖项。但是,一旦我尝试将通道和路由器配置拉出到一个单独的类中,然后在其他地方注入MessageChannel,我就会得到一个循环依赖。我只是能够通过在内部类中配置MessageChannel 来解决它
    猜你喜欢
    • 2017-05-31
    • 2020-07-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-01-13
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多