草庐IT

java - spring 集成中的自定义并发 TcpOutboundGateway

coder 2023-09-19 原文

我正在尝试使用 TcpOutboundGateway 和客户端 TcpConnectionFactory 与外部 TCP 服务器通信。在我的场景中,每个连接都应该与不同的线程相关联(线程上的每个连接可能用于多个请求/响应)。

所以我使用了这个主题中的 ThreadAffinityClientConnectionFactory:Spring Integration tcp client multiple connections

它工作正常,直到我尝试打开超过 4 个并发连接,第五个(及以上)连接因超时而失败。 我发现 org.springframework.integration.ip.tcp.TcpOutboundGatewayhandleRequestMessage 方法中使用信号量来获取连接,所以我覆盖了 TcpOuboundGateway像这样:

public class NoSemaphoreTcpOutboundGateway extends TcpOutboundGateway {

    private volatile AbstractClientConnectionFactory connectionFactory;
    private final Map<String, NoSemaphoreTcpOutboundGateway.AsyncReply> pendingReplies = new ConcurrentHashMap();

    @Override
    public boolean onMessage(Message<?> message) {
        String connectionId = (String)message.getHeaders().get("ip_connectionId");
        if(connectionId == null) {
            this.logger.error("Cannot correlate response - no connection id");
            this.publishNoConnectionEvent(message, (String)null, "Cannot correlate response - no connection id");
            return false;
        }

        if(this.logger.isTraceEnabled()) {
            this.logger.trace("onMessage: " + connectionId + "(" + message + ")");
        }

        NoSemaphoreTcpOutboundGateway.AsyncReply reply = (NoSemaphoreTcpOutboundGateway.AsyncReply)this.pendingReplies.get(connectionId);
        if(reply == null) {
            if(message instanceof ErrorMessage) {
                return false;
            } else {
                String errorMessage = "Cannot correlate response - no pending reply for " + connectionId;
                this.logger.error(errorMessage);
                this.publishNoConnectionEvent(message, connectionId, errorMessage);
                return false;
            }
        } else {
            reply.setReply(message);
            return false;
        }

    }

    @Override
    protected Message handleRequestMessage(Message<?> requestMessage) {
        connectionFactory = (AbstractClientConnectionFactory) this.getConnectionFactory();
        Assert.notNull(this.getConnectionFactory(), this.getClass().getName() + " requires a client connection factory");

        TcpConnection connection = null;
        String connectionId = null;

        Message var7;
        try {
            /*if(!this.isSingleUse()) {
                this.logger.debug("trying semaphore");
                if(!this.semaphore.tryAcquire(this.requestTimeout, TimeUnit.MILLISECONDS)) {
                    throw new MessageTimeoutException(requestMessage, "Timed out waiting for connection");
                }

                haveSemaphore = true;
                if(this.logger.isDebugEnabled()) {
                    this.logger.debug("got semaphore");
                }
            }*/

            connection = this.getConnectionFactory().getConnection();
            NoSemaphoreTcpOutboundGateway.AsyncReply e = new NoSemaphoreTcpOutboundGateway.AsyncReply(10000);
            connectionId = connection.getConnectionId();
            this.pendingReplies.put(connectionId, e);
            if(this.logger.isDebugEnabled()) {
                this.logger.debug("Added pending reply " + connectionId);
            }

            connection.send(requestMessage);

            //connection may be closed after send (in interceptor) if its disconnect message
            if (!connection.isOpen())
                return null;

            Message replyMessage = e.getReply();
            if(replyMessage == null) {
                if(this.logger.isDebugEnabled()) {
                    this.logger.debug("Remote Timeout on " + connectionId);
                }

                this.connectionFactory.forceClose(connection);
                throw new MessageTimeoutException(requestMessage, "Timed out waiting for response");
            }

            if(this.logger.isDebugEnabled()) {
                this.logger.debug("Response " + replyMessage);
            }

            var7 = replyMessage;
        } catch (Exception var11) {
            this.logger.error("Tcp Gateway exception", var11);
            if(var11 instanceof MessagingException) {
                throw (MessagingException)var11;
            }

            throw new MessagingException("Failed to send or receive", var11);
        } finally {
            if(connectionId != null) {
                this.pendingReplies.remove(connectionId);
                if(this.logger.isDebugEnabled()) {
                    this.logger.debug("Removed pending reply " + connectionId);
                }
            }
        }
        return var7;
    }

    private void publishNoConnectionEvent(Message<?> message, String connectionId, String errorMessage) {
        ApplicationEventPublisher applicationEventPublisher = this.connectionFactory.getApplicationEventPublisher();
        if(applicationEventPublisher != null) {
            applicationEventPublisher.publishEvent(new TcpConnectionFailedCorrelationEvent(this, connectionId, new MessagingException(message, errorMessage)));
        }
    }

    private final class AsyncReply {
        private final CountDownLatch latch;
        private final CountDownLatch secondChanceLatch;
        private final long remoteTimeout;
        private volatile Message<?> reply;

