diff --git a/paimon-api/src/main/java/org/apache/paimon/utils/ThreadPoolUtils.java b/paimon-api/src/main/java/org/apache/paimon/utils/ThreadPoolUtils.java index b5c28a19a08a..b7151afbbd53 100644 --- a/paimon-api/src/main/java/org/apache/paimon/utils/ThreadPoolUtils.java +++ b/paimon-api/src/main/java/org/apache/paimon/utils/ThreadPoolUtils.java @@ -229,7 +229,7 @@ private static class SequentialBatchIterator implements CloseableBatchIter private Iterator activeResults = Collections.emptyList().iterator(); private T next; - private boolean closed; + private volatile boolean closed; private SequentialBatchIterator( ExecutorService executor, @@ -243,7 +243,7 @@ private SequentialBatchIterator( @Override public boolean hasNext() { - if (!closed) { + if (!isClosed()) { advanceIfNeeded(); } return next != null; @@ -259,6 +259,10 @@ public T next() { return result; } + private boolean isClosed() { + return closed; + } + private void advanceIfNeeded() { while (next == null) { if (activeResults.hasNext()) { @@ -286,7 +290,7 @@ private void advanceIfNeeded() { private void submitBatch(List batch) { ClassLoader classLoader = Thread.currentThread().getContextClassLoader(); for (U input : batch) { - BatchTask task = new BatchTask<>(processor, input, classLoader); + BatchTask task = new BatchTask<>(this, processor, input, classLoader); executor.execute(task); activeTasks.add(task); } @@ -343,6 +347,7 @@ private static class BatchTask implements Runnable { private static final int CANCELLED = 2; private static final int FINISHED = 3; + private final SequentialBatchIterator iterator; private final Function> processor; private final U input; private final ClassLoader classLoader; @@ -354,7 +359,12 @@ private static class BatchTask implements Runnable { private Throwable failure; private volatile boolean failureReported; - private BatchTask(Function> processor, U input, ClassLoader classLoader) { + private BatchTask( + SequentialBatchIterator iterator, + Function> processor, + U input, + ClassLoader classLoader) { + this.iterator = iterator; this.processor = processor; this.input = input; this.classLoader = classLoader; @@ -363,7 +373,7 @@ private BatchTask(Function> processor, U input, ClassLoader classLoad @Override public void run() { synchronized (this) { - if (state == CANCELLED) { + if (isCanceled()) { state = FINISHED; completion.countDown(); return; @@ -386,6 +396,10 @@ public void run() { } } + private boolean isCanceled() { + return state == CANCELLED || iterator.isClosed(); + } + private synchronized void cancel() { if (state == CREATED) { state = CANCELLED;