-
Notifications
You must be signed in to change notification settings - Fork 2
Expand file tree
/
Copy pathExchangerDemo.java
More file actions
116 lines (98 loc) · 2.8 KB
/
Copy pathExchangerDemo.java
File metadata and controls
116 lines (98 loc) · 2.8 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Exchanger;
public class ExchangerDemo{
static Exchanger<DataBuffer> exchanger = new Exchanger<DataBuffer>();
static DataBuffer initialEmptyBuffer = new DataBuffer();
static DataBuffer initialFullBuffer = new DataBuffer("I");
public static void main(String[] args){
class FillingLoop implements Runnable{
int count =0;
@Override
public void run(){
DataBuffer currentBuffer = initialEmptyBuffer;
try{
while(true){
addToBuffer(currentBuffer);
if(currentBuffer.isFull()){
System.out.println(Thread.currentThread() + " as filling thread wants to exchange");
currentBuffer = exchanger.exchange(currentBuffer);
System.out.println(Thread.currentThread() + " receives an exchange");
}
}
}catch(InterruptedException ie){
System.out.println(Thread.currentThread() + " has been interrupted.");
}
}
void addToBuffer(DataBuffer buffer){
String item = "NI_" + count++;
System.out.println(Thread.currentThread() + " adding: " + item);
buffer.add(item);
try{
Thread.sleep((long)(Math.random()*2000)); // simulating lengthier operation
}catch(InterruptedException ie){
ie.printStackTrace();
}
}
}
class EmptyingLoop implements Runnable{
@Override
public void run(){
DataBuffer currentBuffer = initialFullBuffer;
try{
while(true){
takeFromBuffer(currentBuffer);
if(currentBuffer.isEmpty()){
System.out.println(Thread.currentThread() + " as emptying thread, wants to exchange");
currentBuffer = exchanger.exchange(currentBuffer);
System.out.println(Thread.currentThread() + " receives an exchange");
}
}
}catch(InterruptedException ie){
System.out.println(Thread.currentThread() + " has been interrupted.");
}
}
void takeFromBuffer(DataBuffer buffer){
System.out.println(Thread.currentThread() + " removing: " + buffer.remove());
try{
Thread.sleep((long)(Math.random()*2000)); // simulating lengthier operation
}catch(InterruptedException ie){
ie.printStackTrace();
}
}
}
new Thread(new EmptyingLoop()).start();
new Thread(new FillingLoop()).start();
}
}
class DataBuffer{
private final static int MAX_ITEMS = 10;
private List<String> items = new ArrayList<String>();
DataBuffer(){
}
DataBuffer(String prefix){
for(int i=0; i<MAX_ITEMS;i++){
String item = prefix + i;
System.out.printf("Adding %s%n", item);
items.add(item);
}
}
void add(String s){
if(!isFull()){
items.add(s);
}
}
boolean isEmpty(){
return items.isEmpty();
}
boolean isFull(){
return items.size() == MAX_ITEMS;
}
String remove(){
if(!isEmpty()){
return items.remove(0);
}else{
return null;
}
}
}