关于流程管理,您可以通过混合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);
}
}