【问题标题】:Application level load balancer for pool of processes进程池的应用程序级负载均衡器
【发布时间】:2014-10-05 06:16:00
【问题描述】:

我们有传统的 C++ 单体软件,其作用类似于请求-回复 TCP 服务器。该软件是单线程的,可以同时处理一个请求。目前,我们有固定的此类进程池以并行服务多个客户端。

由于大量消息,客户端会定期遇到请求处理的严重延迟。目前我们有一个想法通过在客户端和工作人员之间引入一种代理来解决这个问题:

我们希望此代理具有以下功能:

  1. 应用程序级负载平衡:通过检查请求上下文和客户端 ID 在工作人员之间传播请求
  2. 控制和监控工作进程的生命周期
  3. 产生额外的工作进程(在不同的 PC 上)来处理峰值

事实上,我们希望它的行为类似于 Java 中的 ExecutorService,但使用 工作进程而不是线程。目前的想法是基于 Jetty 或 Tomcat 服务器在 Java 中实现此 Balancer,内部消息队列和 servlet 将请求转发到工作进程。

但我想知道:是否存在可以自动化此过程的现有解决方案(最好在 Java 中使用)?实现这样的代理最简单的方法是什么?

更新:

我对请求上下文所做的事情 - 嗯,C++ 服务器确实是一个混乱的软件。事实上,每次它接收到不同的上下文时,它都会相应地更新内部缓存以匹配该上下文。例如,如果您请求该服务器向您提供一些英文数据,那么它会将内部缓存重新加载为英文。如果下一个请求是法语的,那么它会再次重新加载缓存。显然我想通过更智能地转发请求来最小化缓存重新加载的次数。

通信协议是自制的(基于TCP/IP),但从中提取上下文部分相对容易。

目前负载平衡是在客户端实现的,因此每个客户端都被配置为了解所有服务器节点,并以循环方式向它们发送请求。这种方法存在几个问题:客户端连接管理复杂,与多个互不了解的客户端一起工作不正确,无法管理节点生命周期。我们无法通过重构解决列出的问题。

很可能我们最终会使用自制的转发解决方案,但我仍然想知道是否有现有产品至少用于流程管理?理想情况下,这将是 Java 应用程序服务器,它可以:

  • 派生子节点(另一个 Java 进程)
  • 监控其生命周期
  • 通过某种协议与他们通信

也许这个功能已经在一些现有的应用服务器中实现了?这将大大简化问题!

【问题讨论】:

  • 你为什么不使用像 haproxy 这样的东西作为负载均衡器。它支持 TCP 端口之间的负载平衡,还能够动态重新加载配置,帮助您添加和删除工作人员。
  • @abyz:谢谢你的建议,听起来很有趣!但是我们需要根据请求内容(例如语言、客户端等)将请求转发给工作人员——我不确定 HAProxy 是否支持。我们需要的是基于可配置规则的“智能”转发。
  • 如果是这种情况恕我直言,您必须根据需要从头开始开发它,并使用 ActiveMQ、ZeroMQ、HornetQ 等某种排队系统来正确处理并发、负载平衡和路由。我认为还需要稍微改变你的工作代码。

标签: java c++ networking architecture


【解决方案1】:

不确定您要如何检查请求上下文和客户端 ID,以及这会如何影响路由。还有什么是请求内容格式?如果不是文本而是二进制,那么事情会变得更加困难。

您希望对流程进行什么样的控制和监控?

虽然您将此与 Java Executor 服务进行比较,但将问题空间从线程扩展到进程是可以的,但当您谈论多台 PC 时,它会呈指数级增长。现在我们在这里讨论集群。

您需要为此专有的 TCP/IP 请求/响应自定义负载平衡。现在,如果我没有过度设计,您可能还想开始考虑节点故障的策略并添加更多节点以进一步扩展您的系统。

当前处理请求映射到池化进程的组件是什么?是否可以将此组件重构为可配置,以便将请求路由到池化进程或一组已配置节点中的一个节点?

【讨论】:

  • 我更新了我的帖子来回答你的问题。是的,我们很可能最终会使用自定义转发,这没关系。我现在的问题主要是关于简化流程管理(假设我们要从 Java 服务器管理一组 Java 节点)。是否可以为此使用现有的 Java 应用程序服务器功能?
【解决方案2】:

您似乎正在寻找消息传递系统。 Apache Camel 有很多组件可以集成不同的协议并添加自定义处理逻辑(使用 XML 或 java API)。 Apache Camel 已经实现了很多 (Enterprise Integration Patterns)。

它与Apache MINA 集成,这也是一个很好的起点。

目前尚不清楚如何在其他计算机上即时启动新实例。我认为您至少需要在这些机器上运行一些代理,您可以请求启动新服务器。

