在线程之间来回传递数据的一种简单方法是使用位于包java.util.concurrent 中的接口BlockingQueue<E> 的实现。
此接口具有向集合中添加具有不同行为的元素的方法:
-
add(E):如果可能添加,否则抛出异常
-
boolean offer(E):如果元素被添加则返回真,否则返回假
-
boolean offer(E, long, TimeUnit): 尝试添加元素,等待指定时间
-
put(E):阻塞调用线程直到元素被添加
它还定义了具有类似行为的元素检索方法:
-
take(): 阻塞直到有可用的元素
-
poll(long, TimeUnit):获取元素或返回null
我最常使用的实现是:ArrayBlockingQueue、LinkedBlockingQueue 和 SynchronousQueue。
第一个 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)!
祝你好运。