【问题标题】:ExecutorService suitable for a huge amount of short-lived tasksExecutorService 适用于海量短期任务
【发布时间】:2012-09-29 01:00:38
【问题描述】:

是否有适合大量非常短暂的任务的 ExecutorService?我设想在切换到同步等待之前在内部尝试忙等待的东西。保持任务的顺序并不重要,但应该可以强制执行内存一致性(所有任务发生在主线程重新获得控制权之前)。

下面发布的测试包含 100,000 个任务,每个任务连续生成 100 个doubles。它接受线程池的大小作为命令行参数,并始终测试串行版本与并行版本。 (如果没有给出命令行参数,则只测试串行版本。)并行版本使用固定大小的线程池,任务分配甚至不是时间测量的一部分。尽管如此,并行版本从不比串行版本快,我已经尝试了多达 80 个线程(在具有 40 个超线程内核的机器上)。为什么?

import java.util.ArrayList;
import java.util.Random;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class ExecutorPerfTest {
    public static final int TASKS = 100000;
    public static final int SUBTASKS = 100;

    static final ThreadLocal<Random> R = new ThreadLocal<Random>() {
        @Override
        protected synchronized Random initialValue() {
            return new Random();
        }
    };

    public class SeqTest implements Runnable {
        @Override
        public void run() {
            Random r = R.get();
            for (int i = 0; i < TASKS; i++)
                for (int j = 0; j < SUBTASKS; j++)
                    r.nextDouble();
        }
    }

    public class ExecutorTest implements Runnable {
        private final class RandomGenerating implements Callable<Double> {
            @Override
            public Double call() {
                double d = 0;
                Random r = R.get();
                for (int j = 0; j < SUBTASKS; j++)
                    d = r.nextDouble();
                return d;
            }
        }

        private final ExecutorService threadPool;
        private ArrayList<Callable<Double>> tasks = new ArrayList<Callable<Double>>(TASKS);

        public ExecutorTest(int nThreads) {
            threadPool = Executors.newFixedThreadPool(nThreads);
            for (int i = 0; i < TASKS; i++)
                tasks.add(new RandomGenerating());
        }

        public void run() {
            try {
                threadPool.invokeAll(tasks);
            } catch (InterruptedException e) {
                e.printStackTrace();
            } finally {
                threadPool.shutdown();
            }
        }
    }

    public static void main(String[] args) {
        ExecutorPerfTest executorPerfTest = new ExecutorPerfTest();
        if (args.length > 0)
            executorPerfTest.start(new String[]{});
        executorPerfTest.start(args);
    }

    private void start(String[] args) {
        final Runnable r;
        if (args.length == 0) {
            r = new SeqTest();
        }
        else {
            final int nThreads = Integer.parseInt(args[0]);
            r = new ExecutorTest(nThreads);
        }
        System.out.printf("Starting\n");
        long t = System.nanoTime();
        r.run();
        long dt = System.nanoTime() - t;
        System.out.printf("Time: %.6fms\n", 1e-6 * dt);
    }
}

【问题讨论】:

  • 我首先要看的是多线程版本使用了多少个内核。鉴于您的结果,如果只有一两个,我不会感到惊讶。如果是这种情况,并且您在 Linux 上运行,请查看“任务集”。如果你在 Windows 上,我不知道。
  • 您是否尝试过从随机生成更改为其他任务?就像计算 sin 函数一样?只是为了检查这是否可能是原因。
  • 您确实意识到,如果您的代码受 CPU 限制,那么超线程只会减慢您的速度。
  • @GreyBeardedGeek:当我将 SUBTASKS 增加到 100'000 时,我可以看到所有内核都已激活。当 SUBTASKS 为 100 时,如果我使用超过 5 个线程,CPU 消耗将保持在大约 500%。 -- taskset 对我有什么帮助?
  • @jtahlborn:不,我没有,谢谢。不过,我希望一个体面的实现能够很好​​地扩展到 40 个“真实”内核。

标签: java multithreading performance concurrency threadpool


【解决方案1】:

Executors.newFixedThreadPool(nThreads) 的调用将创建一个ThreadPoolExecutor,它从LinkedBlockingQueue 读取任务,即。执行器中的所有线程都将锁定在同一个队列上以检索下一个任务。

鉴于每个任务的大小非常小,并且您引用的线程/cpu 数量相对较多,您的程序很可能运行缓慢,因为会发生高度锁争用和上下文切换。

请注意,LinkedBlockingQueue 使用的 ReentrantLock 的实现在尝试在线程放弃和阻塞之前获取锁时已经旋转了很短的时间(最多大约 1us)。

如果您的用例允许,那么您可能想尝试使用 Disruptor 模式,请参阅 http://lmax-exchange.github.com/disruptor/

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2022-08-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-06-13
    相关资源
    最近更新 更多