【问题标题】:java blockingqueue consumer block on full queuejavablockingqueue消费者阻塞在完整队列上
【发布时间】:2016-01-18 10:14:41
【问题描述】:

我正在编写一个小程序,将 Twitter 公共流中的推文放入 HBase 数据库。该程序使用两个线程,一个用于收集推文,一个用于处理它们。 第一个线程使用 twitter4j StatusListener 获取推文并将它们放入容量为 100 的 ArrayBlockingQueue 中。 第二个线程从队列中获取状态,过滤所需的数据并将它们移动到数据库中。 处理比收集状态需要更多时间。

生产者长这样:

public void onStatus(Status status) {
    try {
        this.queue.put(status);
    } catch(Exception ex) {
        ex.printStackTrace();
    }
}

消费者使用take并调用一个函数来处理新的状态:

public void run() {
    try {
        while(true) {
            // Get new status to process
            this.status = this.queue.take();
            this.analyse();
        }
    } catch(Exception ex) {
        ex.printStackTrace();
    }
 }

在主函数中创建并启动了两个线程:

ArrayBlockingQueue<Status> queue_public = new ArrayBlockingQueue<Status>(100);

Thread ta_public = new Thread(new TweetAnalyser(cl.getOptionValue("config"), queue_public));
Thread st_public = new Thread(new RunPublicStream(cl.getOptionValue("config"), queue_public));

ta_public.start();
st_public.start();

程序运行了一段时间没有任何问题,然后突然停止。此时队列已满,消费者似乎无法从中获取新状态。我尝试了几种生产者/消费者模式的变体,但均未成功。不抛出异常。

我不知道要寻找失败。我希望有人能给我一个提示或解决方案。

【问题讨论】:

  • 一段时间是多久?当队列填满时会立即发生故障,还是会高兴一段时间? analyse 中是否有任何类型的 System.exit 类型调用(或从那里调用的方法)?
  • 不,它只是过滤推文的标签、用户名和文本,并将它们放入数据库。
  • 不调用analyse还会失败吗? “停止”到底是什么意思 - JVM 退出,它只是挂起,还是其他什么?
  • 感谢您的帮助!我明天会在没有分析的情况下尝试它,但是我有一个 System.out.println 在 take 和 analyze 命令之间,它没有显示。我所说的停止是指消费者挂起。它不会退出 JVM。如果我在生产者中使用带有超时的优惠,它会收集新的推文,直到我手动退出程序。
  • 分析功能失败,典型的失败是因为打字错误。我在函数中有第二个队列并将名称混合在一起,所以分析本身被阻塞了。

标签: java blockingqueue consumer take


【解决方案1】:

如果使用阻塞队列,请仔细检查代码中的阻塞命令(ArrayBlockingQueue 的 put 和 take),如果使用多个列表,则检查拼写错误。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-01-09
    • 2015-09-25
    • 1970-01-01
    • 1970-01-01
    • 2014-10-21
    相关资源
    最近更新 更多