        private AsyncReply(long remoteTimeout) {
            this.latch = new CountDownLatch(1);
            this.secondChanceLatch = new CountDownLatch(1);
            this.remoteTimeout = remoteTimeout;
        }

        public Message<?> getReply() throws Exception {
            try {
                if(!this.latch.await(this.remoteTimeout, TimeUnit.MILLISECONDS)) {
                    return null;
                }
            } catch (InterruptedException var2) {
                Thread.currentThread().interrupt();
            }
            for(boolean waitForMessageAfterError = true; this.reply instanceof ErrorMessage; waitForMessageAfterError = false) {
                if(!waitForMessageAfterError) {
                    if(this.reply.getPayload() instanceof MessagingException) {
                        throw (MessagingException)this.reply.getPayload();
                    }

                    throw new MessagingException("Exception while awaiting reply", (Throwable)this.reply.getPayload());
                }
                NoSemaphoreTcpOutboundGateway.this.logger.debug("second chance");
                this.secondChanceLatch.await(2L, TimeUnit.SECONDS);
            }
            return this.reply;
        }

        public void setReply(Message<?> reply) {
            if(this.reply == null) {
                this.reply = reply;
                this.latch.countDown();
            } else if(this.reply instanceof ErrorMessage) {
                this.reply = reply;
                this.secondChanceLatch.countDown();
            }
        }
    }
}

SpringContext的配置是这样的:

@Configuration
@ImportResource("gateway.xml")
public class Conf {

@Bean
@Autowired
@ServiceActivator(inputChannel = "clientOutChannel")
public NoSemaphoreTcpOutboundGateway noSemaphoreTcpOutboundGateway(ThreadAffinityClientConnectionFactory cf, DirectChannel clientInChannel){
    NoSemaphoreTcpOutboundGateway gw = new NoSemaphoreTcpOutboundGateway();
    gw.setConnectionFactory(cf);
    gw.setReplyChannel(clientInChannel);
    gw.setRequestTimeout(10000);
    return gw;
}

 <int-ip:tcp-connection-factory
        id="delegateCF"
        type="client"
        host="${remoteService.host}"
        port="${remoteService.port}"
        single-use="true"
        lookup-host="false"
        ssl-context-support="sslContext"
        deserializer="clientDeserializer"
        serializer="clientSerializer"
        interceptor-factory-chain="clientLoggingTcpConnectionInterceptorFactory"
        using-nio="false"/>

delegateCF 被传递给ThreadAffinityClientConnectionFactory 构造函数

那么,问题是:

