标签: blockingqueue

重大的Java BlockingQueue

所以我在生产者/消费者类型应用程序中使用固定大小的BlockingQueue [ArrayBlockingQueue],但我希望用户能够动态更改队列大小.问题是没有BlockingQueue实现允许在创建后更改容量.以前有人见过这个吗?有任何想法吗?

java concurrency blockingqueue

8
推荐指数
1
解决办法
3477
查看次数

如何并行等待多个阻塞队列?

我有两个独立的阻塞队列.客户端通常使用第二个阻塞队列中的第一个来检索要处理的元素.

在某些情况下,客户端对来自两个阻塞队列的元素感兴趣,无论哪个队列首先提供数据.

客户端如何同时等待两个队列?

java concurrency blockingqueue

7
推荐指数
1
解决办法
1141
查看次数

C++ pthread阻塞队列死锁(我认为)

我遇到了pthreads的问题,我认为我遇到了僵局.我创建了一个我认为正在工作的阻塞队列,但在做了一些测试后我发现如果我尝试取消阻塞在blocking_queue上的多个线程,我似乎陷入僵局.

阻塞队列非常简单,如下所示:

template <class T> class Blocking_Queue
{
public:
    Blocking_Queue()
    {
        pthread_mutex_init(&_lock, NULL);
        pthread_cond_init(&_cond, NULL);
    }

    ~Blocking_Queue()
    {
        pthread_mutex_destroy(&_lock);
        pthread_cond_destroy(&_cond);
    }

    void put(T t)
    {
        pthread_mutex_lock(&_lock);
        _queue.push(t);
        pthread_cond_signal(&_cond);
        pthread_mutex_unlock(&_lock);
    }

     T pull()
     {
        pthread_mutex_lock(&_lock);
        while(_queue.empty())
        {
            pthread_cond_wait(&_cond, &_lock);
        }

        T t = _queue.front();
        _queue.pop();

        pthread_mutex_unlock(&_lock);

        return t;
     }

priavte:
    std::queue<T> _queue;
    pthread_cond_t _cond;
    pthread_mutex_t _lock;
}
Run Code Online (Sandbox Code Playgroud)

为了测试,我创建了4个线程来拉动这个阻塞队列.我向阻塞队列添加了一些print语句,每个线程都进入pthread_cond_wait()方法.但是,当我尝试在每个线程上调用pthread_cancel()和pthread_join()时,程序就会挂起.

我还用一个线程对它进行了测试,它完美无缺.

根据文档,pthread_cond_wait()是一个取消点,因此在这些线程上调用cancel会导致它们停止执行(这只适用于1个线程).但是pthread_mutex_lock不是取消点.可能会在调用pthread_cancel()时发生某些事情,取消的线程在终止之前获取互斥锁并且不解锁它,然后当下一个线程被取消时它无法获取互斥锁和死锁?或者还有别的我做错了.

任何建议都很可爱.谢谢 :)

c++ deadlock pthreads blockingqueue

6
推荐指数
1
解决办法
4643
查看次数

Java:Producer = Consumer,如何知道何时停止?

我有几个工人,使用ArrayBlockingQueue.

每个worker从队列中获取一个对象,对其进行处理,结果可以得到几个对象,这些对象将被放入队列中进行进一步处理.所以,工人=生产者+消费者.

工人:

public class Worker implements Runnable
{
    private BlockingQueue<String> processQueue = null;

    public Worker(BlockingQueue<String> processQueue)
    {
        this.processQueue = processQueue;
    }

    public void run()
    {
        try
        {
            do
            {
                String item = this.processQueue.take();
                ArrayList<String> resultItems = this.processItem(item);

                for(String resultItem : resultItems)
                {
                    this.processQueue.put(resultItem);
                }
            }
            while(true);
        }
        catch(Exception)
        {
            ...
        }
    }

    private ArrayList<String> processItem(String item) throws Exception
    {
        ...
    }
}
Run Code Online (Sandbox Code Playgroud)

主要:

