The goal I want to achieve is: serial processing of events with the same id, parallel processing of events with different ids. After learning a lot of reference materials, I found that using CompletableFuture can accomplish the above goals very well. However, in an accidental test I found that runAsync and thenRunAsync caused the program to block, and the program did not deadlock at this time. Does anyone know why this is? Below is my code:
public class Demo {
private static LongAdder longAdder = new LongAdder();
private static ThreadPoolExecutor threadPoolExecutor = new ThreadPoolExecutor(
1 + 1 + 2,
1 + 1 + 2,
0,
TimeUnit.SECONDS,
new SynchronousQueue<>(),
Executors.defaultThreadFactory(),
new ThreadPoolExecutor.AbortPolicy()
);
private static ArrayBlockingQueue<Event> arrayBlockingQueue = new ArrayBlockingQueue<>(20000);
private static ConcurrentHashMap<Integer, CompletableFuture<Void>> dispatch = new ConcurrentHashMap<>();
public static void main(String[] args) {
threadPoolExecutor.execute(() -> {
for (int i = 0; i < 1000; i++) {
try {
Event event = new Event(0, "message" + i);
arrayBlockingQueue.put(event);
} catch (Exception e) {
e.printStackTrace();
}
}
for (int i = 0; i < 1000; i++) {
try {
Event event = new Event(1, "message" + i);
arrayBlockingQueue.put(event);
} catch (Exception e) {
e.printStackTrace();
}
}
});
threadPoolExecutor.execute(() -> {
while (true) {
try {
Event event = arrayBlockingQueue.poll();
if (event == null) {
continue;
}
dispatch.compute(event.getId() % 2, (eventId, completableFuture) ->
(completableFuture == null)
? CompletableFuture.runAsync(() -> {
System.out.println(Thread.currentThread().getName() + " <runAsync>: " + event);
longAdder.increment();
System.out.println("->" + longAdder.longValue());
}, threadPoolExecutor)
: completableFuture.thenRunAsync(() -> {
System.out.println(Thread.currentThread().getName() + " <thenRunAsync>: " + event);
longAdder.increment();
System.out.println("->" + longAdder.longValue());
}, threadPoolExecutor));
} catch (Exception e) {
e.printStackTrace();
}
}
});
}
}
A typical result is:
.
.
.
pool-1-thread-3 <thenRunAsync>: Event(id=1, message=message996)
->1049
pool-1-thread-4 <thenRunAsync>: Event(id=1, message=message997)
->1050
pool-1-thread-3 <thenRunAsync>: Event(id=1, message=message998)
->1051
pool-1-thread-4 <thenRunAsync>: Event(id=1, message=message999)
->1052
Explain that only 1052 events were processed, why is this happening?