标签: blockingqueue

具有可变延迟的ScheduledExecutorService

假设我有一个从java.util.concurrent.BlockingQueue中提取元素并处理它们的任务.

public void scheduleTask(int delay, TimeUnit timeUnit)
{
    scheduledExecutorService.scheduleWithFixedDelay(new Task(queue), 0, delay, timeUnit);
}
Run Code Online (Sandbox Code Playgroud)

如果可以动态更改频率,如何安排/重新安排任务?

  • 我们的想法是获取数据更新流并将它们批量传播到GUI
  • 用户应该能够改变更新的频率

java concurrency executorservice blockingqueue

20
推荐指数
3
解决办法
1万
查看次数

是否可以向ThreadPoolExecutor的BlockingQueue添加任务?

ThreadPoolExecutor的JavaDoc 不清楚是否可以将任务直接添加到BlockingQueue执行程序的后台.文档称调用executor.getQueue()"主要用于调试和监视".

我正ThreadPoolExecutor用自己的方式构建一个BlockingQueue.我保留对队列的引用,以便我可以直接向其添加任务.返回相同的队列,getQueue()因此我假设admonition getQueue()适用于通过我的方式获取的对后备队列的引用.

代码的一般模式是:

int n = ...; // number of threads
queue = new ArrayBlockingQueue<Runnable>(queueSize);
executor = new ThreadPoolExecutor(n, n, 1, TimeUnit.HOURS, queue);
executor.prestartAllCoreThreads();
// ...
while (...) {
    Runnable job = ...;
    queue.offer(job, 1, TimeUnit.HOURS);
}
while (jobsOutstanding.get() != 0) {
    try {
        Thread.sleep(...);
    } catch (InterruptedException e) {
        Thread.currentThread().interrupt();
    }
}
executor.shutdownNow();
Run Code Online (Sandbox Code Playgroud)

queue.offer() VS executor.execute()

据我了解,典型的用途是通过添加任务executor.execute().上面示例中的方法具有阻塞队列的优点,但execute()如果队列已满并立即失败并拒绝我的任务.我也喜欢提交作业与阻塞队列交互; 对我来说,这感觉更"纯粹"的生产者 …

java concurrency producer-consumer executorservice blockingqueue

18
推荐指数
2
解决办法
9474
查看次数

如何阻止BlockingQueue为空?

我正在寻找一种方法来阻止,直到BlockingQueue空.

我知道,在多线程环境中,只要有生产者将物品放入其中BlockingQueue,就会出现队列变空的情况,并且在几纳秒之后它就会充满物品.

但是,如果只有一个生产者,那么它可能希望等待(并阻止),直到队列在停止将项目放入队列后为空.

的Java /伪代码:

// Producer code
BlockingQueue queue = new BlockingQueue();

while (having some tasks to do) {
    queue.put(task);
}

queue.waitUntilEmpty(); // <-- how to do this?

print("Done");
Run Code Online (Sandbox Code Playgroud)

你有什么主意吗?

编辑:我知道包装BlockingQueue和使用额外的条件可以解决问题,我只是问是否有一些预先制定的解决方案和/或更好的替代品.

java concurrency blockingqueue

18
推荐指数
2
解决办法
2万
查看次数

Java BlockingQueue与批处理?

我对与Java BlockingQueue相同的数据结构感兴趣,但它必须能够批处理队列中的对象.换句话说,我希望生产者能够将对象放入队列,但是使用消费者块,take()直到队列达到一定的大小(批量大小).

然后,一旦队列达到批量大小,生产者必须阻止,put()直到消费者已经消耗了队列中的所有元素(在这种情况下,生产者将再次开始生产并且消费者块直到再次到达批次).

是否存在类似的数据结构?或者我应该写它(我不介意),如果有什么东西,我只是不想浪费我的时间.


UPDATE

也许稍微澄清一下事情:

情况总是如下.可以有多个生产者向队列添加项目,但永远不会有多个消费者从队列中获取项目.

