Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
});
}
Expand Down Expand Up @@ -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"));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down
Loading