  • 就并发性而言,可以将 NoSemaphoreTcpOutboundGatewayThreadAffinityClientConnectionFactory 结合使用吗?

最佳答案

看起来你走对了,但同时我认为你不需要自定义 TcpOutboundGateway信号量 逻辑基于:

if (!this.isSingleUse) {
            logger.debug("trying semaphore");
            if (!this.semaphore.tryAcquire(this.requestTimeout, TimeUnit.MILLISECONDS)) {
                throw new MessageTimeoutException(requestMessage, "Timed out waiting for connection");
            }

同时看看Gary对ThreadAffinityClientConnectionFactory的解决方案:

@Bean
public TcpNetClientConnectionFactory delegateCF() {
    TcpNetClientConnectionFactory clientCF = new TcpNetClientConnectionFactory("localhost", 1234);
    clientCF.setSingleUse(true); // so each thread gets his own connection
    return clientCF;
}

@Bean
public ThreadAffinityClientConnectionFactory affinityCF() {
    return new ThreadAffinityClientConnectionFactory(delegateCF());
}

关注评论。您只需要委托(delegate) isSingleUse()

关于java - spring 集成中的自定义并发 TcpOutboundGateway,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/41166718/

有关java - spring 集成中的自定义并发 TcpOutboundGateway的更多相关文章

  1. ruby - 如何从 ruby​​ 中的字符串运行任意对象方法? - 2

    总的来说,我对ruby​​还比较陌生,我正在为我正在创建的对象编写一些rspec测试用例。许多测试用例都非常基础,我只是想确保正确填充和返回值。我想知道是否有办法使用循环结构来执行此操作。不必为我要测试的每个方法都设置一个assertEquals。例如:describeitem,"TestingtheItem"doit"willhaveanullvaluetostart"doitem=Item.new#HereIcoulddotheitem.name.shouldbe_nil#thenIcoulddoitem.category.shouldbe_nilendend但我想要一些方法来使用

  2. ruby - 其他文件中的 Rake 任务 - 2

    我试图在一个项目中使用rake,如果我把所有东西都放到Rakefile中,它会很大并且很难读取/找到东西,所以我试着将每个命名空间放在lib/rake中它自己的文件中,我添加了这个到我的rake文件的顶部:Dir['#{File.dirname(__FILE__)}/lib/rake/*.rake'].map{|f|requiref}它加载文件没问题,但没有任务。我现在只有一个.rake文件作为测试,名为“servers.rake”,它看起来像这样:namespace:serverdotask:testdoputs"test"endend所以当我运行rakeserver:testid时

  3. ruby-on-rails - Ruby net/ldap 模块中的内存泄漏 - 2

    作为我的Rails应用程序的一部分,我编写了一个小导入程序,它从我们的LDAP系统中吸取数据并将其塞入一个用户表中。不幸的是,与LDAP相关的代码在遍历我们的32K用户时泄漏了大量内存,我一直无法弄清楚如何解决这个问题。这个问题似乎在某种程度上与LDAP库有关,因为当我删除对LDAP内容的调用时,内存使用情况会很好地稳定下来。此外,不断增加的对象是Net::BER::BerIdentifiedString和Net::BER::BerIdentifiedArray,它们都是LDAP库的一部分。当我运行导入时,内存使用量最终达到超过1GB的峰值。如果问题存在,我需要找到一些方法来更正我的代

  4. ruby - Facter::Util::Uptime:Module 的未定义方法 get_uptime (NoMethodError) - 2

    我正在尝试设置一个puppet节点,但ruby​​gems似乎不正常。如果我通过它自己的二进制文件(/usr/lib/ruby/gems/1.8/gems/facter-1.5.8/bin/facter)在cli上运行facter,它工作正常,但如果我通过由ruby​​gems(/usr/bin/facter)安装的二进制文件,它抛出:/usr/lib/ruby/1.8/facter/uptime.rb:11:undefinedmethod`get_uptime'forFacter::Util::Uptime:Module(NoMethodError)from/usr/lib/ruby

  5. ruby-on-rails - Rails 3 中的多个路由文件 - 2

    Rails2.3可以选择随时使用RouteSet#add_configuration_file添加更多路由。是否可以在Rails3项目中做同样的事情? 最佳答案 在config/application.rb中:config.paths.config.routes在Rails3.2(也可能是Rails3.1)中,使用:config.paths["config/routes"] 关于ruby-on-rails-Rails3中的多个路由文件,我们在StackOverflow上找到一个类似的问题

  6. ruby-on-rails - Rails - 一个 View 中的多个模型 - 2

    我需要从一个View访问多个模型。以前,我的links_controller仅用于提供以不同方式排序的链接资源。现在我想包括一个部分(我假设)显示按分数排序的顶级用户(@users=User.all.sort_by(&:score))我知道我可以将此代码插入每个链接操作并从View访问它,但这似乎不是“ruby方式”,我将需要在不久的将来访问更多模型。这可能会变得很脏,是否有针对这种情况的任何技术?注意事项:我认为我的应用程序正朝着单一格式和动态页面内容的方向发展,本质上是一个典型的网络应用程序。我知道before_filter但考虑到我希望应用程序进入的方向,这似乎很麻烦。最终从任何

  7. ruby-on-rails - Rails 3.2.1 中 ActionMailer 中的未定义方法 'default_content_type=' - 2

    我在我的项目中添加了一个系统来重置用户密码并通过电子邮件将密码发送给他,以防他忘记密码。昨天它运行良好(当我实现它时)。当我今天尝试启动服务器时,出现以下错误。=>BootingWEBrick=>Rails3.2.1applicationstartingindevelopmentonhttp://0.0.0.0:3000=>Callwith-dtodetach=>Ctrl-CtoshutdownserverExiting/Users/vinayshenoy/.rvm/gems/ruby-1.9.3-p0/gems/actionmailer-3.2.1/lib/action_mailer

  8. ruby-on-rails - Rails 应用程序中的 Rails : How are you using application_controller. rb 是新手吗? - 2

    刚入门rails,开始慢慢理解。有人可以解释或给我一些关于在application_controller中编码的好处或时间和原因的想法吗?有哪些用例。您如何为Rails应用程序使用应用程序Controller?我不想在那里放太多代码,因为据我了解,每个请求都会调用此Controller。这是真的? 最佳答案 ApplicationController实际上是您应用程序中的每个其他Controller都将从中继承的类(尽管这不是强制性的)。我同意不要用太多代码弄乱它并保持干净整洁的态度,尽管在某些情况下ApplicationContr

  9. ruby-on-rails - form_for 中不在模型中的自定义字段 - 2

    我想向我的Controller传递一个参数,它是一个简单的复选框,但我不知道如何在模型的form_for中引入它,这是我的观点:{:id=>'go_finance'}do|f|%>Transferirde:para:Entrada:"input",:placeholder=>"Quantofoiganho?"%>Saída:"output",:placeholder=>"Quantofoigasto?"%>Nota:我想做一个额外的复选框,但我该怎么做,模型中没有一个对象,而是一个要检查的对象,以便在Controller中创建一个ifelse,如果没有检查,请帮助我,非常感谢,谢谢

  10. ruby - 主要 :Object when running build from sublime 的未定义方法 `require_relative' - 2

    我已经从我的命令行中获得了一切,所以我可以运行rubymyfile并且它可以正常工作。但是当我尝试从sublime中运行它时,我得到了undefinedmethod`require_relative'formain:Object有人知道我的sublime设置中缺少什么吗?我正在使用OSX并安装了rvm。 最佳答案 或者,您可以只使用“require”,它应该可以正常工作。我认为“require_relative”仅适用于ruby​​1.9+ 关于ruby-主要:Objectwhenrun

随机推荐