forked from examplehub/Java
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathReentrantLockConditionExampleTest.java
More file actions
82 lines (77 loc) · 2.15 KB
/
ReentrantLockConditionExampleTest.java
File metadata and controls
82 lines (77 loc) · 2.15 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
package com.examplehub.basics.thread;
import java.util.ArrayList;
import java.util.LinkedList;
import java.util.Queue;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.ReentrantLock;
import org.junit.jupiter.api.Test;
class ReentrantLockConditionExampleTest {
static class TaskQueue {
Queue<String> queue = new LinkedList<>();
ReentrantLock lock = new ReentrantLock();
Condition condition = lock.newCondition();
public void addTask(String task) {
lock.lock();
try {
queue.add(task);
condition.signalAll();
} finally {
lock.unlock();
}
}
public String getTask() throws InterruptedException {
lock.lock();
try {
while (queue.isEmpty()) {
condition.await();
}
return queue.remove();
} finally {
lock.unlock();
}
}
}
@Test
void test() throws InterruptedException {
var tasks = new TaskQueue();
var threads = new ArrayList<Thread>();
for (int i = 0; i < 5; i++) {
Thread thread =
new Thread(
() -> {
while (true) {
try {
String task = tasks.getTask();
System.out.println("execute: " + task);
} catch (InterruptedException e) {
System.out.println(
"thread " + Thread.currentThread().getName() + " interrupted");
return;
}
}
});
thread.start();
threads.add(thread);
}
var addThread =
new Thread(
() -> {
for (int i = 0; i < 10; ++i) {
String task = "task - " + Math.random();
System.out.println("added: " + task);
tasks.addTask(task);
try {
Thread.sleep(100);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
});
addThread.start();
addThread.join();
Thread.sleep(100);
for (var thread : threads) {
thread.interrupt();
}
}
}