Concurrency Basics
Worker Pipeline
A simple worker pipeline separates preparing work from consuming work. The example stays deterministic by using an in-memory queue and waiting for both threads to finish.
queue
A queue stores items in the order they should be processed.
producer consumer
One thread can produce items while another thread consumes them.
Producer and Consumer
WorkerPipeline.java
Replay: real traced execution (multi-file project)
import java.util.ArrayDeque;
import java.util.Queue;
public class WorkerPipeline {
static class Mailbox {
private final Queue<String> items = new ArrayDeque<>();
synchronized void add(String item) {
items.add(item);
System.out.println("added=" + item);
}
synchronized String take() {
String item = items.remove();
System.out.println("took=" + item);
return item;
}
}
public static void main(String[] args) throws InterruptedException {
int itemCount = 2;
Mailbox mailbox = new Mailbox();
Thread producer = new Thread(() -> {
for (int i = 1; i <= itemCount; i++) {
mailbox.add("job-" + i);
}
});
producer.start();
producer.join();
Thread consumer = new Thread(() -> {
for (int i = 1; i <= itemCount; i++) {
String item = mailbox.take();
System.out.println("processed=" + item);
}
});
consumer.start();
consumer.join();
}
}
import java.util.ArrayDeque;
import java.util.Queue;
public class WorkerPipeline {
static class Mailbox {
private final Queue<String> items = new ArrayDeque<>();
synchronized void add(String item) {
items.add(item);
System.out.println("added=" + item);
}
synchronized String take() {
String item = items.remove();
System.out.println("took=" + item);
return item;
}
}
public static void main(String[] args) throws InterruptedException {
int itemCount = 3;
Mailbox mailbox = new Mailbox();
Thread producer = new Thread(() -> {
for (int i = 1; i <= itemCount; i++) {
mailbox.add("job-" + i);
}
});
producer.start();
producer.join();
Thread consumer = new Thread(() -> {
for (int i = 1; i <= itemCount; i++) {
String item = mailbox.take();
System.out.println("processed=" + item);
}
});
consumer.start();
consumer.join();
}
}
itemCount ← 2, mailbox ← ⟨WorkerPipeline$Mailbox A⟩, producer ← Thread[#19,Thread-0,5,main]
20public static void main(String[] args) throws InterruptedException {21 int itemCount→ 2 = 2; //@itemCount=322 Mailbox mailbox→ ⟨WorkerPipeline$Mailbox A⟩ = new Mailbox();2324 Thread producer→ Thread[#19,Thread-0,5,main] = new Thread(() -> {25 for (int i = 1; i <= itemCount; i++) {26 mailbox.add("job-" + i);27 }28 });2930 producer.start();31 producer.join();producer.join();
30producer.start();31producer.join();synchronized void add(String item)
pass 1 of 28synchronized void add(String itemjob-1) {9 items.add(itemjob-1);10 System.out.println("added=" + itemjob-1);11}outputadded=job-1mailbox.add("job-" + i);
25for (int i = 1; i <= itemCount; i++) {26 mailbox.add("job-" + i);27}synchronized void add(String item)
pass 2 of 28synchronized void add(String itemjob-2) {9 items.add(itemjob-2);10 System.out.println("added=" + itemjob-2);11}outputadded=job-2mailbox.add("job-" + i);
25for (int i = 1; i <= itemCount; i++) {26 mailbox.add("job-" + i);27}consumer ← Thread[#20,Thread-1,5,main]
30producer.start();31producer.join();3233Thread consumer→ Thread[#20,Thread-1,5,main] = new Thread(() -> {34 for (int i = 1; i <= itemCount; i++) {35 String item = mailbox.take();36 System.out.println("processed=" + item);37 }38});3940consumer.start();41consumer.join();item ← job-1
pass 1 of 213synchronized String take() {14 String item→ job-1 = items.remove();15 System.out.println("took=" + itemjob-1);16 return item;outputtook=job-1return item;
15 System.out.println("took=" + item);16 return itemjob-1;17 }18}1920public static void main(String[] args) throws InterruptedException {21 int itemCount = 2; //@itemCount=322 Mailbox mailbox = new Mailbox();2324 Thread producer = new Thread(() -> {25 for (int i = 1; i <= itemCount; i++) {26 mailbox.add("job-" + i);27 }28 });2930 producer.start();31 producer.join();3233 Thread consumer = new Thread(() -> {34 for (int i = 1; i <= itemCount; i++) {35 String item = mailbox.take();36 System.out.println("processed=" + item);37 }38 });3940 consumer.start();41 consumer.join();42}String item = mailbox.take();
34for (int i = 1; i <= itemCount; i++) {35 String item = mailbox.take();36 System.out.println("processed=" + item);37}outputprocessed=job-1item ← job-2
pass 2 of 213synchronized String take() {14 String item→ job-2 = items.remove();15 System.out.println("took=" + itemjob-2);16 return itemjob-2;17}outputtook=job-2String item = mailbox.take();
34for (int i = 1; i <= itemCount; i++) {35 String item = mailbox.take();36 System.out.println("processed=" + item);37}outputprocessed=job-2consumer.join();
40 consumer.start();41 consumer.join();42}
itemCount ← 3, mailbox ← ⟨WorkerPipeline$Mailbox A⟩, producer ← Thread[#19,Thread-0,5,main]
20public static void main(String[] args) throws InterruptedException {21 int itemCount→ 3 = 3;22 Mailbox mailbox→ ⟨WorkerPipeline$Mailbox A⟩ = new Mailbox();2324 Thread producer→ Thread[#19,Thread-0,5,main] = new Thread(() -> {25 for (int i = 1; i <= itemCount; i++) {26 mailbox.add("job-" + i);27 }28 });2930 producer.start();31 producer.join();producer.join();
30producer.start();31producer.join();synchronized void add(String item)
pass 1 of 38synchronized void add(String itemjob-1) {9 items.add(itemjob-1);10 System.out.println("added=" + itemjob-1);11}outputadded=job-1All 3 passes — pass 1 is the card above pass item1 job-1 2 job-2 3 job-3 mailbox.add("job-" + i);
25for (int i = 1; i <= itemCount; i++) {26 mailbox.add("job-" + i);27}mailbox.add("job-" + i);
25for (int i = 1; i <= itemCount; i++) {26 mailbox.add("job-" + i);27}mailbox.add("job-" + i);
25for (int i = 1; i <= itemCount; i++) {26 mailbox.add("job-" + i);27}consumer ← Thread[#20,Thread-1,5,main]
30 producer.start();31 producer.join();3233 Thread consumer→ Thread[#20,Thread-1,5,main] = new Thread(() -> {34 for (int i = 1; i <= itemCount; i++) {35 String item = mailbox.take();36 System.out.println("processed=" + item);37 }38 });3940 consumer.start();41 consumer.join();42}item ← job-1
pass 1 of 313synchronized String take() {14 String item→ job-1 = items.remove();15 System.out.println("took=" + itemjob-1);16 return itemjob-1;17}outputtook=job-1All 3 passes — pass 1 is the card above pass item1 job-1 2 job-2 3 job-3 String item = mailbox.take();
34for (int i = 1; i <= itemCount; i++) {35 String item = mailbox.take();36 System.out.println("processed=" + item);37}outputprocessed=job-1String item = mailbox.take();
34for (int i = 1; i <= itemCount; i++) {35 String item = mailbox.take();36 System.out.println("processed=" + item);37}outputprocessed=job-2String item = mailbox.take();
34for (int i = 1; i <= itemCount; i++) {35 String item = mailbox.take();36 System.out.println("processed=" + item);37}outputprocessed=job-3consumer.join();
40 consumer.start();41 consumer.join();42}