From 0d91a4ab9e395468e1eafc26aeaa321fc524dacf Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 26 Sep 2026 18:56:37 +0800 Subject: [PATCH] [core] Surface JDBC lock acquisition failures MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit lockAcquire converted every SQLException of the lock INSERT into false, so anything but a held lock — a lock key exceeding the column size, a missing lock table, a connection failover — silently spun for the whole acquire timeout and then failed without a cause. Treat only constraint violations as a held lock and rethrow the rest. Assisted-by: GLM-5.3 --- .../jdbc/AbstractDistributedLockDialect.java | 23 +++++++++++++++- .../apache/paimon/jdbc/JdbcCatalogTest.java | 26 +++++++++++++++++++ 2 files changed, 48 insertions(+), 1 deletion(-) diff --git a/paimon-core/src/main/java/org/apache/paimon/jdbc/AbstractDistributedLockDialect.java b/paimon-core/src/main/java/org/apache/paimon/jdbc/AbstractDistributedLockDialect.java index b8bf3e62d853..94f3b790792f 100644 --- a/paimon-core/src/main/java/org/apache/paimon/jdbc/AbstractDistributedLockDialect.java +++ b/paimon-core/src/main/java/org/apache/paimon/jdbc/AbstractDistributedLockDialect.java @@ -60,7 +60,13 @@ public boolean lockAcquire(JdbcClientPool connections, String lockId, long timeo preparedStatement.setLong(2, timeoutMillSeconds / 1000); return preparedStatement.executeUpdate() > 0; } catch (SQLException ex) { - return false; + // only a constraint violation means the lock is held; swallowing + // anything else turns a broken lock table into a silent full-timeout + // spin that hides the root cause + if (isConstraintViolation(ex)) { + return false; + } + throw ex; } }); } @@ -96,4 +102,19 @@ public int tryReleaseTimedOutLock(JdbcClientPool connections, String lockId) } public abstract String getTryReleaseTimedOutLock(); + + private static boolean isConstraintViolation(SQLException ex) { + // SQLState 23xxx covers MySQL/PostgreSQL; SQLite reports constraint failures + // without a standard SQLState + if (ex.getSQLState() != null && ex.getSQLState().startsWith("23")) { + return true; + } + if (ex instanceof java.sql.SQLIntegrityConstraintViolationException) { + return true; + } + String message = ex.getMessage(); + return message != null + && (message.toLowerCase(java.util.Locale.ROOT).contains("constraint") + || message.toLowerCase(java.util.Locale.ROOT).contains("duplicate")); + } } diff --git a/paimon-core/src/test/java/org/apache/paimon/jdbc/JdbcCatalogTest.java b/paimon-core/src/test/java/org/apache/paimon/jdbc/JdbcCatalogTest.java index 55808e1e426e..f3365ff05dbe 100644 --- a/paimon-core/src/test/java/org/apache/paimon/jdbc/JdbcCatalogTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/jdbc/JdbcCatalogTest.java @@ -202,6 +202,32 @@ public void testAcquireLockFail() throws SQLException, InterruptedException { .isFalse(); } + @Test + public void testBrokenLockInsertSurfacesRootCause() throws Exception { + // an INSERT failure other than a constraint violation is a configuration or + // connection error, not a held lock: the old catch-all swallowed it and spun + // for the whole acquire timeout + java.sql.SQLException failure = new java.sql.SQLException("connection reset"); + java.sql.Connection connection = org.mockito.Mockito.mock(java.sql.Connection.class); + org.mockito.Mockito.when( + connection.prepareStatement(org.mockito.ArgumentMatchers.anyString())) + .thenThrow(failure); + JdbcClientPool pool = org.mockito.Mockito.mock(JdbcClientPool.class); + org.mockito.Mockito.doAnswer( + invocation -> + ((org.apache.paimon.client.ClientPool.Action< + ?, java.sql.Connection, ?>) + invocation.getArgument(0)) + .run(connection)) + .when(pool) + .run(org.mockito.ArgumentMatchers.any()); + MysqlDistributedLockDialect dialect = new MysqlDistributedLockDialect(); + + assertThatThrownBy(() -> dialect.lockAcquire(pool, "jdbc.broken.lock", 1000)) + .isInstanceOf(java.sql.SQLException.class) + .hasMessageContaining("connection reset"); + } + @Test public void testCleanTimeoutLockAndAcquireLock() throws SQLException, InterruptedException { String lockId = "jdbc.testDb.testTable";