Repository navigation
[Fix-18274][Registry] Handle expired JDBC heartbeat sessions #18416
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
qiuyanjun888
wants to merge
21
commits into
apache:dev
Choose a base branch
from
qiuyanjun888:Fix-18274-alt-20260715T021630Z
base: dev
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Open
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 2cf868b
[Fix-18274][Registry] Disconnect expired JDBC sessions immediately
qiuyanjun888 cbdfd1d
[Fix-18274][Registry] Preserve stopped state during heartbeat
qiuyanjun888 79cc4a0
[Fix-18274][Registry] Make heartbeat state transitions atomic
qiuyanjun888 eed6305
Merge branch 'dev' into Fix-18274-alt-20260715T021630Z
SbloodyS 93e0b2f
Merge branch 'dev' into Fix-18274-alt-20260715T021630Z
SbloodyS 7a87311
Merge remote-tracking branch 'upstream/dev' into Fix-18274-alt-202607…
qiuyanjun888 fddfbe1
[Fix-18274][Registry] Address JDBC registry state review
qiuyanjun888 71e8a14
[Fix-18274][Registry] Address JDBC registry review comments
qiuyanjun888 d7d0910
Merge branch 'dev' into Fix-18274-alt-20260715T021630Z
SbloodyS 21b4067
docs: add JDBC registry heartbeat state design
qiuyanjun888 1774519
docs: add JDBC registry heartbeat implementation plan
qiuyanjun888 ae7f23e
[Fix-18274][Registry] Make JDBC heartbeat state transitions safe
qiuyanjun888 fb69158
[Fix-18274][Registry] Stabilize JDBC heartbeat test baseline
qiuyanjun888 3c85344
docs: add JDBC registry concurrency design
qiuyanjun888 4ef9952
[Fix-18274][Registry] Fix JDBC registry concurrency races
qiuyanjun888 6e18b6d
Merge branch 'dev' into Fix-18274-alt-20260715T021630Z
qiuyanjun888 95c141a
[Fix-18274][Docs] Remove internal JDBC registry design documents
qiuyanjun888 0cb0955
[Fix-18274][Registry] Revert out-of-scope concurrency changes
qiuyanjun888 9a6b9e1
[Fix-18274][Registry] Handle JDBC heartbeat state races
qiuyanjun888 dc41136
[Fix-18274][Registry] Guard expired heartbeat cleanup against renewal
qiuyanjun888 File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Some comments aren't visible on the classic Files Changed page.
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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; | ||
|
|
@@ -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<>(); | ||
|
|
||
|
|
@@ -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, | ||
|
|
@@ -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; | ||
| } | ||
|
|
@@ -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, | ||
|
|
@@ -168,7 +172,7 @@ public void deregisterClient(IJdbcRegistryClient jdbcRegistryClient) { | |
|
|
||
| @Override | ||
| public JdbcRegistryServerState getServerState() { | ||
| return jdbcRegistryServerState; | ||
| return serverState.get(); | ||
| } | ||
|
|
||
| @Override | ||
|
|
@@ -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) | ||
|
|
@@ -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() | ||
|
|
@@ -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; | ||
| } | ||
|
|
@@ -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()); | ||
| 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
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Is all state synchronization intended to handle concurrency issues caused by |
||
| } | ||
| 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) { | ||
|
|
@@ -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) { | ||
|
|
||
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.