|
|
@ -29,11 +29,13 @@ import java.sql.ResultSet; |
|
|
|
import java.sql.SQLException; |
|
|
|
import java.sql.SQLException; |
|
|
|
import java.sql.SQLIntegrityConstraintViolationException; |
|
|
|
import java.sql.SQLIntegrityConstraintViolationException; |
|
|
|
import java.sql.Statement; |
|
|
|
import java.sql.Statement; |
|
|
|
import java.sql.Timestamp; |
|
|
|
|
|
|
|
import java.util.ArrayList; |
|
|
|
import java.util.ArrayList; |
|
|
|
import java.util.Collection; |
|
|
|
import java.util.Collection; |
|
|
|
|
|
|
|
import java.util.Iterator; |
|
|
|
import java.util.List; |
|
|
|
import java.util.List; |
|
|
|
import java.util.stream.Collectors; |
|
|
|
import java.util.Optional; |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
import lombok.NonNull; |
|
|
|
|
|
|
|
|
|
|
|
import org.slf4j.Logger; |
|
|
|
import org.slf4j.Logger; |
|
|
|
import org.slf4j.LoggerFactory; |
|
|
|
import org.slf4j.LoggerFactory; |
|
|
@ -72,7 +74,8 @@ public class MysqlOperator implements AutoCloseable { |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
public List<MysqlRegistryData> queryAllMysqlRegistryData() throws SQLException { |
|
|
|
public List<MysqlRegistryData> queryAllMysqlRegistryData() throws SQLException { |
|
|
|
String sql = "select id, `key`, data, type, create_time, last_update_time from t_ds_mysql_registry_data"; |
|
|
|
String sql = |
|
|
|
|
|
|
|
"select id, `key`, data, type, last_term, create_time, last_update_time from t_ds_mysql_registry_data"; |
|
|
|
try ( |
|
|
|
try ( |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sql); |
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sql); |
|
|
@ -84,6 +87,7 @@ public class MysqlOperator implements AutoCloseable { |
|
|
|
.key(resultSet.getString("key")) |
|
|
|
.key(resultSet.getString("key")) |
|
|
|
.data(resultSet.getString("data")) |
|
|
|
.data(resultSet.getString("data")) |
|
|
|
.type(resultSet.getInt("type")) |
|
|
|
.type(resultSet.getInt("type")) |
|
|
|
|
|
|
|
.lastTerm(resultSet.getLong("last_term")) |
|
|
|
.createTime(resultSet.getTimestamp("create_time")) |
|
|
|
.createTime(resultSet.getTimestamp("create_time")) |
|
|
|
.lastUpdateTime(resultSet.getTimestamp("last_update_time")) |
|
|
|
.lastUpdateTime(resultSet.getTimestamp("last_update_time")) |
|
|
|
.build(); |
|
|
|
.build(); |
|
|
@ -93,24 +97,75 @@ public class MysqlOperator implements AutoCloseable { |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
public long insertOrUpdateEphemeralData(String key, String value) throws SQLException { |
|
|
|
public Long insertOrUpdateEphemeralData(String key, String value) throws SQLException { |
|
|
|
|
|
|
|
Optional<MysqlRegistryData> mysqlRegistryDataOptional = selectByKey(key); |
|
|
|
|
|
|
|
if (mysqlRegistryDataOptional.isPresent()) { |
|
|
|
|
|
|
|
long id = mysqlRegistryDataOptional.get().getId(); |
|
|
|
|
|
|
|
if (!updateValueById(id, value)) { |
|
|
|
|
|
|
|
throw new SQLException(String.format("update registry value failed, key: %s, value: %s", key, value)); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
return id; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
MysqlRegistryData mysqlRegistryData = MysqlRegistryData.builder() |
|
|
|
|
|
|
|
.key(key) |
|
|
|
|
|
|
|
.data(value) |
|
|
|
|
|
|
|
.type(DataType.EPHEMERAL.getTypeValue()) |
|
|
|
|
|
|
|
.lastTerm(System.currentTimeMillis()) |
|
|
|
|
|
|
|
.build(); |
|
|
|
|
|
|
|
return insertMysqlRegistryData(mysqlRegistryData); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
private Optional<MysqlRegistryData> selectByKey(@NonNull String key) throws SQLException { |
|
|
|
String sql = |
|
|
|
String sql = |
|
|
|
"INSERT INTO t_ds_mysql_registry_data (`key`, data, type, create_time, last_update_time) VALUES (?, ?, ?, current_timestamp, current_timestamp)" |
|
|
|
"select id, `key`, data, type, create_time, last_update_time from t_ds_mysql_registry_data where `key` = ?"; |
|
|
|
+ |
|
|
|
try ( |
|
|
|
"ON DUPLICATE KEY UPDATE data=?, last_update_time=current_timestamp"; |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
// put a ephemeralData
|
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sql)) { |
|
|
|
|
|
|
|
preparedStatement.setString(1, key); |
|
|
|
|
|
|
|
try (ResultSet resultSet = preparedStatement.executeQuery()) { |
|
|
|
|
|
|
|
if (resultSet.next()) { |
|
|
|
|
|
|
|
return Optional.of( |
|
|
|
|
|
|
|
MysqlRegistryData.builder() |
|
|
|
|
|
|
|
.id(resultSet.getLong("id")) |
|
|
|
|
|
|
|
.key(resultSet.getString("key")) |
|
|
|
|
|
|
|
.data(resultSet.getString("data")) |
|
|
|
|
|
|
|
.type(resultSet.getInt("type")) |
|
|
|
|
|
|
|
.createTime(resultSet.getTimestamp("create_time")) |
|
|
|
|
|
|
|
.lastUpdateTime(resultSet.getTimestamp("last_update_time")) |
|
|
|
|
|
|
|
.build()); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
return Optional.empty(); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
private boolean updateValueById(long id, String value) throws SQLException { |
|
|
|
|
|
|
|
String sql = "update t_ds_mysql_registry_data set data = ?, last_term = ? where id = ?"; |
|
|
|
|
|
|
|
try ( |
|
|
|
|
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
|
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sql)) { |
|
|
|
|
|
|
|
preparedStatement.setString(1, value); |
|
|
|
|
|
|
|
preparedStatement.setLong(2, System.currentTimeMillis()); |
|
|
|
|
|
|
|
preparedStatement.setLong(3, id); |
|
|
|
|
|
|
|
return preparedStatement.executeUpdate() > 0; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
private long insertMysqlRegistryData(@NonNull MysqlRegistryData mysqlRegistryData) throws SQLException { |
|
|
|
|
|
|
|
String sql = |
|
|
|
|
|
|
|
"INSERT INTO t_ds_mysql_registry_data (`key`, data, type, last_term) VALUES (?, ?, ?, ?)"; |
|
|
|
try ( |
|
|
|
try ( |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
PreparedStatement preparedStatement = |
|
|
|
PreparedStatement preparedStatement = |
|
|
|
connection.prepareStatement(sql, Statement.RETURN_GENERATED_KEYS)) { |
|
|
|
connection.prepareStatement(sql, Statement.RETURN_GENERATED_KEYS)) { |
|
|
|
preparedStatement.setString(1, key); |
|
|
|
preparedStatement.setString(1, mysqlRegistryData.getKey()); |
|
|
|
preparedStatement.setString(2, value); |
|
|
|
preparedStatement.setString(2, mysqlRegistryData.getData()); |
|
|
|
preparedStatement.setInt(3, DataType.EPHEMERAL.getTypeValue()); |
|
|
|
preparedStatement.setInt(3, mysqlRegistryData.getType()); |
|
|
|
preparedStatement.setString(4, value); |
|
|
|
preparedStatement.setLong(4, mysqlRegistryData.getLastTerm()); |
|
|
|
int insertCount = preparedStatement.executeUpdate(); |
|
|
|
int insertCount = preparedStatement.executeUpdate(); |
|
|
|
ResultSet generatedKeys = preparedStatement.getGeneratedKeys(); |
|
|
|
ResultSet generatedKeys = preparedStatement.getGeneratedKeys(); |
|
|
|
if (insertCount < 1 || !generatedKeys.next()) { |
|
|
|
if (insertCount < 1 || !generatedKeys.next()) { |
|
|
|
throw new SQLException("Insert or update ephemeral data error"); |
|
|
|
throw new SQLException("Insert ephemeral data error, data: " + mysqlRegistryData); |
|
|
|
} |
|
|
|
} |
|
|
|
return generatedKeys.getLong(1); |
|
|
|
return generatedKeys.getLong(1); |
|
|
|
} |
|
|
|
} |
|
|
@ -118,18 +173,21 @@ public class MysqlOperator implements AutoCloseable { |
|
|
|
|
|
|
|
|
|
|
|
public long insertOrUpdatePersistentData(String key, String value) throws SQLException { |
|
|
|
public long insertOrUpdatePersistentData(String key, String value) throws SQLException { |
|
|
|
String sql = |
|
|
|
String sql = |
|
|
|
"INSERT INTO t_ds_mysql_registry_data (`key`, data, type, create_time, last_update_time) VALUES (?, ?, ?, current_timestamp, current_timestamp)" |
|
|
|
"INSERT INTO t_ds_mysql_registry_data (`key`, data, type, last_term) VALUES (?, ?, ?, ?)" |
|
|
|
+ |
|
|
|
+ |
|
|
|
"ON DUPLICATE KEY UPDATE data=?, last_update_time=current_timestamp"; |
|
|
|
"ON DUPLICATE KEY UPDATE data=?, last_term=?"; |
|
|
|
// put a persistent Data
|
|
|
|
// put a persistent Data
|
|
|
|
try ( |
|
|
|
try ( |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
PreparedStatement preparedStatement = |
|
|
|
PreparedStatement preparedStatement = |
|
|
|
connection.prepareStatement(sql, Statement.RETURN_GENERATED_KEYS)) { |
|
|
|
connection.prepareStatement(sql, Statement.RETURN_GENERATED_KEYS)) { |
|
|
|
|
|
|
|
long term = System.currentTimeMillis(); |
|
|
|
preparedStatement.setString(1, key); |
|
|
|
preparedStatement.setString(1, key); |
|
|
|
preparedStatement.setString(2, value); |
|
|
|
preparedStatement.setString(2, value); |
|
|
|
preparedStatement.setInt(3, DataType.PERSISTENT.getTypeValue()); |
|
|
|
preparedStatement.setInt(3, DataType.PERSISTENT.getTypeValue()); |
|
|
|
preparedStatement.setString(4, value); |
|
|
|
preparedStatement.setLong(4, term); |
|
|
|
|
|
|
|
preparedStatement.setString(5, value); |
|
|
|
|
|
|
|
preparedStatement.setLong(6, term); |
|
|
|
int insertCount = preparedStatement.executeUpdate(); |
|
|
|
int insertCount = preparedStatement.executeUpdate(); |
|
|
|
ResultSet generatedKeys = preparedStatement.getGeneratedKeys(); |
|
|
|
ResultSet generatedKeys = preparedStatement.getGeneratedKeys(); |
|
|
|
if (insertCount < 1 || !generatedKeys.next()) { |
|
|
|
if (insertCount < 1 || !generatedKeys.next()) { |
|
|
@ -176,8 +234,7 @@ public class MysqlOperator implements AutoCloseable { |
|
|
|
try ( |
|
|
|
try ( |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sql)) { |
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sql)) { |
|
|
|
preparedStatement.setTimestamp(1, |
|
|
|
preparedStatement.setLong(1, System.currentTimeMillis() - expireTimeWindow); |
|
|
|
new Timestamp(System.currentTimeMillis() - expireTimeWindow)); |
|
|
|
|
|
|
|
int i = preparedStatement.executeUpdate(); |
|
|
|
int i = preparedStatement.executeUpdate(); |
|
|
|
if (i > 0) { |
|
|
|
if (i > 0) { |
|
|
|
logger.info("Clear expire lock, size: {}", i); |
|
|
|
logger.info("Clear expire lock, size: {}", i); |
|
|
@ -188,11 +245,11 @@ public class MysqlOperator implements AutoCloseable { |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
public void clearExpireEphemeralDate() { |
|
|
|
public void clearExpireEphemeralDate() { |
|
|
|
String sql = "delete from t_ds_mysql_registry_data where last_update_time < ? and type = ?"; |
|
|
|
String sql = "delete from t_ds_mysql_registry_data where last_term < ? and type = ?"; |
|
|
|
try ( |
|
|
|
try ( |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sql)) { |
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sql)) { |
|
|
|
preparedStatement.setTimestamp(1, new Timestamp(System.currentTimeMillis() - expireTimeWindow)); |
|
|
|
preparedStatement.setLong(1, System.currentTimeMillis() - expireTimeWindow); |
|
|
|
preparedStatement.setInt(2, DataType.EPHEMERAL.getTypeValue()); |
|
|
|
preparedStatement.setInt(2, DataType.EPHEMERAL.getTypeValue()); |
|
|
|
int i = preparedStatement.executeUpdate(); |
|
|
|
int i = preparedStatement.executeUpdate(); |
|
|
|
if (i > 0) { |
|
|
|
if (i > 0) { |
|
|
@ -205,7 +262,7 @@ public class MysqlOperator implements AutoCloseable { |
|
|
|
|
|
|
|
|
|
|
|
public MysqlRegistryData getData(String key) throws SQLException { |
|
|
|
public MysqlRegistryData getData(String key) throws SQLException { |
|
|
|
String sql = |
|
|
|
String sql = |
|
|
|
"SELECT id, `key`, data, type, create_time, last_update_time FROM t_ds_mysql_registry_data WHERE `key` = ?"; |
|
|
|
"SELECT id, `key`, data, type, last_term, create_time, last_update_time FROM t_ds_mysql_registry_data WHERE `key` = ?"; |
|
|
|
try ( |
|
|
|
try ( |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sql)) { |
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sql)) { |
|
|
@ -219,6 +276,7 @@ public class MysqlOperator implements AutoCloseable { |
|
|
|
.key(resultSet.getString("key")) |
|
|
|
.key(resultSet.getString("key")) |
|
|
|
.data(resultSet.getString("data")) |
|
|
|
.data(resultSet.getString("data")) |
|
|
|
.type(resultSet.getInt("type")) |
|
|
|
.type(resultSet.getInt("type")) |
|
|
|
|
|
|
|
.lastTerm(resultSet.getLong("last_term")) |
|
|
|
.createTime(resultSet.getTimestamp("create_time")) |
|
|
|
.createTime(resultSet.getTimestamp("create_time")) |
|
|
|
.lastUpdateTime(resultSet.getTimestamp("last_update_time")) |
|
|
|
.lastUpdateTime(resultSet.getTimestamp("last_update_time")) |
|
|
|
.build(); |
|
|
|
.build(); |
|
|
@ -265,13 +323,14 @@ public class MysqlOperator implements AutoCloseable { |
|
|
|
*/ |
|
|
|
*/ |
|
|
|
public MysqlRegistryLock tryToAcquireLock(String key) throws SQLException { |
|
|
|
public MysqlRegistryLock tryToAcquireLock(String key) throws SQLException { |
|
|
|
String sql = |
|
|
|
String sql = |
|
|
|
"INSERT INTO t_ds_mysql_registry_lock (`key`, lock_owner, last_term, last_update_time, create_time) VALUES (?, ?, current_timestamp, current_timestamp, current_timestamp)"; |
|
|
|
"INSERT INTO t_ds_mysql_registry_lock (`key`, lock_owner, last_term) VALUES (?, ?, ?)"; |
|
|
|
try ( |
|
|
|
try ( |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
PreparedStatement preparedStatement = |
|
|
|
PreparedStatement preparedStatement = |
|
|
|
connection.prepareStatement(sql, Statement.RETURN_GENERATED_KEYS)) { |
|
|
|
connection.prepareStatement(sql, Statement.RETURN_GENERATED_KEYS)) { |
|
|
|
preparedStatement.setString(1, key); |
|
|
|
preparedStatement.setString(1, key); |
|
|
|
preparedStatement.setString(2, MysqlRegistryConstant.LOCK_OWNER); |
|
|
|
preparedStatement.setString(2, MysqlRegistryConstant.LOCK_OWNER); |
|
|
|
|
|
|
|
preparedStatement.setLong(3, System.currentTimeMillis()); |
|
|
|
preparedStatement.executeUpdate(); |
|
|
|
preparedStatement.executeUpdate(); |
|
|
|
try (ResultSet resultSet = preparedStatement.getGeneratedKeys()) { |
|
|
|
try (ResultSet resultSet = preparedStatement.getGeneratedKeys()) { |
|
|
|
if (resultSet.next()) { |
|
|
|
if (resultSet.next()) { |
|
|
@ -299,7 +358,7 @@ public class MysqlOperator implements AutoCloseable { |
|
|
|
.id(resultSet.getLong("id")) |
|
|
|
.id(resultSet.getLong("id")) |
|
|
|
.key(resultSet.getString("key")) |
|
|
|
.key(resultSet.getString("key")) |
|
|
|
.lockOwner(resultSet.getString("lock_owner")) |
|
|
|
.lockOwner(resultSet.getString("lock_owner")) |
|
|
|
.lastTerm(resultSet.getTimestamp("last_term")) |
|
|
|
.lastTerm(resultSet.getLong("last_term")) |
|
|
|
.lastUpdateTime(resultSet.getTimestamp("last_update_time")) |
|
|
|
.lastUpdateTime(resultSet.getTimestamp("last_update_time")) |
|
|
|
.createTime(resultSet.getTimestamp("create_time")) |
|
|
|
.createTime(resultSet.getTimestamp("create_time")) |
|
|
|
.build(); |
|
|
|
.build(); |
|
|
@ -322,24 +381,38 @@ public class MysqlOperator implements AutoCloseable { |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
public boolean updateEphemeralDataTerm(Collection<Long> ephemeralDateIds) throws SQLException { |
|
|
|
public boolean updateEphemeralDataTerm(Collection<Long> ephemeralDateIds) throws SQLException { |
|
|
|
String sql = "update t_ds_mysql_registry_data set `last_update_time` = current_timestamp() where `id` IN (?)"; |
|
|
|
StringBuilder sb = new StringBuilder("update t_ds_mysql_registry_data set `last_term` = ? where `id` IN ("); |
|
|
|
String ids = ephemeralDateIds.stream().map(String::valueOf).collect(Collectors.joining(",")); |
|
|
|
Iterator<Long> iterator = ephemeralDateIds.iterator(); |
|
|
|
|
|
|
|
for (int i = 0; i < ephemeralDateIds.size(); i++) { |
|
|
|
|
|
|
|
sb.append(iterator.next()); |
|
|
|
|
|
|
|
if (i != ephemeralDateIds.size() - 1) { |
|
|
|
|
|
|
|
sb.append(","); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
sb.append(")"); |
|
|
|
try ( |
|
|
|
try ( |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sql)) { |
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sb.toString())) { |
|
|
|
preparedStatement.setString(1, ids); |
|
|
|
preparedStatement.setLong(1, System.currentTimeMillis()); |
|
|
|
return preparedStatement.executeUpdate() > 0; |
|
|
|
return preparedStatement.executeUpdate() > 0; |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
public boolean updateLockTerm(List<Long> lockIds) throws SQLException { |
|
|
|
public boolean updateLockTerm(List<Long> lockIds) throws SQLException { |
|
|
|
String sql = |
|
|
|
StringBuilder sb = |
|
|
|
"update t_ds_mysql_registry_lock set `last_term` = current_timestamp and `last_update_time` = current_timestamp where `id` IN (?)"; |
|
|
|
new StringBuilder("update t_ds_mysql_registry_lock set `last_term` = ? where `id` IN ("); |
|
|
|
String ids = lockIds.stream().map(String::valueOf).collect(Collectors.joining(",")); |
|
|
|
Iterator<Long> iterator = lockIds.iterator(); |
|
|
|
|
|
|
|
for (int i = 0; i < lockIds.size(); i++) { |
|
|
|
|
|
|
|
sb.append(iterator.next()); |
|
|
|
|
|
|
|
if (i != lockIds.size() - 1) { |
|
|
|
|
|
|
|
sb.append(","); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
sb.append(")"); |
|
|
|
try ( |
|
|
|
try ( |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
Connection connection = dataSource.getConnection(); |
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sql)) { |
|
|
|
PreparedStatement preparedStatement = connection.prepareStatement(sb.toString())) { |
|
|
|
preparedStatement.setString(1, ids); |
|
|
|
preparedStatement.setLong(1, System.currentTimeMillis()); |
|
|
|
return preparedStatement.executeUpdate() > 0; |
|
|
|
return preparedStatement.executeUpdate() > 0; |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|