【问题标题】:Using custom thread pool and executors in Jgroups在 Jgroups 中使用自定义线程池和执行器
【发布时间】:2016-02-28 14:22:21
【问题描述】:

如果我们查看 Jgroup 的 DefaultThreadFactory 内部,则会出现以下代码

protected Thread newThread(Runnable r,String name,String addr,String cluster_name) { String thread_name=getNewThreadName(name, addr, cluster_name); Thread retval=new Thread(r, thread_name); retval.setDaemon(createDaemons); return retval; }

由于使用了新线程,所以我相信在托管服务器环境中这可能会导致问题,也不是一个好的做法。

如果我只是将默认线程工厂和执行器替换为 WebSphere 的托管工厂和执行器,Jgroups 的行为是否仍然相同?

任何指针都会有所帮助..?

更新

我的意图是将 JGroups 与 WebSphere AS 8.5 一起使用。我渴望没有任何非托管线程。我的主要用例是领导选举和一些消息传递。它将用于管理 Spring Integration 轮询器并确保只有一个轮询器在集群中运行。

WAS 8.5 仍然使用 CommonJ api 进行工作管理。

我正在使用 Spring 来抽象任务执行器和调度器。

最初很容易将线程池替换为任务执行器,因为它们共享 Executor api。

TaskScheduler 必须适应您的 TimeScheduler 接口。 它们非常相似,也许从 ScheduledExecutorService 扩展可能是这里的一个选项。我实现了你的接口并委托给 Springs TaskScheduler。

主要问题在于 ThreadFactory。 CommonJ 没有这个概念。为此,我创建了一个 ThreadWrapper,它封装了 Runnable 并在调用“线程的”启动方法时委托给 TaskExecutor。我忽略了线程重命名功能,因为这不会有任何效果。

public Thread newThread(Runnable command) {
    log.debug("newThread");
    RunnableWrapper wrappedCommand = new RunnableWrapper(command);
    return new ThreadWrapper(taskExecutor, wrappedCommand);
}

public synchronized void start() {
    try { 
        taskExecutor.execute(runnableWrapper);          
    } catch (Exception e) {
        throw new UnableToStartException(e);
    }
}

这是我遇到问题的地方。问题在于运输。在许多情况下,使用一些内部可运行文件的 run 方法,例如TP的DiagnosticsHandler、TransferQueueBundler和GMS的ViewHandler都有while语句检查线程。

public class DiagnosticsHandler implements Runnable {
    public void run() {
        byte[] buf;
        DatagramPacket packet;
        while(Thread.currentThread().equals(thread)) {
            //...
        }
    }
}

protected class TransferQueueBundler extends BaseBundler implements Runnable {
    public void run() {
        while(Thread.currentThread() == bundler_thread) {
            //...
        }
    }
}

class ViewHandler implements Runnable {
    public void run() {
        long start_time, wait_time;  // ns
        long timeout=TimeUnit.NANOSECONDS.convert(max_bundling_time, TimeUnit.MILLISECONDS);
        List<Request> requests=new LinkedList<>();
        while(Thread.currentThread().equals(thread) && !suspended) {
            //...
        }
    }
}

这与我们的线程包装不合作。如果这可以改变,以便在存储的线程上调用 equals 方法,则可以覆盖它。

正如您从各种 sn-ps 中看到的那样,有各种实现和保护级别,从包、受保护和公共不等。这增加了扩展类的难度。

完成所有这些后,它仍然没有完全消除非托管线程的问题。

我正在使用创建协议栈的属性文件方法。一旦设置了属性,这将初始化协议栈。移除底层协议创建的 Timer 线程。 TimeScheduler 必须在初始化堆栈之前设置。

一旦完成,所有线程都被管理。

您对如何更轻松地实现这一点有什么建议吗?

【问题讨论】:

    标签: java websphere websphere-8 jgroups


    【解决方案1】:

    是的,您可以注入您的 on 线程池,有关详细信息,请参阅 [1]。

    [1]http://www.jgroups.org/manual/index.html#_replacing_the_default_and_oob_thread_pools

    【讨论】:

    • 感谢您的回复我已经阅读了这篇博文。如果您能告诉我用 JSR 236 并发实用程序替换这些将提供所需的效果,我将非常高兴。我对线程工厂中使用的其他变量(例如 BaseName 等)有点困惑
    • 也许我们在这里混淆了线程池和线程工厂。你只对后者感兴趣吗?在这种情况下,子类 DefaultThreadFactory 用于在需要时在池中创建线程,并提供线程的命名(例如 OOB-Thread11)。我不知道 JSR 236 是否可以用来代替线程池。只要他们定义了一个 Executor,事情就应该起作用。
    • 无论节点是否为协调者,都应该使用注入的工厂。您是否有一些示例代码显示出了什么问题?
    • 嗨,Bela 感谢您的回复...我看到工厂被覆盖但行为发生了变化.. 需要注意的是,我使用 FILE_PING 进行发现,替换工厂后..节点无法相互连接并且每次添加新节点时都会创建一个新集群...我看到的第二个问题是计时器线程工厂没有被覆盖...我收到此警告 10 次,然后创建了一个新集群...。 .org.jgroups.protocols.pbcast.ClientGmsImpl joinInternal WARNING: 31430: JOIN(31430) sent to 14136 timed out (after 2000 ms), on try 1 ...
    • 我不知道 CommonJ 库,但您只需要返回一个线程,我可以在该线程上调用 isAlive() 等。我建议将您的代码打包成一个我可以运行的可编译项目了解您要做什么。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-08-27
    • 2017-12-29
    • 2018-07-28
    • 1970-01-01
    • 2017-06-28
    • 2017-11-23
    相关资源
    最近更新 更多