Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
0a31733
[Fix-18274][Registry] Handle expired JDBC heartbeat sessions
qiuyanjun888 Jul 15, 2026
2cf868b
[Fix-18274][Registry] Disconnect expired JDBC sessions immediately
qiuyanjun888 Jul 16, 2026
cbdfd1d
[Fix-18274][Registry] Preserve stopped state during heartbeat
qiuyanjun888 Jul 20, 2026
79cc4a0
[Fix-18274][Registry] Make heartbeat state transitions atomic
qiuyanjun888 Jul 30, 2026
eed6305
Merge branch 'dev' into Fix-18274-alt-20260715T021630Z
SbloodyS Sep 1, 2026
93e0b2f
Merge branch 'dev' into Fix-18274-alt-20260715T021630Z
SbloodyS Sep 2, 2026
7a87311
Merge remote-tracking branch 'upstream/dev' into Fix-18274-alt-202607…
qiuyanjun888 Sep 12, 2026
fddfbe1
[Fix-18274][Registry] Address JDBC registry state review
qiuyanjun888 Sep 12, 2026
71e8a14
[Fix-18274][Registry] Address JDBC registry review comments
qiuyanjun888 Sep 14, 2026
d7d0910
Merge branch 'dev' into Fix-18274-alt-20260715T021630Z
SbloodyS Sep 14, 2026
21b4067
docs: add JDBC registry heartbeat state design
qiuyanjun888 Sep 26, 2026
1774519
docs: add JDBC registry heartbeat implementation plan
qiuyanjun888 Sep 26, 2026
ae7f23e
[Fix-18274][Registry] Make JDBC heartbeat state transitions safe
qiuyanjun888 Sep 26, 2026
fb69158
[Fix-18274][Registry] Stabilize JDBC heartbeat test baseline
qiuyanjun888 Sep 26, 2026
3c85344
docs: add JDBC registry concurrency design
qiuyanjun888 Sep 26, 2026
4ef9952
[Fix-18274][Registry] Fix JDBC registry concurrency races
qiuyanjun888 Sep 26, 2026
6e18b6d
Merge branch 'dev' into Fix-18274-alt-20260715T021630Z
qiuyanjun888 Sep 27, 2026
95c141a
[Fix-18274][Docs] Remove internal JDBC registry design documents
qiuyanjun888 Sep 27, 2026
0cb0955
[Fix-18274][Registry] Revert out-of-scope concurrency changes
qiuyanjun888 Sep 27, 2026
9a6b9e1
[Fix-18274][Registry] Handle JDBC heartbeat state races
qiuyanjun888 Sep 27, 2026
dc41136
[Fix-18274][Registry] Guard expired heartbeat cleanup against renewal
qiuyanjun888 Sep 28, 2026
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 @@ -19,6 +19,8 @@

import org.apache.dolphinscheduler.plugin.registry.jdbc.model.DO.JdbcRegistryClientHeartbeat;

import org.apache.ibatis.annotations.Delete;
import org.apache.ibatis.annotations.Param;
import org.apache.ibatis.annotations.Select;

import java.util.List;
Expand All @@ -30,4 +32,8 @@ public interface JdbcRegistryClientHeartbeatMapper extends BaseMapper<JdbcRegist
@Select("select * from t_ds_jdbc_registry_client_heartbeat")
List<JdbcRegistryClientHeartbeat> selectAll();

@Delete("delete from t_ds_jdbc_registry_client_heartbeat "
+ "where id = #{id} and last_heartbeat_time = #{lastHeartbeatTime}")
int deleteByIdAndLastHeartbeatTime(@Param("id") Long id, @Param("lastHeartbeatTime") Long lastHeartbeatTime);

}
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,10 @@ public void deleteByIds(Collection<Long> clientIds) {
jdbcRegistryClientHeartbeatMapper.deleteBatchIds(clientIds);
}

