【问题标题】:How to handle message using Spring Cloud Stream app starter TCP如何使用 Spring Cloud Stream app starter TCP 处理消息
【发布时间】:2017-02-20 14:13:11
【问题描述】:

我想使用Spring Cloud Stream App Starter TCP Source project(maven artifact)以便能够通过套接字/端口接收 TCP 消息,处理它们,然后将结果推送到消息代理(例如 RabbitMQ )。

这个 TCP 源项目似乎完全符合我的要求,但它会自动将接收到的消息发送到输出通道。那么,是否有一种干净的方法可以仍然使用 TCP 源项目,但拦截 TCP 传入消息以在内部对其进行转换,然后再将它们输出到我的消息代理?

【问题讨论】:

    标签: java spring tcp spring-cloud-stream


    【解决方案1】:

    aggregation

    您使用源和处理器创建聚合应用程序。

    Spring Cloud Stream 支持将多个应用程序聚合在一起,直接连接它们的输入和输出通道,并避免通过代理交换消息的额外成本。从 Spring Cloud Stream 1.0 版本开始,仅以下类型的应用程序支持聚合:

    源、汇、处理器 ...

    它们可以通过创建一系列互连应用程序来聚合在一起,其中序列中一个元素的输出通道连接到下一个元素的输入通道(如果存在)。序列可以从源或处理器开始,它可以包含任意数量的处理器,并且必须以处理器或接收器结束。

    编辑

    作为源自动装配问题的解决方法,您可以尝试类似...

    @EnableBinding(Source.class)
    @EnableConfigurationProperties(TcpSourceProperties.class)
    public class MyTcpSourceConfiguration {
    
        @Autowired
        private Source channels;
    
        @Autowired
        private TcpSourceProperties properties;
    
        @Bean
        public TcpReceivingChannelAdapter adapter(
                @Qualifier("tcpSourceConnectionFactory") AbstractConnectionFactory connectionFactory) {
            TcpReceivingChannelAdapter adapter = new TcpReceivingChannelAdapter();
            adapter.setConnectionFactory(connectionFactory);
            adapter.setOutputChannelName("toMyProcessor");
            return adapter;
        }
    
        @ServiceActivator(inputChannel = "toMyProcessor", outputChannel = Source.OUTPUT)
        public byte[] myProcessor(byte[] fromTcp) {
            ...
        }
    
        @Bean
        public TcpConnectionFactoryFactoryBean tcpSourceConnectionFactory(
                @Qualifier("tcpSourceDecoder") AbstractByteArraySerializer decoder) throws Exception {
            TcpConnectionFactoryFactoryBean factoryBean = new TcpConnectionFactoryFactoryBean();
            factoryBean.setType("server");
            factoryBean.setPort(this.properties.getPort());
            factoryBean.setUsingNio(this.properties.isNio());
            factoryBean.setUsingDirectBuffers(this.properties.isUseDirectBuffers());
            factoryBean.setLookupHost(this.properties.isReverseLookup());
            factoryBean.setDeserializer(decoder);
            factoryBean.setSoTimeout(this.properties.getSocketTimeout());
            return factoryBean;
        }
    
        @Bean
        public EncoderDecoderFactoryBean tcpSourceDecoder() {
            EncoderDecoderFactoryBean factoryBean = new EncoderDecoderFactoryBean(this.properties.getDecoder());
            factoryBean.setMaxMessageSize(this.properties.getBufferSize());
            return factoryBean;
        }
    
    }
    

    【讨论】:

    • 好的,我确实成功地创建了聚合以内部接收 TCP 源产生的消息。但是,当我尝试将我的接收器更改为处理器时。我收到一条错误消息“org.springframework.cloud.stream.app.tcp.source.TcpSourceConfiguration 中的字段通道需要一个 bean,但找到了 2 个:1- 源 2- .Processor:以编程方式注册的单例”。你知道如何解决这个问题吗?
    • 是的,这是个问题;我已经打开了 issue if you want to track it 一个解决方案是实现您自己的源代码 - 复制 TcpSourceConfiguration 并直接在源代码中添加您的处理。我将详细编辑我的答案。
    • 对不起,上面的代码似乎不起作用。感谢您的帮助!
    • 糟糕 - 抱歉 - 从 @ServiceActivator 中删除 @Bean
    • 我推送了一个完全正常工作的app to my sandbox repo
    猜你喜欢
    • 2019-04-19
    • 2019-06-09
    • 2021-04-27
    • 1970-01-01
    • 1970-01-01
    • 2017-07-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多