【问题标题】:Java BlockingQueue latency high on LinuxLinux 上的 Java BlockingQueue 延迟很高
【发布时间】:2011-06-02 19:47:12
【问题描述】:

我正在使用 BlockingQueue:s(尝试 ArrayBlockingQueue 和 LinkedBlockingQueue)在我当前正在处理的应用程序中的不同线程之间传递对象。性能和延迟在这个应用程序中相对重要,所以我很好奇使用 BlockingQueue 在两个线程之间传递对象需要多少时间。为了衡量这一点,我编写了一个带有两个线程(一个消费者和一个生产者)的简单程序,我让生产者将时间戳(使用 System.nanoTime() 获取)传递给消费者,请参见下面的代码。

我记得在某个论坛的某个地方读到过,其他人尝试过这个过程大约需要 10 微秒(不知道在什么操作系统和硬件上),所以当它花费大约 30 微秒时我并不感到惊讶我在我的 Windows 7 机器上(英特尔 E7500 核心 2 双核 CPU,2.93GHz),同时在后台运行许多其他应用程序。然而,当我在我们更快的 Linux 服务器(两个 Intel X5677 3.46GHz 四核 CPU,运行内核为 2.6.26-2-amd64 的 Debian 5)上进行相同的测试时,我感到非常惊讶。我预计延迟会低于我的 windows 盒子,但相反它要高得多 - ~75 - 100 微秒!两项测试均使用 Sun 的 Hotspot JVM 版本 1.6.0-23 完成。

有没有其他人在 Linux 上做过类似的测试并得到类似的结果?或者有谁知道为什么它在 Linux 上慢得多(硬件更好),难道是线程切换在 Linux 上比 Windows 慢得多吗?如果是这样的话,Windows 似乎更适合某种应用程序。非常感谢任何帮助我理解相对较高的数字的帮助。

编辑:
在 DaveC 发表评论后,我还做了一个测试,我将 JVM(在 Linux 机器上)限制为单个内核(即所有线程在同一个内核上运行)。这极大地改变了结果——延迟降至 20 微秒以下,即优于 Windows 机器上的结果。我还做了一些测试,我将生产者线程限制在一个内核上,将消费者线程限制在另一个内核上(尝试将它们都放在同一个套接字和不同的套接字上),但这似乎没有帮助 - 延迟仍然是 ~75微秒。顺便说一句,这个测试应用程序几乎是我在执行测试时在机器上运行的所有内容。

有谁知道这些结果是否有意义?如果生产者和消费者在不同的内核上运行,它真的应该慢很多吗?任何意见都非常感谢。

再次编辑(1 月 6 日):
我尝试了对代码和运行环境的不同更改:

  1. 我将 Linux 内核升级到 2.6.36.2(从 2.6.26.2)。内核升级后,测量时间从升级前的 75-100 变为 60 微秒,变化非常小。为生产者和消费者线程设置 CPU 亲和性没有任何效果,除非将它们限制在同一个核心。在同一内核上运行时,测得的延迟为 13 微秒。

  2. 在原始代码中,我让生产者在每次迭代之间休眠 1 秒,以便给消费者足够的时间来计算经过的时间并将其打印到控制台。如果我删除对 Thread.sleep () 的调用,而是让生产者和消费者在每次迭代中都调用 barrier.await() (消费者在将经过的时间打印到控制台后调用它),测量的延迟从60 微秒到 10 微秒以下。如果在同一个内核上运行线程,延迟会低于 1 微秒。谁能解释为什么这会显着减少延迟?我的第一个猜测是,更改的效果是生产者在消费者调用 queue.take() 之前调用 queue.put(),因此消费者永远不必阻塞,但是在使用了 ArrayBlockingQueue 的修改版本之后,我发现这个猜测是错误的——消费者实际上阻止了。如果您有其他猜测,请告诉我。 (顺便说一句,如果我让生产者同时调用 Thread.sleep() 和 barrier.await(),延迟保持在 60 微秒)。

  3. 我还尝试了另一种方法——我没有调用 queue.take(),而是调用 queue.poll(),超时时间为 100 微秒。这将平均延迟降低到 10 微秒以下,但 CPU 密集度当然要高得多(但 CPU 密集度可能比忙等待时要少?)。