现在,问题是这些设置有多个并行和串行.换句话说,生产者为多个队列生产物品,而消费者本身也可以是生产者.这可以更容易地被视为生产者,消费者 - 生产者和最终消费者的有向图.

生产者应该阻塞直到队列为空(@Peter Lawrey)的原因是因为每个都将在一个线程中运行.如果你让它们只是在空间可用的情况下生成,你最终会遇到太多线程试图同时处理太多东西的情况.

也许将它与执行服务相结合可以解决问题?

java queue producer-consumer blockingqueue

17
推荐指数
1
解决办法
1万
查看次数

Java中golang通道的等价物

我有一个要求,我需要从一组阻塞队列中读取.阻塞队列由我正在使用的库创建.我的代码必须从队列中读取.我不想为每个阻塞队列创建一个读者线程.相反,我想使用单个线程(或者可能最多使用2/3线程)轮询它们的数据可用性.由于某些阻塞队列可能长时间没有数据,而其中一些阻塞队列可能会获得数据突发.轮询具有较小超时的队列将起作用,但这根本不高效,因为它仍然需要在所有队列上保持循环,即使其中一些队列长时间没有数据.基本上,我正在寻找一个选择/ epoll(用于套接字)类型的阻塞队列机制.任何线索都非常感谢.

尽管如此,在Go中执行此操作非常简单.下面的代码模拟了与channel和goroutines相同的内容:

package main

import "fmt"
import "time"
import "math/rand"

func sendMessage(sc chan string) {
    var i int

    for {
        i =  rand.Intn(10)
        for ; i >= 0 ; i-- {
            sc <- fmt.Sprintf("Order number %d",rand.Intn(100))
        }
        i = 1000 + rand.Intn(32000);
        time.Sleep(time.Duration(i) * time.Millisecond)
    }
}

