This is an automated email from the ASF dual-hosted git repository. chibenwa pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/james-project.git
commit 9f898f0796a5fc29f7fdaa0a3572a3d137df2980 Author: Quan Tran <[email protected]> AuthorDate: Sat Sep 26 21:45:56 2026 +0700 JAMES-2586 Cancel the Postgres query upon jOOQ reactive timeout Upon TimeoutException, `PostgresExecutor` now asks Postgres to cancel the running query through r2dbc-postgresql `cancelRequest()` (the pooled connection being unwrapped to the driver connection) before propagating the error and releasing the connection. The cancel request goes through a dedicated connection, so it does not depend on the busy one. The connection thus becomes usable again within milliseconds instead of staying busy until the query completes on its own, which prevents a single slow query from cascading timeouts onto unrelated queries. This applies to every execution method, including each page of `executeRowsPaginated`. Co-Authored-By: Claude Opus 5.5 <[email protected]> --- .../backends/postgres/utils/PostgresExecutor.java | 45 +++++++++++++++++++--- 1 file changed, 39 insertions(+), 6 deletions(-) diff --git a/backends-common/postgres/src/main/java/org/apache/james/backends/postgres/utils/PostgresExecutor.java b/backends-common/postgres/src/main/java/org/apache/james/backends/postgres/utils/PostgresExecutor.java index baee1f95e6..dc7db5c768 100644 --- a/backends-common/postgres/src/main/java/org/apache/james/backends/postgres/utils/PostgresExecutor.java +++ b/backends-common/postgres/src/main/java/org/apache/james/backends/postgres/utils/PostgresExecutor.java @@ -51,8 +51,10 @@ import org.slf4j.LoggerFactory; import com.google.common.annotations.VisibleForTesting; +import io.r2dbc.postgresql.api.PostgresqlConnection; import io.r2dbc.spi.Connection; import io.r2dbc.spi.R2dbcBadGrammarException; +import io.r2dbc.spi.Wrapped; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.util.retry.Retry; @@ -121,7 +123,7 @@ public class PostgresExecutor { connection -> dslContext(connection) .flatMap(queryFunction) .timeout(postgresConfiguration.getJooqReactiveTimeout()) - .doOnError(TimeoutException.class, e -> LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, e)) + .onErrorResume(TimeoutException.class, e -> handleTimeout(connection, e)) .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF) .filter(preparedStatementConflictException())) .then(), @@ -153,7 +155,7 @@ public class PostgresExecutor { Flux<Record> recordFlux = dslContext(connection) .flatMapMany(queryFunction) .timeout(postgresConfiguration.getJooqReactiveTimeout()) - .doOnError(TimeoutException.class, e -> LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, e)) + .onErrorResume(TimeoutException.class, e -> handleTimeout(connection, e)) .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF) .filter(preparedStatementConflictException())); @@ -210,7 +212,7 @@ public class PostgresExecutor { connection -> dslContext(connection) .flatMapMany(queryFunction) .timeout(postgresConfiguration.getJooqReactiveTimeout()) - .doOnError(TimeoutException.class, e -> LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, e)) + .onErrorResume(TimeoutException.class, e -> handleTimeout(connection, e)) .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF) .filter(preparedStatementConflictException())), jamesPostgresConnectionFactory::closeConnection))); @@ -222,7 +224,7 @@ public class PostgresExecutor { connection -> dslContext(connection) .flatMap(queryFunction.andThen(Mono::from)) .timeout(postgresConfiguration.getJooqReactiveTimeout()) - .doOnError(TimeoutException.class, e -> LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, e)) + .onErrorResume(TimeoutException.class, e -> handleTimeout(connection, e)) .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF) .filter(preparedStatementConflictException())), jamesPostgresConnectionFactory::closeConnection))); @@ -240,7 +242,7 @@ public class PostgresExecutor { connection -> dslContext(connection) .flatMap(queryFunction) .timeout(postgresConfiguration.getJooqReactiveTimeout()) - .doOnError(TimeoutException.class, e -> LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, e)) + .onErrorResume(TimeoutException.class, e -> handleTimeout(connection, e)) .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF) .filter(preparedStatementConflictException())) .map(Record1::value1), @@ -258,7 +260,7 @@ public class PostgresExecutor { connection -> dslContext(connection) .flatMap(queryFunction) .timeout(postgresConfiguration.getJooqReactiveTimeout()) - .doOnError(TimeoutException.class, e -> LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, e)) + .onErrorResume(TimeoutException.class, e -> handleTimeout(connection, e)) .retryWhen(Retry.backoff(MAX_RETRY_ATTEMPTS, MIN_BACKOFF) .filter(preparedStatementConflictException())), jamesPostgresConnectionFactory::closeConnection))); @@ -288,6 +290,37 @@ public class PostgresExecutor { jamesPostgresConnectionFactory.close().block(); } + private <T> Mono<T> handleTimeout(Connection connection, TimeoutException timeoutException) { + LOGGER.error(JOOQ_TIMEOUT_ERROR_LOG, timeoutException); + return cancelRunningQuery(connection) + .then(Mono.error(timeoutException)); + } + + /** + * Cancelling the reactive pipeline does not stop the query on the Postgres server side: the connection stays busy + * until the query completes, and is handed back to the pool in that state. Asking Postgres to cancel the running query + * ensures the connection is quickly usable again. + */ + private Mono<Void> cancelRunningQuery(Connection connection) { + return unwrapPostgresqlConnection(connection) + .map(postgresqlConnection -> postgresqlConnection.cancelRequest() + .onErrorResume(e -> { + LOGGER.warn("Failed to cancel the timed out Postgres query", e); + return Mono.empty(); + })) + .orElse(Mono.empty()); + } + + private Optional<PostgresqlConnection> unwrapPostgresqlConnection(Connection connection) { + if (connection instanceof PostgresqlConnection postgresqlConnection) { + return Optional.of(postgresqlConnection); + } + if (connection instanceof Wrapped<?> wrapped && wrapped.unwrap() instanceof PostgresqlConnection postgresqlConnection) { + return Optional.of(postgresqlConnection); + } + return Optional.empty(); + } + private Predicate<Throwable> preparedStatementConflictException() { return throwable -> throwable.getCause() instanceof R2dbcBadGrammarException && throwable.getMessage().contains("prepared statement") --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
