CompletableFuture.runAsync and thenRunAsync caused the program to block

Viewed 205

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?

0 Answers
Related