【讨论】:

  • 感谢您的回答。实际上,正如我在问题末尾所描述的那样,我主要是在寻找简化流程管理的框架。
【解决方案3】:

由于您坚持使用 TCP/IP 协议并且想要基于内容的路由 - 您必须自己做一些事情。

我认为您可以获得一些现成的集成平台,并为您的协议和路由处理程序编写适配器。

根据我的经验:Apache Camel - 是可行的方法。

作为设计的第一步,我将采取:

  1. Camel 作为集成总线
  2. Apache Active MQ 作为消息代理 (JMS)
  3. Mysql 或 Postgre SQL 或 HSQLDB 作为数据库(基于负载和大小要求)
  4. Server 的 TCP/IP 协议设计适配器从客户端获取连接并将数据发送到 JMS
  5. 设计路由组件,它从 JMS 获取消息并将其转发到所需的队列,并根据以下条件将其转发到服务器之一:服务器上下文、负载、队列长度。
  6. 设计端点组件,从 JMS 获取消息并发送到真实服务器,获取响应并将其发送到 JMS。

Camel 服务器的作用是构建流程中所有步骤的工作流: 它可以从JMS获取消息并调用java方法,从调用中获取返回数据并推送到JMS等。所以你不需要自己做。

JMS 的作用 - 在每种类型的多个节点之间平衡工作负载:适配器、路由器和端点。

悬而未决的问题仍然存在:

  1. 监控和自动启动服务器上的节点池 - 这个问题值得单独讨论。

【讨论】:

  • 我已经编辑了我的问题,强调我实际上是在寻找工具/库/框架来简化流程管理。对于交流部分-我一点问题都没有。主要问题是您在“未解决的问题仍然存在”中输入的内容
【解决方案4】:

关于流程管理,您可以通过混合Apache Commons Exec 库的功能轻松实现您的目标,Apache Commons Exec 库可以帮助生成新的工人实例,Apache Commons Pool 库将管理正在运行的实例。

实现非常简单,因为公共池将确保您一次可以使用一个对象,直到它返回到池中。如果对象没有返回到池中,公共池将为您生成新实例。您可以通过添加看门狗服务(来自 apache commons exec)来控制工作人员的生命周期 - 看门狗可以杀死一段时间未使用的实例,或者您也可以使用公共池本身,例如通过调用 pool.clearOldest()。您还可以通过调用 pool.getNumActive() 查看当前处理了多少请求(有多少工作人员处于活动状态)。更多内容请参考 GenericKeyedObjectPool 的 javadoc。

可以通过一个在 Tomcat 服务器上运行的简单 servlet 来实现。这个 servlet 将实例化池并通过调用 pool.borowObject(parameters) 简单地向池请求新的工作人员。在参数内部,您定义了应让您的工作人员处理请求的特征(在您的情况下,参数应包括语言)。如果没有这样的工人可用(例如没有法语工人),池将为您生成新的工人。此外,如果有一个工作人员但该工作人员当前正在处理另一个请求,则池还将为您生成一个新工作人员(因此您将有两个工作人员处理相同的语言)。当您调用 pool.returnObject(parameters, instance) 时,Worker 将准备好处理新请求。

整个实现只用了不到 200 行代码(完整代码见下文)。该代码包括工作进程被外部杀死或崩溃的情况(请参阅 WorkersFactory.activateObject())。

恕我直言:使用 Apache Cammel 对您来说不是一个好选择,因为它太大了,而且它被设计为不同消息格式之间的中介总线。您不需要进行转换,也不需要处理不同格式的消息。寻求简单的解决方案。

package com.myapp;

import org.apache.commons.exec.CommandLine;
import org.apache.commons.exec.DefaultExecutor;
import org.apache.commons.exec.ExecuteWatchdog;
import org.apache.commons.pool2.BaseKeyedPooledObjectFactory;
import org.apache.commons.pool2.KeyedPooledObjectFactory;
import org.apache.commons.pool2.PooledObject;
import org.apache.commons.pool2.impl.DefaultPooledObject;
import org.apache.commons.pool2.impl.GenericKeyedObjectPool;

import javax.servlet.ServletException;
import javax.servlet.http.HttpServletRequest;
import javax.servlet.http.HttpServletResponse;
import java.io.IOException;
import java.util.Objects;

public class BalancingServlet extends javax.servlet.http.HttpServlet {

    private final WorkersPool workersPool;

    public BalancingServlet() {
        workersPool = new WorkersPool(new WorkersFactory());
    }


    protected void doPost(HttpServletRequest request, HttpServletResponse response) throws ServletException, IOException {

    }