public class Test
{
    public static void main(String[] args) throws Exception
    {
        new Test().run();
    }

    private …
Run Code Online (Sandbox Code Playgroud)

java producer-consumer blockingqueue

6
推荐指数
1
解决办法
2557
查看次数

以生产者/消费者模式暂停消费者

我有生产者和消费者联系BlockingQueue.

消费者从队列中等待记录并处理它:

Record r = mQueue.take();
process(r);
Run Code Online (Sandbox Code Playgroud)

我需要从其他线程暂停这个过程一段时间.怎么实现呢?

现在我认为实现它,但它似乎是不好的解决方案:

private Object mLock = new Object();
private boolean mLocked = false;

public void lock() {
    mLocked = true;
}

public void unlock() {
    mLocked = false;
    mLock.notify();

}

public void run() {
    ....
            Record r = mQueue.take();
            if (mLocked) {
                mLock.wait();
            }
            process(r);
}
Run Code Online (Sandbox Code Playgroud)

java multithreading producer-consumer blockingqueue

6
推荐指数
1
解决办法
238
查看次数

Java:线程生产者消费者等待生成数据的最有效方法是什么

使用BlockingQueue消耗生成的数据时,等待数据出现的最有效方法是什么?

场景:

步骤1)数据列表将是添加时间戳的数据存储.这些时间戳需要按最接近当前时间优先级排序.此列表可能为空.线程将时间戳插入其中.生产

步骤2)我想在另一个线程中使用此处的数据,该线程将从数据中获取时间戳并检查它们是否在当前时间之后.消费者然后生产

步骤3)如果它们在当前时间之后,则将它们发送到另一个线程以供消费和处理.在此处理时间戳数据后,从步骤1数据存储中删除.消费然后编辑原始列表.

在下面的代码中,数据字段引用步骤1中的数据存储.结果是在当前时间之后已发送的时间戳列表.步骤2.然后将结果消耗步骤3.

private BlockingQueue<LocalTime> data;
private final LinkedBlockingQueue<Result> results = new LinkedBlockingQueue<Result>();

@Override
public void run() {
  while (!data.isEmpty()) {
    for (LocalTime dataTime : data) {
      if (new LocalTime().isAfter(dataTime)) {
        results.put(result);
      }
    }
  }
}
Run Code Online (Sandbox Code Playgroud)

问题 等待数据列表中可能可能为空的数据的最有效方法是什么?专注于:

while (!data.isEmpty())
Run Code Online (Sandbox Code Playgroud)

以前的问题.

java multithreading producer-consumer blockingqueue java-7

6
推荐指数
1
解决办法
708
查看次数

JMSException InterruptedIOException - 生产者线程被中断

我正在获得JMS异常,似乎队列没有退出或者它没有完成任务.

消息是异步的,它在大多数情况下工作正常,但有时会低于异常.似乎听众在另一边继续听,但在制作人一方得到了这个例外.

javax.jms.JMSException: java.io.InterruptedIOException
at org.apache.activemq.util.JMSExceptionSupport.create(JMSExceptionSupport.java:62)
at org.apache.activemq.ActiveMQConnection.syncSendPacket(ActiveMQConnection.java:1266)
at org.apache.activemq.ActiveMQConnection.ensureConnectionInfoSent(ActiveMQConnection.java:1350)
at org.apache.activemq.ActiveMQConnection.start(ActiveMQConnection.java:495)
at com.vtech.mqservice.response.SendResponse.sendResponseToQueue(SendResponse.java:44)


Caused by: java.io.InterruptedIOException
at org.apache.activemq.transport.WireFormatNegotiator.oneway(WireFormatNegotiator.java:102)
at org.apache.activemq.transport.MutexTransport.oneway(MutexTransport.java:40)
at org.apache.activemq.transport.ResponseCorrelator.asyncRequest(ResponseCorrelator.java:74)
at org.apache.activemq.transport.ResponseCorrelator.request(ResponseCorrelator.java:79)
at org.apache.activemq.ActiveMQConnection.syncSendPacket(ActiveMQConnection.java:1244)
... 0 more
Run Code Online (Sandbox Code Playgroud)