public boolean deleteByIdAndLastHeartbeatTime(Long id, Long lastHeartbeatTime) {
return jdbcRegistryClientHeartbeatMapper.deleteByIdAndLastHeartbeatTime(id, lastHeartbeatTime) == 1;
}

public boolean updateById(JdbcRegistryClientHeartbeatDTO jdbcRegistryClientHeartbeatDTO) {
JdbcRegistryClientHeartbeat jdbcRegistryClientHeartbeat =
JdbcRegistryClientHeartbeatDTO.toJdbcRegistryClientHeartbeat(jdbcRegistryClientHeartbeatDTO);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.stream.Collectors;

import lombok.SneakyThrows;
Expand All @@ -70,7 +71,8 @@ public class JdbcRegistryServer implements IJdbcRegistryServer {

private final JdbcRegistryLockManager jdbcRegistryLockManager;

private JdbcRegistryServerState jdbcRegistryServerState;
private final AtomicReference<JdbcRegistryServerState> serverState =
new AtomicReference<>(JdbcRegistryServerState.INIT);

private final List<IJdbcRegistryClient> jdbcRegistryClients = new CopyOnWriteArrayList<>();

Expand All @@ -81,7 +83,7 @@ public class JdbcRegistryServer implements IJdbcRegistryServer {

private final ScheduledExecutorService schedulerThreadExecutor;

private Long lastSuccessHeartbeat;
private volatile long lastSuccessHeartbeat;

public JdbcRegistryServer(JdbcRegistryDataRepository jdbcRegistryDataRepository,
JdbcRegistryLockRepository jdbcRegistryLockRepository,
Expand All @@ -100,13 +102,12 @@ public JdbcRegistryServer(JdbcRegistryDataRepository jdbcRegistryDataRepository,
transactionTemplate, schedulerThreadExecutor);
this.jdbcRegistryLockManager = new JdbcRegistryLockManager(
jdbcRegistryProperties, jdbcRegistryLockRepository);
this.jdbcRegistryServerState = JdbcRegistryServerState.INIT;
lastSuccessHeartbeat = System.currentTimeMillis();
}

@Override
public void start() {
if (jdbcRegistryServerState != JdbcRegistryServerState.INIT) {
public synchronized void start() {
if (serverState.get() != JdbcRegistryServerState.INIT) {
// The server is already started or stopped, will not start again.
return;
}
Expand All @@ -119,7 +120,10 @@ public void start() {
jdbcRegistryProperties.getSessionTimeout().toMillis(),
TimeUnit.MILLISECONDS);
jdbcRegistryDataManager.start();
jdbcRegistryServerState = JdbcRegistryServerState.STARTED;
if (!serverState.compareAndSet(JdbcRegistryServerState.INIT, JdbcRegistryServerState.STARTED)) {
log.warn("The JdbcRegistryServer state changed before startup completed: {}", serverState.get());
return;
}
doTriggerOnConnectedListener();
schedulerThreadExecutor.scheduleWithFixedDelay(
this::refreshClientsHeartbeat,
Expand Down Expand Up @@ -168,7 +172,7 @@ public void deregisterClient(IJdbcRegistryClient jdbcRegistryClient) {

@Override
public JdbcRegistryServerState getServerState() {
return jdbcRegistryServerState;
return serverState.get();
}

@Override
Expand Down Expand Up @@ -256,7 +260,22 @@ public void releaseJdbcRegistryLock(Long clientId, String lockKey) {

@Override
public void close() {
jdbcRegistryServerState = JdbcRegistryServerState.STOPPED;
synchronized (this) {
while (true) {
JdbcRegistryServerState currentState = serverState.get();
if (currentState == JdbcRegistryServerState.STOPPED) {
log.warn("The JdbcRegistryServer is already STOPPED.");
return;
}
if (serverState.compareAndSet(currentState, JdbcRegistryServerState.STOPPED)) {
break;
}
// A heartbeat can change the state without this monitor. Do not drop the close request.
log.debug("Failed to stop JdbcRegistryServer from state {}, current state is {}, retrying",
currentState,
serverState.get());
}
}
schedulerThreadExecutor.shutdown();
List<Long> clientIds = jdbcRegistryClients.stream()
.map(IJdbcRegistryClient::getJdbcRegistryClientIdentify)
Expand All @@ -269,23 +288,26 @@ public void close() {

private void purgeInvalidJdbcRegistryMetadata() {
final StopWatch stopWatch = StopWatch.createStarted();
if (jdbcRegistryServerState == JdbcRegistryServerState.STOPPED) {
JdbcRegistryServerState currentState = getServerState();
if (currentState == JdbcRegistryServerState.STOPPED
|| currentState == JdbcRegistryServerState.DISCONNECTED) {
return;
}
// remove the client which is already dead from the registry, and remove it's related data and lock.
final List<JdbcRegistryClientHeartbeatDTO> jdbcRegistryClients = jdbcRegistryClientRepository.queryAll();
final Set<Long> deadJdbcRegistryClientIds = jdbcRegistryClients
final Set<Long> deletedJdbcRegistryClientIds = jdbcRegistryClients
.stream()
.filter(JdbcRegistryClientHeartbeatDTO::isDead)
.filter(jdbcRegistryClient -> jdbcRegistryClientRepository.deleteByIdAndLastHeartbeatTime(
jdbcRegistryClient.getId(), jdbcRegistryClient.getLastHeartbeatTime()))
.map(JdbcRegistryClientHeartbeatDTO::getId)
.collect(Collectors.toSet());
doPurgeJdbcRegistryClientInDB(deadJdbcRegistryClientIds);

// remove the data and lock which client is not exist.
final Set<Long> existJdbcRegistryClientIds = jdbcRegistryClients
.stream()
.map(JdbcRegistryClientHeartbeatDTO::getId)
.filter(id -> !deadJdbcRegistryClientIds.contains(id))
.filter(id -> !deletedJdbcRegistryClientIds.contains(id))
.collect(Collectors.toSet());
jdbcRegistryDataManager.getAllJdbcRegistryData()
.stream()
Expand Down Expand Up @@ -321,8 +343,11 @@ private void refreshClientsHeartbeat() {
if (CollectionUtils.isEmpty(jdbcRegistryClients)) {
return;
}
if (jdbcRegistryServerState == JdbcRegistryServerState.STOPPED) {
log.warn("The JdbcRegistryServer is STOPPED, will not refresh clients: {} heartbeat.",
JdbcRegistryServerState currentState = getServerState();
if (currentState == JdbcRegistryServerState.STOPPED
|| currentState == JdbcRegistryServerState.DISCONNECTED) {
log.warn("The JdbcRegistryServer is {}, will not refresh clients: {} heartbeat.",
currentState,
CollectionUtils.collect(jdbcRegistryClients, IJdbcRegistryClient::getJdbcRegistryClientIdentify));
return;
}
Expand All @@ -341,38 +366,69 @@ private void refreshClientsHeartbeat() {
}
JdbcRegistryClientHeartbeatDTO clone = jdbcRegistryClientHeartbeatDTO.clone();
clone.setLastHeartbeatTime(now);
jdbcRegistryClientRepository.updateById(jdbcRegistryClientHeartbeatDTO);
if (!jdbcRegistryClientRepository.updateById(clone)) {
log.error("The client heartbeat has expired: {}", jdbcRegistryClientHeartbeatDTO.getId());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
log.error("The client heartbeat has expired: {}", jdbcRegistryClientHeartbeatDTO.getId());
log.error("Refresh client: {} heartbeat failed, the client might has been deleted", jdbcRegistryClientHeartbeatDTO.getId());

throw new IllegalStateException(
"The client heartbeat record no longer exists: " + jdbcRegistryClientHeartbeatDTO.getId());
}
jdbcRegistryClientHeartbeatDTO.setLastHeartbeatTime(clone.getLastHeartbeatTime());
}
if (jdbcRegistryServerState == JdbcRegistryServerState.SUSPENDED) {
jdbcRegistryServerState = JdbcRegistryServerState.STARTED;
doTriggerReconnectedListener();
currentState = serverState.get();
boolean reconnected = currentState == JdbcRegistryServerState.SUSPENDED;
if (reconnected) {
if (!serverState.compareAndSet(JdbcRegistryServerState.SUSPENDED, JdbcRegistryServerState.STARTED)) {
log.debug("Failed to reconnect JdbcRegistryServer; current state is {}", serverState.get());
return;
}
} else if (currentState != JdbcRegistryServerState.STARTED) {
return;
}
// Serialize heartbeat side effects with close(), even if close wins after the state transition.
synchronized (this) {
if (serverState.get() != JdbcRegistryServerState.STARTED) {
return;
}
lastSuccessHeartbeat = now;
if (reconnected) {
doTriggerReconnectedListener();
}
Comment on lines +376 to +394

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is all state synchronization intended to handle concurrency issues caused by close? That seems overly complicated—shouldn’t we be able to ignore concurrency issues caused by close? Currently, it seems that only a service shutdown triggers a close, right?

}
lastSuccessHeartbeat = now;
log.debug("Success refresh clients: {} heartbeat.",
CollectionUtils.collect(jdbcRegistryClients, IJdbcRegistryClient::getJdbcRegistryClientIdentify));
} catch (Exception ex) {
log.error("Failed to refresh the client's term", ex);
switch (jdbcRegistryServerState) {
case STARTED:
jdbcRegistryServerState = JdbcRegistryServerState.SUSPENDED;
break;
case SUSPENDED:
if (System.currentTimeMillis() - lastSuccessHeartbeat > jdbcRegistryProperties.getSessionTimeout()
.toMillis()) {
jdbcRegistryServerState = JdbcRegistryServerState.DISCONNECTED;
currentState = serverState.get();
if (currentState != JdbcRegistryServerState.STARTED
&& currentState != JdbcRegistryServerState.SUSPENDED) {
return;
}
long sessionTimeoutMillis = jdbcRegistryProperties.getSessionTimeout().toMillis();
if (System.currentTimeMillis() - lastSuccessHeartbeat > sessionTimeoutMillis) {
if (!serverState.compareAndSet(currentState, JdbcRegistryServerState.DISCONNECTED)) {
log.debug("Failed to disconnect JdbcRegistryServer from state {}, current state is {}",
currentState,
serverState.get());
return;
}
synchronized (this) {
if (serverState.get() == JdbcRegistryServerState.DISCONNECTED) {
doTriggerOnDisConnectedListener();
}
break;
default:
break;
}
} else if (currentState == JdbcRegistryServerState.STARTED
&& !serverState.compareAndSet(JdbcRegistryServerState.STARTED, JdbcRegistryServerState.SUSPENDED)) {
log.debug("Failed to suspend JdbcRegistryServer; current state is {}", serverState.get());
return;
}
}
}

private void doTriggerReconnectedListener() {
log.info("Trigger:onReconnected listener.");
connectionStateListeners.forEach(listener -> {
if (serverState.get() != JdbcRegistryServerState.STARTED) {
return;
}
try {
listener.onReconnected();
} catch (Exception ex) {
Expand All @@ -395,6 +451,9 @@ private void doTriggerOnConnectedListener() {
private void doTriggerOnDisConnectedListener() {
log.info("Trigger:onDisConnected listener.");
connectionStateListeners.forEach(listener -> {
if (serverState.get() != JdbcRegistryServerState.DISCONNECTED) {
return;
}
try {
listener.onDisConnected();
} catch (Exception ex) {
Expand Down
Loading
Loading