    protected void doGet(HttpServletRequest request, HttpServletResponse response) throws ServletException, IOException {
        response.getWriter().println("Balancing");

        String language = request.getParameter("language");
        String someOtherParam = request.getParameter("other");
        WorkerParameters workerParameters = new WorkerParameters(language, someOtherParam);

        String requestSpecificParam1 = request.getParameter("requestParam1");
        String requestSpecificParam2 = request.getParameter("requestParam2");

        try {
            WorkerInstance workerInstance = workersPool.borrowObject(workerParameters);
            workerInstance.handleRequest(requestSpecificParam1, requestSpecificParam2);
            workersPool.returnObject(workerParameters, workerInstance);

        } catch (Exception e) {
            e.printStackTrace();
        }


    }
}

class WorkerParameters {
    private final String workerLangauge;
    private final String someOtherParam;

    WorkerParameters(String workerLangauge, String someOtherParam) {
        this.workerLangauge = workerLangauge;
        this.someOtherParam = someOtherParam;
    }

    public String getWorkerLangauge() {
        return workerLangauge;
    }

    public String getSomeOtherParam() {
        return someOtherParam;
    }

    @Override
    public boolean equals(Object o) {
        if (this == o) return true;
        if (o == null || getClass() != o.getClass()) return false;

        WorkerParameters that = (WorkerParameters) o;

        return Objects.equals(this.workerLangauge, that.workerLangauge) && Objects.equals(this.someOtherParam, that.someOtherParam);
    }

    @Override
    public int hashCode() {
        return Objects.hash(workerLangauge, someOtherParam);
    }
}

class WorkerInstance {
    private final Thread thread;
    private WorkerParameters workerParameters;

    public WorkerInstance(final WorkerParameters workerParameters) {
        this.workerParameters = workerParameters;

        // launch the process here   
        System.out.println("Spawing worker for language: " + workerParameters.getWorkerLangauge());

        // use commons Exec to spawn your process using command line here

        // something like


        thread = new Thread(new Runnable() {
            @Override
            public void run() {
                try {
                    String line = "C:/Windows/notepad.exe" ;
                    final CommandLine cmdLine = CommandLine.parse(line);

                    final DefaultExecutor executor = new DefaultExecutor();
                    executor.setExitValue(0);
//                    ExecuteWatchdog watchdog = new ExecuteWatchdog(60000); // if you want to kill process running too long
//                    executor.setWatchdog(watchdog);

                    int exitValue = executor.execute(cmdLine);
                    System.out.println("process finished with exit code: " + exitValue);
                } catch (IOException e) {
                    throw new RuntimeException("Problem while executing application for language: " + workerParameters.getWorkerLangauge(), e);
                }


            }
        });

        thread.start();


        System.out.println("Process spawned for language: " + workerParameters.getWorkerLangauge());


    }

    public void handleRequest(String someRequestParam1, String someRequestParam2) {
        System.out.println("Handling request for extra params: " + someRequestParam1 + ", " + someRequestParam2);

        // communicate with your application using parameters here

        // communcate via tcp or whatever protovol you want using extra parameters: someRequestParam1, someRequestParam2


    }

    public boolean isRunning() {
        return thread.isAlive();
    }


}

class WorkersFactory extends BaseKeyedPooledObjectFactory<WorkerParameters, WorkerInstance> {

    @Override
    public WorkerInstance create(WorkerParameters parameters) throws Exception {
        return new WorkerInstance(parameters);
    }

    @Override
    public PooledObject<WorkerInstance> wrap(WorkerInstance worker) {
        return new DefaultPooledObject<WorkerInstance>(worker);
    }

    @Override
    public void activateObject(WorkerParameters worker, PooledObject<WorkerInstance> p)
            throws Exception {
        System.out.println("Activating worker for lang: " + worker.getWorkerLangauge());

        if  (! p.getObject().isRunning()) {
            System.out.println("Worker for lang: " + worker.getWorkerLangauge() + " stopped working, needs to respawn it");
            throw new RuntimeException("Worker for lang: " + worker.getWorkerLangauge() + " stopped working, needs to respawn it");
        }
    }

    @Override
    public void passivateObject(WorkerParameters worker, PooledObject<WorkerInstance> p)
            throws Exception {
        System.out.println("Passivating worker for lang: " + worker.getWorkerLangauge());
    }

}

class WorkersPool extends GenericKeyedObjectPool<WorkerParameters, WorkerInstance> {

    public WorkersPool(KeyedPooledObjectFactory<WorkerParameters, WorkerInstance> factory) {
        super(factory);
    }
}

【讨论】:

  • 感谢您的详细回答。听起来这正是我想要的!我会调查这个产品并试一试
  • 良好的池管理库和出色的描述!同意骆驼不是为游泳池管理而设计的。
猜你喜欢
  • 1970-01-01
  • 2018-02-26
  • 1970-01-01
  • 1970-01-01
  • 2018-10-15
  • 2017-07-28
  • 1970-01-01
  • 2012-01-11
  • 2017-09-27
相关资源
最近更新 更多