请帮我确定导致生产者线程被中断的原因.

我将activemq版本升级到最新版本并将更新调查结果.

请指出我正确的方向?

更新:正在使用的ActiveMQ版本是activemq-all-5.3.0.jar

java activemq-classic jms rabbitmq blockingqueue

6
推荐指数
1
解决办法
1477
查看次数

使用ArrayBlockingQueue会使进程变慢

我刚刚使用ArrayBlockingQueue进行多线程处理.但它似乎放慢了速度而不是加速.你能帮助我吗?我基本上是导入一个文件(大约300k行)并解析它们并将它们存储在数据库中

public class CellPool {
private static class RejectedHandler implements RejectedExecutionHandler {
    @Override
    public void rejectedExecution(Runnable arg0, ThreadPoolExecutor arg1) {
      System.err.println(Thread.currentThread().getName() + " execution rejected: " + arg0);     
    }
  }

  private static class Task implements Runnable {
    private JSONObject obj;

    public Task(JSONObject obj) {
      this.obj = obj;
    }

    @Override
    public void run() {
      try {
        Thread.sleep(1);
        runThis(obj);
      } catch (InterruptedException e) {
        e.printStackTrace();
      }
    }

    public void runThis(JSONObject obj) {
        //where the rows are parsed and stored in the DB, …
Run Code Online (Sandbox Code Playgroud)

java multithreading blockingqueue threadpoolexecutor

6
推荐指数
1
解决办法
500
查看次数

Python:Kombu + RabbitMQ死锁 - 队列被阻止或阻塞

问题

我有一个RabbitMQ服务器,作为我的一个系统的队列中心.在过去一周左右,它的制作人每隔几个小时就会完全停止.

我试过了什么

蛮力

  • 停止消费者会释放锁定几分钟,但随后阻止返回.
  • 重启RabbitMQ解决了几个小时的问题.
  • 我有一些自动脚本可以完成丑陋的重启,但显然远非正确的解决方案.

分配更多内存

cantSleepNow的回答之后,我将分配给RabbitMQ内存增加到90%.服务器有16GB的内存,消息数量不是很高(每天数百万),所以这似乎不是问题.

从命令行:

sudo rabbitmqctl set_vm_memory_high_watermark 0.9
Run Code Online (Sandbox Code Playgroud)

并与/etc/rabbitmq/rabbitmq.config:

[
   {rabbit,
   [
     {loopback_users, []},
     {vm_memory_high_watermark, 0.9}
   ]
   }
].
Run Code Online (Sandbox Code Playgroud)

代码与设计

我为所有消费者和生产者使用Python.

生产者

生产者是提供呼叫的API服务器.每当呼叫到达时,都会打开一个连接,发送一条消息并关闭连接.

from kombu import Connection

def send_message_to_queue(host, port, queue_name, message):
    """Sends a single message to the queue."""
    with Connection('amqp://guest:guest@%s:%s//' % (host, port)) as conn:
        simple_queue = conn.SimpleQueue(name=queue_name, no_ack=True)
        simple_queue.put(message)
        simple_queue.close()
Run Code Online (Sandbox Code Playgroud)

消费者

消费者彼此略有不同,但通常使用以下模式 - 打开连接,并等待消息到达.连接可以长时间保持打开状态(比如几天).

with Connection('amqp://whatever:whatever@whatever:whatever//') as conn:
    while True:
        queue = …
Run Code Online (Sandbox Code Playgroud)

python deadlock rabbitmq blockingqueue kombu

6
推荐指数
1
解决办法
1063
查看次数

BlockingQueue 和 putAll

有人知道为什么 java 的 BlockingQueue 没有 putAll 方法吗?这样的方法有问题吗?有什么好的方法可以解决这个问题而不必完全重新实现 BlockingQueue?

java blockingqueue

5
推荐指数
1
解决办法
690
查看次数