func sendNum(c chan int) {
    var i int 
    for  {
        i = rand.Intn(16);
        for ; i >=  0; i-- {
            time.Sleep(20 * time.Millisecond)
            c <- rand.Intn(65534)
        }
        i = 1000 + rand.Intn(24000); …
Run Code Online (Sandbox Code Playgroud)

java concurrency multithreading go blockingqueue

16
推荐指数
2
解决办法
5605
查看次数

BlockingQueue的drainTo()方法的线程安全性

BlockingQueue的文档说批量操作不是线程安全的,尽管它没有明确提到方法drainTo().

BlockingQueue实现是线程安全的.所有排队方法都使用内部锁或其他形式的并发控制以原子方式实现其效果.但是,除非在实现中另有说明,否则批量收集操作addAll,containsAll,retainAll和removeAll不一定以原子方式执行.因此,例如,在仅添加c中的一些元素之后,addAll(c)可能会失败(抛出异常).

drainTo()方法的文档指定无法以线程安全的方式修改阻塞BlockingQueue元素的集合.但是,它没有提到有关drainTo()操作是线程安全的任何事情.

从此队列中删除所有可用元素,并将它们添加到给定集合中.此操作可能比重复轮询此队列更有效.尝试向集合c添加元素时遇到的故障可能导致在抛出关联的异常时元素既不在集合中,也不在集合中.尝试将队列排放到自身会导致IllegalArgumentException.此外,如果在操作正在进行时修改指定的集合,则此操作的行为是不确定的.

那么,drainTo()方法是线程安全的吗?换句话说,如果一个线程在阻塞队列上调用了drainTo()方法,而另一个线程在同一个队列上调用了add()或put(),那么在两个操作结束时队列的状态是否一致?

java thread-safety blockingqueue

14
推荐指数
2
解决办法
6415
查看次数

队列已满,在阻塞队列的深度,需要澄清

从文件内容填充队列时,深度似乎不会增加,因为此实现中未添加元素.

    BlockingQueue<String> q = new SynchronousQueue<String>();
            ...
        fstream = new FileInputStream("/path/to/file.txt");
            ...
        while ((line = br.readLine()) != null) {
            if (q.offer(line))
                System.out.println("Depth: " + q.size()); //0
        }
Run Code Online (Sandbox Code Playgroud)

替换offeradd,抛出异常

Exception in thread "main" java.lang.IllegalStateException: Queue full
  ...
Run Code Online (Sandbox Code Playgroud)

我做错了什么?插入第一个元素后,为什么队列立即满了?

java blockingqueue

14
推荐指数
2
解决办法
2万
查看次数

ThreadPoolExecutor策略

我正在尝试使用ThreadPoolExecutor来安排任务,但遇到了一些问题.这是它陈述的行为:

  1. 如果运行的corePoolSize线程少于corePoolSize,则Executor总是更喜欢添加新线程而不是排队.
  2. 如果corePoolSize或更多线程正在运行,则Executor总是更喜欢排队请求而不是添加新线程.
  3. 如果请求无法排队,则会创建一个新线程,除非它超过maximumPoolSize,在这种情况下,该任务将被拒绝.

我想要的行为是这样的:

  1. 与上述相同
  2. 如果正在运行corePoolSize以上但不到maximumPoolSize线程,则更喜欢在排队时添加新线程,并在添加新线程时使用空闲线程.
  3. 与上述相同

基本上我不希望任何任务被拒绝; 我希望他们在无限队列中排队.但我确实想拥有maximumPoolSize线程.如果我使用无界队列,它在击中coreSize后就不会生成线程.如果我使用有界队列,它会拒绝任务.有没有办法解决?

我现在正在考虑的是在SynchronousQueue上运行ThreadPoolExecutor,但不直接向它提供任务 - 而是将它们提供给单独的无界LinkedBlockingQueue.然后另一个线程从LinkedBlockingQueue进入Executor,如果一个被拒绝,它只是再次尝试,直到它被拒绝.这看起来像是一种痛苦而且有点像黑客 - 有更清洁的方法吗?

java concurrency multithreading executor blockingqueue

13
推荐指数
1
解决办法
5281
查看次数

在java中实现自己的阻塞队列

我知道这个问题之前已被多次询问和回答,但我无法弄清楚互联网上的例子,比如这个那个.

这两种解决方案都检查阻塞队列的数组/队列/链表的空白,notifyAll以及put()方法中的等待线程,反之亦然get().一个评论在第二环节强调了这一情况,并提到这是没有必要的.

所以问题是; 检查队列是否为空,对我来说似乎有点奇怪 完全通知所有等待的线程.有任何想法吗?

提前致谢.

java multithreading synchronization java.util.concurrent blockingqueue

13
推荐指数
2
解决办法
2万
查看次数

Java(Android)多线程进程

我正在开发应用程序(Matt的traceroute windows版本http://winmtr.net/),它创建多个线程,每个线程都有自己的进程(执行ping命令).ThreadPoolExecutor一段时间后关闭所有线程(例如10秒)

ThreadPoolExecutor 使用阻塞队列(在执行之前保存任务)

int NUMBER_OF_CORES = Runtime.getRuntime().availableProcessors();
ThreadPoolExecutor poolExecutor = new ThreadPoolExecutor(
    NUMBER_OF_CORES * 2, NUMBER_OF_CORES * 2 + 2, 10L, TimeUnit.SECONDS, 
    new LinkedBlockingQueue<Runnable>()
);
Run Code Online (Sandbox Code Playgroud)

PingThread.java

private class PingThread extends Thread {

    @Override
    public void run() {
        long pingStartedAt = System.currentTimeMillis();
        // PingRequest is custom object
        PingRequest request = buildPingRequest(params);

        if (!isCancelled() && !Thread.currentThread().isInterrupted()) {

            // PingResponse is custom object

            // Note:
            // executePingRequest uses PingRequest to create a command 
            // which than create …
Run Code Online (Sandbox Code Playgroud)

java multithreading android blockingqueue threadpoolexecutor

13
推荐指数
2
解决办法
1699
查看次数