通过I/O在线程间进行通信通常很有用。提供线程功能的类库以“管道”的形式对线程间的 I/O 提供了支持。它们在Java I/O 类库中的对应物就是PipedWriter(允许任务向管道写)和PipedReader(允许不同的任务从同一个管道中读取)。这个模型可以看做是“生产者-消费者”问题的变体,这里的管道就是一个封装好的解决方案。管道基本上是一个阻塞队列, 存在于多个引入BlockingQueue之前的Java版本中。
下面是一个简单的例子,两个任务使用一个管道进行通信:
1import java.io.IOException; 2import java.io.PipedReader; 3import java.io.PipedWriter; 4import java.util.Random; 5import java.util.concurrent.ExecutorService; 6import java.util.concurrent.Executors; 7import java.util.concurrent.TimeUnit; 8 9/** 10 * 发送端 11 */ 12class Sender implements Runnable { 13 private Random rand = new Random(47); 14 private PipedWriter writer = new PipedWriter(); 15 public PipedWriter getWriter() { return writer; } 16 @Override 17 public void run() { 18 try { 19 while(true) { 20 for (char c = 'A'; c < 'z'; c++) { 21 writer.write(c); 22 TimeUnit.MILLISECONDS.sleep(rand.nextInt(500)); 23 } 24 } 25 } catch (IOException e) { 26 System.out.println(e + " Sender write Exception"); 27 } catch (InterruptedException e) { 28 System.out.println(e + " Sender sleep Interrupted"); 29 } 30 } 31} 32 33/** 34 * 接收端 35 */ 36class Receiver implements Runnable { 37 private PipedReader reader; 38 public Receiver(Sender sender) throws IOException { 39 reader = new PipedReader(sender.getWriter()); 40 } 41 @Override 42 public void run() { 43 int count = 0; 44 try { 45 while(true) { 46 //在读取到内容之前,会一直阻塞 47 char s = (char)reader.read(); 48 System.out.print("Read: " + s + ", "); 49 if (++count % 5 == 0) { 50 System.out.println(); 51 } 52 } 53 } catch (IOException e) { 54 System.out.println(e + " Receiver read Exception."); 55 } 56 } 57} 58 59public class PipedIO { 60 public static void main(String[] args) throws Exception { 61 Sender sender = new Sender(); 62 Receiver receiver = new Receiver(sender); 63 ExecutorService exec = Executors.newCachedThreadPool(); 64 exec.execute(sender); 65 exec.execute(receiver); 66 TimeUnit.SECONDS.sleep(5); 67 exec.shutdownNow(); 68 } 69}
执行结果(可能的结果):
1Read: A, Read: B, Read: C, Read: D, Read: E, 2Read: F, Read: G, Read: H, Read: I, Read: J, 3Read: K, Read: L, Read: M, Read: N, Read: O, 4Read: P, Read: Q, Read: R, Read: S, Read: T, 5Read: U, java.io.InterruptedIOException Receiver read Exception. 6java.lang.InterruptedException: sleep interrupted Sender sleep Interrupted
Sender和Receiver代表了需要互相通信的两个任务。Sender创建了一个PipedWriter,它是一个单独的对象;但是对于Receiver,PipedReader的建立必须在构造器中与一个PipedWriter相关联。就是说,PipedReader与PipedWriter的构造可以通过如下两种方式:
1//方式一:先构造PipedReader,再通过它构造PipedWriter。 2PipedReader reader = new PipedReader(); 3PipedWriter writer = new PipedWriter(reader); 4 5//方式二:先构造PipedWriter,再通过它构造PipedReader。 6PipedWriter writer2 = new PipedWriter(); 7PipedReader reader2 = new PipedReader(writer2);
Sender把数据放进Writer,然后休眠一段时间(随机数)。然而,Receiver没有sleep()和wait。但当它调用read()时,**如果没有更多的数据,**管道将自动阻塞。
注意Sender和Receiver是在main()中启动的,即对象构造彻底完毕之后。如果你启动了一个没有构造完毕的对象,在不同的平台上管道可能会产生不一致的行为(注意,BlockingQueue使用起来更加健壮而容易)。
在shutdownNow()被调用时,可以看到PipedReader与普通I/O之间最重要的差异——PipedReader是可以中断的。如果你将reader.read()替换为System.in.read(),那么interrupt()将不能打断read()调用。