This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 09975973df [common] Fix permit accounting of
SemaphoredDelegatingExecutor under interruption and rejection (#9511)
09975973df is described below
commit 09975973df647c7d4be1f63aace98e31d65dcfeb
Author: YangJie <[email protected]>
AuthorDate: Wed Sep 2 02:42:16 2026 -0400
[common] Fix permit accounting of SemaphoredDelegatingExecutor under
interruption and rejection (#9511)
---
.../paimon/utils/SemaphoredDelegatingExecutor.java | 69 +++++-
.../utils/SemaphoredDelegatingExecutorTest.java | 234 +++++++++++++++++++++
2 files changed, 298 insertions(+), 5 deletions(-)
diff --git
a/paimon-common/src/main/java/org/apache/paimon/utils/SemaphoredDelegatingExecutor.java
b/paimon-common/src/main/java/org/apache/paimon/utils/SemaphoredDelegatingExecutor.java
index bdbb23796b..a332eed224 100644
---
a/paimon-common/src/main/java/org/apache/paimon/utils/SemaphoredDelegatingExecutor.java
+++
b/paimon-common/src/main/java/org/apache/paimon/utils/SemaphoredDelegatingExecutor.java
@@ -26,8 +26,10 @@ import java.util.List;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;
+import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
/**
* A {@link ForwardingExecutorService} to delegate tasks to limit the number
of tasks executed
@@ -81,7 +83,13 @@ public class SemaphoredDelegatingExecutor extends
ForwardingExecutorService {
return Futures.immediateFailedFuture(e);
}
- return super.submit(new CallableWithPermitRelease(task));
+ CallableWithPermitRelease<T> wrapped = new
CallableWithPermitRelease<>(task);
+ try {
+ return super.submit(wrapped);
+ } catch (RejectedExecutionException e) {
+ wrapped.releasePermit();
+ throw e;
+ }
}
@Override
@@ -93,7 +101,13 @@ public class SemaphoredDelegatingExecutor extends
ForwardingExecutorService {
return Futures.immediateFailedFuture(e);
}
- return super.submit(new RunnableWithPermitRelease(task), result);
+ RunnableWithPermitRelease wrapped = new
RunnableWithPermitRelease(task);
+ try {
+ return super.submit(wrapped, result);
+ } catch (RejectedExecutionException e) {
+ wrapped.releasePermit();
+ throw e;
+ }
}
@Override
@@ -105,18 +119,36 @@ public class SemaphoredDelegatingExecutor extends
ForwardingExecutorService {
return Futures.immediateFailedFuture(e);
}
- return super.submit(new RunnableWithPermitRelease(task));
+ RunnableWithPermitRelease wrapped = new
RunnableWithPermitRelease(task);
+ try {
+ return super.submit(wrapped);
+ } catch (RejectedExecutionException e) {
+ wrapped.releasePermit();
+ throw e;
+ }
}
@Override
public void execute(Runnable command) {
+ boolean acquired = true;
try {
this.queueingPermits.acquire();
} catch (InterruptedException e) {
+ // Semaphore.acquire() throws as soon as the caller carries an
interrupt flag, even
+ // when permits are free, and execute() has no channel for
reporting that the task
+ // was dropped. Run it anyway, as this class always has, but
remember that no permit
+ // backs this task so its wrapper does not hand back one that was
never taken.
Thread.currentThread().interrupt();
+ acquired = false;
}
- super.execute(new RunnableWithPermitRelease(command));
+ RunnableWithPermitRelease wrapped = new
RunnableWithPermitRelease(command, acquired);
+ try {
+ super.execute(wrapped);
+ } catch (RejectedExecutionException e) {
+ wrapped.releasePermit();
+ throw e;
+ }
}
public int getAvailablePermits() {
@@ -146,9 +178,15 @@ public class SemaphoredDelegatingExecutor extends
ForwardingExecutorService {
private class RunnableWithPermitRelease implements Runnable {
private final Runnable delegated;
+ private final AtomicBoolean permitHeld;
RunnableWithPermitRelease(Runnable delegated) {
+ this(delegated, true);
+ }
+
+ RunnableWithPermitRelease(Runnable delegated, boolean permitHeld) {
this.delegated = delegated;
+ this.permitHeld = new AtomicBoolean(permitHeld);
}
@Override
@@ -156,6 +194,19 @@ public class SemaphoredDelegatingExecutor extends
ForwardingExecutorService {
try {
this.delegated.run();
} finally {
+ releasePermit();
+ }
+ }
+
+ /**
+ * Hands the permit back, at most once, and only if one was acquired
for this task. A
+ * delegate that runs the task inline (for example {@link
+ * java.util.concurrent.ThreadPoolExecutor.CallerRunsPolicy}) may both
run the wrapper and
+ * let a {@link RejectedExecutionException} out of the same call, so
the submitting method
+ * and {@link #run()} can each reach this.
+ */
+ void releasePermit() {
+ if (this.permitHeld.compareAndSet(true, false)) {
SemaphoredDelegatingExecutor.this.queueingPermits.release();
}
}
@@ -164,6 +215,7 @@ public class SemaphoredDelegatingExecutor extends
ForwardingExecutorService {
private class CallableWithPermitRelease<T> implements Callable<T> {
private final Callable<T> delegated;
+ private final AtomicBoolean permitHeld = new AtomicBoolean(true);
CallableWithPermitRelease(Callable<T> delegated) {
this.delegated = delegated;
@@ -175,10 +227,17 @@ public class SemaphoredDelegatingExecutor extends
ForwardingExecutorService {
try {
result = this.delegated.call();
} finally {
- SemaphoredDelegatingExecutor.this.queueingPermits.release();
+ releasePermit();
}
return result;
}
+
+ /** Hands the permit back, at most once. See {@link
RunnableWithPermitRelease}. */
+ void releasePermit() {
+ if (this.permitHeld.compareAndSet(true, false)) {
+ SemaphoredDelegatingExecutor.this.queueingPermits.release();
+ }
+ }
}
}
diff --git
a/paimon-common/src/test/java/org/apache/paimon/utils/SemaphoredDelegatingExecutorTest.java
b/paimon-common/src/test/java/org/apache/paimon/utils/SemaphoredDelegatingExecutorTest.java
new file mode 100644
index 0000000000..8dc70dceee
--- /dev/null
+++
b/paimon-common/src/test/java/org/apache/paimon/utils/SemaphoredDelegatingExecutorTest.java
@@ -0,0 +1,234 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.utils;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+/** Test for {@link SemaphoredDelegatingExecutor}. */
+public class SemaphoredDelegatingExecutorTest {
+
+ private static final long TIMEOUT_SECONDS = 10;
+
+ @Test
+ public void testInterruptedExecuteRunsTaskWithoutInflatingPermits() throws
Exception {
+ ExecutorService delegate = Executors.newSingleThreadExecutor();
+ try {
+ SemaphoredDelegatingExecutor executor =
+ new SemaphoredDelegatingExecutor(delegate, 0, true);
+ AtomicBoolean ran = new AtomicBoolean(false);
+ AtomicReference<Throwable> thrown = new AtomicReference<>();
+ AtomicBoolean interrupted = new AtomicBoolean(false);
+ CountDownLatch finished = new CountDownLatch(1);
+
+ Thread submitter =
+ new Thread(
+ () -> {
+ try {
+ executor.execute(() -> ran.set(true));
+
interrupted.set(Thread.currentThread().isInterrupted());
+ } catch (Throwable t) {
+ thrown.set(t);
+ } finally {
+ finished.countDown();
+ }
+ });
+ // Daemon: if a regression ever made the permit wait
uninterruptible, the await
+ // below still fails, and this thread must not keep the surefire
fork alive.
+ submitter.setDaemon(true);
+ submitter.start();
+ awaitWaitingOnPermit(executor);
+
+ // Interrupt once, after the submitter is parked on the semaphore:
the flag
+ // asserted below can then only have been restored by execute()
itself.
+ submitter.interrupt();
+ assertThat(finished.await(TIMEOUT_SECONDS,
TimeUnit.SECONDS)).isTrue();
+
+ assertThat(thrown.get()).isNull();
+ assertThat(interrupted.get()).isTrue();
+
+ // Drain the delegate: the task is submitted without a permit, so
it has to run, and
+ // the count has to stay where it was rather than gain a permit
nobody acquired.
+ delegate.shutdown();
+ assertThat(delegate.awaitTermination(TIMEOUT_SECONDS,
TimeUnit.SECONDS)).isTrue();
+ assertThat(ran.get()).isTrue();
+ assertThat(executor.getAvailablePermits()).isZero();
+ } finally {
+ delegate.shutdownNow();
+ }
+ }
+
+ @Test
+ public void testExecuteWithInterruptFlagAlreadySetKeepsPermitCount()
throws Exception {
+ ExecutorService delegate = Executors.newSingleThreadExecutor();
+ try {
+ SemaphoredDelegatingExecutor executor =
+ new SemaphoredDelegatingExecutor(delegate, 2, true);
+ AtomicBoolean ran = new AtomicBoolean(false);
+ AtomicReference<Throwable> thrown = new AtomicReference<>();
+ AtomicBoolean interrupted = new AtomicBoolean(false);
+ CountDownLatch finished = new CountDownLatch(1);
+
+ // Semaphore.acquire() throws the moment the caller carries an
interrupt flag, even
+ // with both permits free, which is the state a Flink or Spark
task is in while it is
+ // being cancelled. The task still has to run and the count still
has to balance.
+ Thread submitter =
+ new Thread(
+ () -> {
+ Thread.currentThread().interrupt();
+ try {
+ executor.execute(() -> ran.set(true));
+
interrupted.set(Thread.currentThread().isInterrupted());
+ } catch (Throwable t) {
+ thrown.set(t);
+ } finally {
+ finished.countDown();
+ }
+ });
+ submitter.setDaemon(true);
+ submitter.start();
+ assertThat(finished.await(TIMEOUT_SECONDS,
TimeUnit.SECONDS)).isTrue();
+
+ assertThat(thrown.get()).isNull();
+ assertThat(interrupted.get()).isTrue();
+
+ delegate.shutdown();
+ assertThat(delegate.awaitTermination(TIMEOUT_SECONDS,
TimeUnit.SECONDS)).isTrue();
+ assertThat(ran.get()).isTrue();
+ assertThat(executor.getAvailablePermits()).isEqualTo(2);
+ } finally {
+ delegate.shutdownNow();
+ }
+ }
+
+ @Test
+ public void testRejectedByDelegateReleasesPermit() {
+ ExecutorService delegate = Executors.newSingleThreadExecutor();
+ delegate.shutdownNow();
+ SemaphoredDelegatingExecutor executor = new
SemaphoredDelegatingExecutor(delegate, 1, true);
+
+ assertThat(executor.getAvailablePermits()).isEqualTo(1);
+
+ assertThatThrownBy(() -> executor.execute(() -> {}))
+ .isInstanceOf(RejectedExecutionException.class);
+ assertThat(executor.getAvailablePermits()).isEqualTo(1);
+
+ assertThatThrownBy(() -> executor.submit(() -> null))
+ .isInstanceOf(RejectedExecutionException.class);
+ assertThat(executor.getAvailablePermits()).isEqualTo(1);
+
+ assertThatThrownBy(() -> executor.submit(() -> {}))
+ .isInstanceOf(RejectedExecutionException.class);
+ assertThat(executor.getAvailablePermits()).isEqualTo(1);
+
+ assertThatThrownBy(() -> executor.submit(() -> {}, "result"))
+ .isInstanceOf(RejectedExecutionException.class);
+ assertThat(executor.getAvailablePermits()).isEqualTo(1);
+ }
+
+ @Test
+ public void testInlineExecutionReleasesPermitOnlyOnce() {
+ // corePoolSize 1 with a queue of 1: once the worker is busy and the
queue is full,
+ // CallerRunsPolicy runs the next task in the calling thread.
+ ThreadPoolExecutor delegate =
+ new ThreadPoolExecutor(
+ 1,
+ 1,
+ 0L,
+ TimeUnit.MILLISECONDS,
+ new LinkedBlockingQueue<>(1),
+ new ThreadPoolExecutor.CallerRunsPolicy());
+ CountDownLatch block = new CountDownLatch(1);
+ try {
+ delegate.execute(
+ () -> {
+ try {
+ block.await();
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ });
+ delegate.execute(() -> {});
+ SemaphoredDelegatingExecutor executor =
+ new SemaphoredDelegatingExecutor(delegate, 1, true);
+
+ // The wrapper runs inline and releases the permit in its finally,
and the task's
+ // own rejection then comes back out of execute(): releasing again
would inflate
+ // the semaphore past permitCount.
+ assertThatThrownBy(
+ () ->
+ executor.execute(
+ () -> {
+ throw new
RejectedExecutionException("from task");
+ }))
+ .isInstanceOf(RejectedExecutionException.class)
+ .hasMessage("from task");
+ assertThat(executor.getAvailablePermits()).isEqualTo(1);
+ } finally {
+ block.countDown();
+ delegate.shutdownNow();
+ }
+ }
+
+ @Test
+ public void testNormalExecutionKeepsPermitsBalanced() throws Exception {
+ ExecutorService delegate = Executors.newCachedThreadPool();
+ try {
+ SemaphoredDelegatingExecutor executor =
+ new SemaphoredDelegatingExecutor(delegate, 2, true);
+ AtomicInteger completed = new AtomicInteger();
+
+ for (int i = 0; i < 5; i++) {
+ executor.execute(completed::incrementAndGet);
+ }
+
+ delegate.shutdown();
+ assertThat(delegate.awaitTermination(TIMEOUT_SECONDS,
TimeUnit.SECONDS)).isTrue();
+ assertThat(completed.get()).isEqualTo(5);
+ assertThat(executor.getAvailablePermits()).isEqualTo(2);
+ } finally {
+ delegate.shutdownNow();
+ }
+ }
+
+ /** Waits until the submitter thread is queued on the semaphore, bounded
so it cannot hang. */
+ private static void awaitWaitingOnPermit(SemaphoredDelegatingExecutor
executor)
+ throws InterruptedException {
+ long deadline = System.nanoTime() +
TimeUnit.SECONDS.toNanos(TIMEOUT_SECONDS);
+ while (executor.getWaitingCount() == 0 && System.nanoTime() <
deadline) {
+ Thread.sleep(1);
+ }
+ assertThat(executor.getWaitingCount())
+ .as("submitter should be parked on the semaphore")
+ .isEqualTo(1);
+ }
+}