diff --git a/concurrent/src/main/java/concurrent/model/QueueConsumer.java b/concurrent/src/main/java/concurrent/model/QueueConsumer.java new file mode 100644 index 0000000..930c28a --- /dev/null +++ b/concurrent/src/main/java/concurrent/model/QueueConsumer.java @@ -0,0 +1,28 @@ +package concurrent.model; + +import lombok.AllArgsConstructor; +import lombok.SneakyThrows; + +import java.util.concurrent.BlockingQueue; + +@AllArgsConstructor +public class QueueConsumer extends Thread { + BlockingQueue blockingQueue; + Integer consumerNumber; + + @SneakyThrows + @Override + public void run() { + work(); + } + + public void work() throws InterruptedException { + while (true) { + Integer take = blockingQueue.take(); + System.out.printf("consumer-queue-{%s}-{%d}\n", consumerNumber, take); + Thread.sleep(1000); + } + } + + +} diff --git a/concurrent/src/main/java/concurrent/model/QueueProducer.java b/concurrent/src/main/java/concurrent/model/QueueProducer.java new file mode 100644 index 0000000..02e04a2 --- /dev/null +++ b/concurrent/src/main/java/concurrent/model/QueueProducer.java @@ -0,0 +1,31 @@ +package concurrent.model; + +import lombok.AllArgsConstructor; +import lombok.SneakyThrows; + +import java.util.Random; +import java.util.concurrent.BlockingQueue; + +@AllArgsConstructor +public class QueueProducer extends Thread { + BlockingQueue blockingQueue; + Integer producerNumber; + + @SneakyThrows + @Override + public void run() { + work(); + } + + public void work() throws InterruptedException { + while (true) { + int abs = Math.abs(new Random().nextInt()); + blockingQueue.put(abs); + System.out.printf("producer-queue-{%s}-{%d}\n", producerNumber, abs); + Thread.sleep(1000); + } + + } + + +} diff --git a/concurrent/src/main/java/concurrent/model/WaitConsumer.java b/concurrent/src/main/java/concurrent/model/WaitConsumer.java new file mode 100644 index 0000000..c868bb4 --- /dev/null +++ b/concurrent/src/main/java/concurrent/model/WaitConsumer.java @@ -0,0 +1,41 @@ +package concurrent.model; + +import lombok.AllArgsConstructor; +import lombok.SneakyThrows; + +import java.util.Queue; + +@AllArgsConstructor +public class WaitConsumer extends Thread { + private final Object LOCK; + private final Queue buffer; + + @SneakyThrows + @Override + public void run() { + work(); + } + + public void work() throws InterruptedException { + while (true) { + try { + Thread.sleep(1000); + } + catch (InterruptedException e) { + e.printStackTrace(); + } + synchronized (LOCK) { + while (buffer.isEmpty()) { + LOCK.wait(); + } + Integer poll = buffer.poll(); + System.out.println(Thread.currentThread() + .getName() + "消费者消费" + poll); + LOCK.notifyAll(); + } + + } + } + + +} diff --git a/concurrent/src/main/java/concurrent/model/WaitProducer.java b/concurrent/src/main/java/concurrent/model/WaitProducer.java new file mode 100644 index 0000000..f28b8d3 --- /dev/null +++ b/concurrent/src/main/java/concurrent/model/WaitProducer.java @@ -0,0 +1,43 @@ +package concurrent.model; + +import lombok.AllArgsConstructor; +import lombok.RequiredArgsConstructor; + +import java.util.Queue; +import java.util.Random; + +@AllArgsConstructor +public class WaitProducer extends Thread { + Integer maxSize = 0; + private final Object LOCK; + private final Queue buffer; + + @Override + public void run() { + work(); + } + + public void work() { + while (true) { + synchronized (LOCK) { + while (buffer.size() == maxSize) { + try { + LOCK.wait(); + } + catch (Exception e) { + e.printStackTrace(); + } + } + int abs = Math.abs(new Random().nextInt()); + buffer.offer(abs); + System.out.println(Thread.currentThread() + .getName() + "生产者生产,目前总共有"); + LOCK.notifyAll(); + } + + } + + } + + +} diff --git a/concurrent/src/test/java/concurrent/model/ModelTest.java b/concurrent/src/test/java/concurrent/model/ModelTest.java new file mode 100644 index 0000000..cd55b00 --- /dev/null +++ b/concurrent/src/test/java/concurrent/model/ModelTest.java @@ -0,0 +1,50 @@ +package concurrent.model; + +import org.junit.jupiter.api.Test; + +import java.util.LinkedList; +import java.util.Queue; +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.LinkedBlockingDeque; + +public class ModelTest { + @Test + public void ArrayBlockingQueueModelTest() throws InterruptedException { + ArrayBlockingQueue arrayBlockingQueue = new ArrayBlockingQueue<>(5000); + for (int i = 0; i < 1000; i++){ + new QueueProducer(arrayBlockingQueue, i).start(); + } + for (int i = 0; i < 200000; i++){ + new QueueConsumer(arrayBlockingQueue, i).start(); + } + + Thread.sleep(100000000); + } + + @Test + public void LinkedBlockingQueueModelTest() throws InterruptedException { + LinkedBlockingDeque arrayBlockingQueue = new LinkedBlockingDeque<>(1000); + for (int i = 0; i < 1000; i++){ + new QueueProducer(arrayBlockingQueue, i).start(); + } + for (int i = 0; i < 200000; i++){ + new QueueConsumer(arrayBlockingQueue, i).start(); + } + + Thread.sleep(100000000); + } + + @Test + public void WaitModelTest() throws InterruptedException { + Object o = new Object(); + Queue queue = new LinkedList(); + Integer maxSize = 100; + for (int i = 0; i < 15; i++){ + new WaitProducer(maxSize, o, queue).start(); + } + for (int i = 0; i < 2; i++){ + new WaitConsumer(o, queue).start(); + } + Thread.sleep(10000000); + } +}