再次编辑(1 月 10 日) - 问题已解决:
ninjalj 建议约 60 微秒的延迟是由于 CPU 必须从更深的睡眠状态中唤醒 - 他完全正确!在 BIOS 中禁用 C 状态后,延迟减少到

...

import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CyclicBarrier;

public class QueueTest {

    ArrayBlockingQueue<Long> queue = new ArrayBlockingQueue<Long>(10);
    Thread consumerThread;
    CyclicBarrier barrier = new CyclicBarrier(2);
    static final int RUNS = 500000;
    volatile int sleep = 1000;

    public void start() {
        consumerThread = new Thread(new Runnable() {
            @Override
            public void run() {
                try {
                    barrier.await();
                    for(int i = 0; i < RUNS; i++) {
                        consume();

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

        try {
            barrier.await();
        } catch (Exception e) { e.printStackTrace(); }

        for(int i = 0; i < RUNS; i++) {
            try {
                if(sleep > 0)
                    Thread.sleep(sleep);
                produce();

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

    public void produce() {
        try {
            queue.put(System.nanoTime());
        } catch (InterruptedException e) {
        }
    }

    public void consume() {
        try {
            long t = queue.take();
            long now = System.nanoTime();
            long time = (now - t) / 1000; // Divide by 1000 to get result in microseconds
            if(sleep > 0) {
                System.out.println("Time: " + time);
            }

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

    }

    public static void main(String[] args) {
        QueueTest test = new QueueTest();
        System.out.println("Starting...");
        // Run first once, ignoring results
        test.sleep = 0;
        test.start();
        // Run again, printing the results
        System.out.println("Starting again...");
        test.sleep = 1000;
        test.start();
    }
}

【问题讨论】:

  • 你试过在限制 jvm 只使用一个 cpu 的 linux 机器上进行测试吗?可能有助于确定时间的去向
  • 有趣 - 我试图通过使用命令 'taskset 0x00000001 java QueueTest' 启动应用程序来将其限制在特定的 CPU 上,延迟从大约 75-100 减少到 ~20 微秒!我不确定我是否理解这种方式......
  • @Johan:这些时候,您在多次迭代中报告相同吗?CyclicBarrier 用于协调处理独立任务的线程。您的任务虽然不是独立的。您有生产者和消费者都在等待屏障然后(一旦两个线程都到达障碍点)它们基本上开始在阻塞队列上同步。您可以看到各种调度组合的交错报告各种延迟。
  • 是的,我打印出每次迭代,即使结果不同,范围是 75-100 微秒。我首先以不同的方式使用了 CyclicBarrier,但在上面的代码中它并不是真正必要的(长睡眠类型确保生产者不会获取时间戳并尝试在消费者准备好之前将其放入队列,至少在第一次迭代之后不会)
  • 使用 Sun JVM 1.6.0-21(32 位)在 RHEL 5.5 双核 E6750 2.66G GHz 上运行此程序需要 18-25 微秒。你在这两种情况下都运行 64 位 jvm 吗?

标签: java linux multithreading latency


【解决方案1】:

您的测试不能很好地衡量队列切换延迟,因为您有一个线程读取队列,该线程同步写入System.out(在它所在的位置时执行字符串和长连接),然后再进行。要正确衡量这一点,您需要将此活动移出该线程,并在获取线程中做尽可能少的工作。

您最好只在接受者中进行计算(当时-现在)并将结果添加到其他集合中,该集合由另一个输出结果的线程定期排出。我倾向于通过添加到通过 AtomicReference 访问的适当大小的数组支持结构来做到这一点(因此,报告线程只需在该引用上使用该存储结构的另一个实例 getAndSet 即可获取最新一批结果;例如 make 2列表,设置一个为活动的,每个 xsa 线程唤醒并交换主动和被动线程)。然后,您可以报告一些分布而不是每个结果(例如十分位数范围),这意味着您不会在每次运行时生成大量日志文件并为您打印有用的信息。

FWIW 我同意 Peter Lawrey 所说的时间,如果延迟真的很关键,那么您需要考虑使用适当的 cpu 亲和性忙等待(即为该线程专用一个核心)

1 月 6 日之后编辑

如果我删除对 Thread.sleep () 的调用,而是让生产者和消费者在每次迭代中都调用 barrier.await()(消费者在将经过的时间打印到控制台后调用它),测量的延迟从 60 微秒减少到 10 微秒以下。如果在同一个内核上运行线程,延迟会低于 1 微秒。谁能解释为什么这会显着降低延迟?

您正在查看java.util.concurrent.locks.LockSupport#park(和对应的unpark)和Thread#sleep 之间的区别。大多数 j.u.c.东西是建立在LockSupport 上的(通常通过ReentrantLock 提供的AbstractQueuedSynchronizer 或直接提供),并且这个(在热点中)解析为sun.misc.Unsafe#park(和unpark),这往往最终落入pthread(posix 线程)库。通常pthread_cond_broadcast 用于唤醒,pthread_cond_waitpthread_cond_timedwait 用于BlockingQueue#take 之类的事情。

我不能说我曾经看过 Thread#sleep 是如何实际实现的(因为我从来没有遇到过不是基于条件的等待的低延迟),但我想它会导致它调度程序以比 pthread 信号机制更激进的方式降级,这就是延迟差异的原因。

【讨论】:

  • 感谢您的意见。在这个特定的测试中,我不认为对 System.out 的同步写入应该是一个问题,因为我让生产者线程等待 2 秒,然后它才会在队列上放置一个新的时间戳......除非我错过了这里有什么?您使用两个列表和 AtomicReferences 记录时间戳的解决方案听起来像是在我的“真实”应用程序中记录延迟的好方法。
  • 是的,公平点,我错过了生产前的睡眠。时间以唤醒时间为主。您应该能够通过添加更改向队列提供时间戳的速率的功能来看到这一点,除非盒子上发生了奇怪的事情,否则延迟会随着您增加提供速率而趋于下降。
  • @Matt:您的观点非常好。但同样如此,为什么相同的代码在 Windows 和 Linux 之间差异如此之大仍然无法解释
  • 这主要是一个 CPU 问题,如果你在 n 个硬件平台上运行它,只是改变操作系统(并为每个操作系统使用“等效”配置),那么你应该将 CPU 视为关键因素。在这种情况下,基准测试运行的差异很大。操作系统、硬件、服务器活动 3. linux 的数字,IMO,显然是异常的(不知道 Windows 的数字是否合理)。没有提到 debian 盒子上的活动或如何配置调度和内核非常旧。在测量如此大的延迟时,这些都可能是重要因素。
  • @Johan:可能是因为 CPU 必须从更深的睡眠状态中唤醒。 C2/C3 状态的延迟似乎在该范围内(Proliant DL360 上的 C0/C1/C2/C3 为 0/1/64/96 uS)。
【解决方案2】:

@彼得劳瑞

某些操作使用操作系统调用(例如锁定/循环屏障)

那些不是操作系统(内核)调用。通过简单的 CAS 实现(在 x86 上也带有空闲内存栅栏)

还有一个:除非你知道为什么(你使用它),否则不要使用 ArrayBlockingQueue。

@OP: 看看 ThreadPoolExecutor,它提供了优秀的生产者/消费者框架。

在下方编辑

为了减少延迟(排除繁忙的等待),将队列更改为 SynchronousQueue 在启动消费者之前添加以下内容

...
consumerThread.setPriority(Thread.MAX_PRIORITY);
consumerThread.start();

这是你能得到的最好的。


编辑2: 这里有同步。队列。并且不打印结果。

package t1;

import java.math.BigDecimal;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.SynchronousQueue;

public class QueueTest {

    static final int RUNS = 250000;

    final SynchronousQueue<Long> queue = new SynchronousQueue<Long>();

    int sleep = 1000;

    long[] results  = new long[0];
    public void start(final int runs) throws Exception {
        results = new long[runs];
        final CountDownLatch barrier = new CountDownLatch(1);
        Thread consumerThread = new Thread(new Runnable() {
            @Override
            public void run() {
                barrier.countDown();
                try {

                    for(int i = 0; i < runs; i++) {                        
                        results[i] = consume(); 

                    }
                } catch (Exception e) {
                    return;
                } 
            }
        });
        consumerThread.setPriority(Thread.MAX_PRIORITY);
        consumerThread.start();


        barrier.await();
        final long sleep = this.sleep;
        for(int i = 0; i < runs; i++) {
            try {                
                doProduce(sleep);

            } catch (Exception e) {
                return;
            }
        }
    }

    private void doProduce(final long sleep) throws InterruptedException {
        produce();
    }

    public void produce() throws InterruptedException {
        queue.put(new Long(System.nanoTime()));//new Long() is faster than value of
    }

    public long consume() throws InterruptedException {
        long t = queue.take();
        long now = System.nanoTime();
        return now-t;
    }

    public static void main(String[] args) throws Throwable {           
        QueueTest test = new QueueTest();
        System.out.println("Starting + warming up...");
        // Run first once, ignoring results
        test.sleep = 0;
        test.start(15000);//10k is the normal warm-up for -server hotspot
        // Run again, printing the results
        System.gc();
        System.out.println("Starting again...");
        test.sleep = 1000;//ignored now
        Thread.yield();
        test.start(RUNS);
        long sum = 0;
        for (long elapsed: test.results){
            sum+=elapsed;
        }
        BigDecimal elapsed = BigDecimal.valueOf(sum, 3).divide(BigDecimal.valueOf(test.results.length), BigDecimal.ROUND_HALF_UP);        
        System.out.printf("Avg: %1.3f micros%n", elapsed); 
    }
}

【讨论】:

    【解决方案3】:

    如果延迟很关键并且您不需要严格的 FIFO 语义,那么您可能需要考虑 JSR-166 的 LinkedTransferQueue。它支持消除,以便相反的操作可以交换值,而不是在队列数据结构上同步。这种方法有助于减少争用,实现并行交换,并避免线程休眠/唤醒惩罚。

    【讨论】:

    • 谢谢,我会研究一下LinkedTransferQueue,看看它是否适合我的应用程序。
    【解决方案4】:

    如果可以的话,我会只使用一个 ArrayBlockingQueue。当我使用它时,Linux 上的延迟在 8-18 微秒之间。一些注意事项。

    • 成本主要是唤醒线程所需的时间。当你唤醒一个线程时,它的数据/代码不会在缓存中,所以你会发现,如果你对线程唤醒后发生的事情进行计时,这可能比你重复运行相同的事情要长 2-5 倍。
    • 某些操作使用操作系统调用(例如锁定/循环屏障),在低延迟情况下,这些调用通常比忙等待更昂贵。我建议尝试忙着等待您的制作人,而不是使用 CyclicBarrier。您也可以忙着等待您的消费者,但这在实际系统上可能会非常昂贵。

    【讨论】:

    • 感谢您的回答。我意识到大部分时间都花在唤醒消费者线程上,但我认为 Linux 中的上下文切换会比我得到的数字便宜得多。我不确定我是否完全理解您的第二点 - CyclicBarrier 仅在此处使用一次(并非真正必要),而不是在发送新时间戳时的每次迭代中。
    猜你喜欢
    • 2017-07-05
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-06-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多