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";