【问题标题】:Piping data between threads with Java使用 Java 在线程之间传递数据
【发布时间】:2011-04-09 04:36:32
【问题描述】:

我正在编写一个模拟电影院的多线程应用程序。每个参与的人都是自己的线程,并发必须完全由信号量完成。我遇到的唯一问题是如何基本上链接线程以便它们可以通信(例如通过管道)。

例如:

Customer[1] 是一个线程,它获取一个信号量,让它走到票房。现在客户[1] 必须告诉票房代理他们想看电影“X”。然后 BoxOfficeAgent[1] 也是一个线程,必须检查以确保电影未满,然后卖票或告诉 Customer[1] 选择另一部电影。

如何在保持信号量并发的同时来回传递数据?

另外,我可以从 java.util.concurrent 使用的唯一类是 Semaphore 类。

【问题讨论】:

  • 我要给出的主要提示是不要让“管道”这个词产生心理障碍。多想想有一个包含信息的“盒子”,然后在盒子上放一个便利贴,以便在盒子里有什么有趣的东西可以查看时告诉其他线程。

标签: java multithreading semaphore piping


【解决方案1】:

在线程之间来回传递数据的一种简单方法是使用位于包java.util.concurrent 中的接口BlockingQueue<E> 的实现。

此接口具有向集合中添加具有不同行为的元素的方法:

  • add(E):如果可能添加,否则抛出异常
  • boolean offer(E):如果元素被添加则返回真,否则返回假
  • boolean offer(E, long, TimeUnit): 尝试添加元素,等待指定时间
  • put(E):阻塞调用线程直到元素被添加

它还定义了具有类似行为的元素检索方法:

  • take(): 阻塞直到有可用的元素
  • poll(long, TimeUnit):获取元素或返回null

我最常使用的实现是:ArrayBlockingQueueLinkedBlockingQueueSynchronousQueue

第一个 ArrayBlockingQueue 具有固定大小,由传递给其构造函数的参数定义。

第二个,LinkedBlockingQueue,大小没有限制。它将始终接受任何元素,即offer 将立即返回true,add 永远不会抛出异常。

第三个,对我来说也是最有趣的一个,SynchronousQueue,就是一个管道。你可以把它想象成一个大小为 0 的队列。它永远不会保留一个元素:如果有其他线程试图从中检索元素,这个队列只会接受元素。相反,只有在有另一个线程试图推送元素时,检索操作才会返回一个元素。

为了满足作业仅使用信号量完成同步的要求,你可以从我给你的关于 SynchronousQueue 的描述中得到启发,然后写一些非常相似的东西:

class Pipe<E> {
  private E e;

  private final Semaphore read = new Semaphore(0);
  private final Semaphore write = new Semaphore(1);

  public final void put(final E e) {
    write.acquire();
    this.e = e;
    read.release();
  }

  public final E take() {
    read.acquire();
    E e = this.e;
    write.release();
    return e;
  }
}

请注意,这个类的行为与我描述的 SynchronousQueue 类似。

一旦方法 put(E) 被调用,它就会获取写入信号量,该信号量将留空,以便对同一方法的另一个调用将在其第一行阻塞。然后,此方法存储对正在传递的对象的引用,并释放读取的信号量。此版本将使调用take() 方法的任何线程都可以继续进行。

take() 方法的第一步自然是获取读取信号量,以禁止任何其他线程同时检索元素。在元素被检索并保存在局部变量中之后(练习:如果删除该行 E e = this.e 会发生什么?),该方法释放写入信号量,以便put(E)方法可以被任何线程再次调用,并返回本地变量中保存的内容。

作为重要说明,请注意对正在传递的对象的引用保存在私有字段中,方法take()put(E)都是final。这是最重要的,而且经常被忽略。如果这些方法不是最终的(或者更糟糕的是,该字段不是私有的),继承类将能够改变 take()put(E) 违反合同的行为。

最后,您可以通过使用try {} finally {} 来避免在take() 方法中声明局部变量的需要,如下所示:

class Pipe<E> {
  // ...
  public final E take() {
    try {
      read.acquire();
      return e;
    } finally {
      write.release();
    }
  }
}

这里,这个例子的重点是为了展示try/finally 的使用,这在没有经验的开发人员中是不会被注意到的。显然,在这种情况下,并没有真正的收获。

哦,该死的,我已经为你完成了大部分作业。在报应方面——为了测试你对信号量的了解——为什么不实现 BlockingQueue 合约定义的其他一些方法呢?例如,您可以实现 offer(E) 方法和 take(E, long, TimeUnit)!

祝你好运。

【讨论】:

  • 这不能满足“并发必须完全由信号量完成”的作业要求。 (当然在现实生活中使用高级并发实用程序之一是最好的选择。)
  • 确实!一开始我错过了。
  • 是的,我可以从 java.util.concurrent 中使用的唯一项目是信号量。不能使用其他线程安全的类......那么这应该用管道来完成吗?如果是的话怎么做?
  • @JustinY17:看看我添加到答案中的 Pipe 类!然后,根据我给出的两个,尝试实现 BlockingQueue 定义的其他方法!
【解决方案2】:

考虑一下带有读/写锁的共享内存。

  1. 创建一个缓冲区来放置消息。
  2. 应使用锁/信号量来控制对缓冲区的访问。
  3. 将此缓冲区用于线程间通信。

问候

PKV

【讨论】:

    猜你喜欢
    • 2016-05-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-10-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多