This is an automated email from the ASF dual-hosted git repository.
sandynz pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shardingsphere-elasticjob.git
The following commit(s) were added to refs/heads/master by this push:
new 8ad78c684 Use direct reads for lifecycle server status (#2530)
8ad78c684 is described below
commit 8ad78c68410e96f806ef4b1aa56641d245fd8f9e
Author: Hongsheng Zhong <[email protected]>
AuthorDate: Wed Jul 29 09:55:42 2026 +0800
Use direct reads for lifecycle server status (#2530)
What: Read mutable server status directly in job and server statistics
APIs, with regression coverage for stale cached values.
Why: Lifecycle management queries should reflect authoritative registry
state after enable or disable operations instead of an asynchronously updated
local cache.
---
.../lifecycle/internal/statistics/JobStatisticsAPIImpl.java | 4 ++--
.../internal/statistics/ServerStatisticsAPIImpl.java | 2 +-
.../internal/statistics/JobStatisticsAPIImplTest.java | 11 +++++++----
.../internal/statistics/ServerStatisticsAPIImplTest.java | 13 +++++++++----
4 files changed, 19 insertions(+), 11 deletions(-)
diff --git
a/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/JobStatisticsAPIImpl.java
b/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/JobStatisticsAPIImpl.java
index 80200edab..7c82fb517 100644
---
a/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/JobStatisticsAPIImpl.java
+++
b/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/JobStatisticsAPIImpl.java
@@ -98,7 +98,7 @@ public final class JobStatisticsAPIImpl implements
JobStatisticsAPI {
List<String> serversPath =
regCenter.getChildrenKeys(jobNodePath.getServerNodePath());
int disabledServerCount = 0;
for (String each : serversPath) {
- if
(JobBriefInfo.JobStatus.DISABLED.name().equals(regCenter.get(jobNodePath.getServerNodePath(each))))
{
+ if
(JobBriefInfo.JobStatus.DISABLED.name().equals(regCenter.getDirectly(jobNodePath.getServerNodePath(each))))
{
disabledServerCount++;
}
}
@@ -147,7 +147,7 @@ public final class JobStatisticsAPIImpl implements
JobStatisticsAPI {
private JobBriefInfo.JobStatus getJobStatusByJobNameAndIp(final String
jobName, final String ip) {
JobNodePath jobNodePath = new JobNodePath(jobName);
- String status = regCenter.get(jobNodePath.getServerNodePath(ip));
+ String status =
regCenter.getDirectly(jobNodePath.getServerNodePath(ip));
if ("DISABLED".equalsIgnoreCase(status)) {
return JobBriefInfo.JobStatus.DISABLED;
} else {
diff --git
a/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ServerStatisticsAPIImpl.java
b/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ServerStatisticsAPIImpl.java
index 954664eab..a1da756ba 100644
---
a/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ServerStatisticsAPIImpl.java
+++
b/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ServerStatisticsAPIImpl.java
@@ -59,7 +59,7 @@ public final class ServerStatisticsAPIImpl implements
ServerStatisticsAPI {
for (String each :
regCenter.getChildrenKeys(jobNodePath.getServerNodePath())) {
servers.putIfAbsent(each, new ServerBriefInfo(each));
ServerBriefInfo serverInfo = servers.get(each);
- if
("DISABLED".equalsIgnoreCase(regCenter.get(jobNodePath.getServerNodePath(each))))
{
+ if
("DISABLED".equalsIgnoreCase(regCenter.getDirectly(jobNodePath.getServerNodePath(each))))
{
serverInfo.getDisabledJobsNum().incrementAndGet();
}
serverInfo.getJobNames().add(jobName);
diff --git
a/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/JobStatisticsAPIImplTest.java
b/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/JobStatisticsAPIImplTest.java
index e46763e20..8099ac923 100644
---
a/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/JobStatisticsAPIImplTest.java
+++
b/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/JobStatisticsAPIImplTest.java
@@ -77,7 +77,7 @@ class JobStatisticsAPIImplTest {
void assertGetOKJobBriefInfoWithPartialDisabledServer() {
when(regCenter.getDirectly("/test_job/config")).thenReturn(LifecycleYamlConstants.getSimpleJobYaml("test_job",
"desc"));
when(regCenter.getChildrenKeys("/test_job/servers")).thenReturn(Arrays.asList("ip1",
"ip2"));
- when(regCenter.get("/test_job/servers/ip1")).thenReturn("DISABLED");
+
when(regCenter.getDirectly("/test_job/servers/ip1")).thenReturn("DISABLED");
when(regCenter.getChildrenKeys("/test_job/instances")).thenReturn(Arrays.asList("ip1@-@defaultInstance",
"ip2@-@defaultInstance"));
when(regCenter.getChildrenKeys("/test_job/sharding")).thenReturn(Arrays.asList("0",
"1"));
when(regCenter.get("/test_job/sharding/0/instance")).thenReturn("ip1@-@defaultInstance");
@@ -90,8 +90,10 @@ class JobStatisticsAPIImplTest {
void assertGetDisabledJobBriefInfo() {
when(regCenter.getDirectly("/test_job/config")).thenReturn(LifecycleYamlConstants.getSimpleJobYaml("test_job",
"desc"));
when(regCenter.getChildrenKeys("/test_job/servers")).thenReturn(Arrays.asList("ip1",
"ip2"));
- when(regCenter.get("/test_job/servers/ip1")).thenReturn("DISABLED");
- when(regCenter.get("/test_job/servers/ip2")).thenReturn("DISABLED");
+ when(regCenter.get("/test_job/servers/ip1")).thenReturn("ENABLED");
+ when(regCenter.get("/test_job/servers/ip2")).thenReturn("ENABLED");
+
when(regCenter.getDirectly("/test_job/servers/ip1")).thenReturn("DISABLED");
+
when(regCenter.getDirectly("/test_job/servers/ip2")).thenReturn("DISABLED");
when(regCenter.getChildrenKeys("/test_job/instances")).thenReturn(Arrays.asList("ip1@-@defaultInstance",
"ip2@-@defaultInstance"));
JobBriefInfo jobBrief = jobStatisticsAPI.getJobBriefInfo("test_job");
assertThat(jobBrief.getStatus(), is(JobBriefInfo.JobStatus.DISABLED));
@@ -155,7 +157,8 @@ class JobStatisticsAPIImplTest {
when(regCenter.getChildrenKeys("/")).thenReturn(Arrays.asList("test_job_1",
"test_job_2", "test_job_3"));
when(regCenter.isExisted("/test_job_1/servers/ip1")).thenReturn(true);
when(regCenter.isExisted("/test_job_2/servers/ip1")).thenReturn(true);
- when(regCenter.get("/test_job_2/servers/ip1")).thenReturn("DISABLED");
+ when(regCenter.get("/test_job_2/servers/ip1")).thenReturn("ENABLED");
+
when(regCenter.getDirectly("/test_job_2/servers/ip1")).thenReturn("DISABLED");
when(regCenter.getChildrenKeys("/test_job_1/instances")).thenReturn(Collections.singletonList("ip1@-@defaultInstance"));
when(regCenter.get("/test_job_1/instances/ip1@-@defaultInstance")).thenReturn("jobInstanceId:
ip1@-@defaultInstance\nserverIp: ip1\n");
int i = 0;
diff --git
a/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ServerStatisticsAPIImplTest.java
b/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ServerStatisticsAPIImplTest.java
index f1e71fe25..1d7b40a42 100644
---
a/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ServerStatisticsAPIImplTest.java
+++
b/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ServerStatisticsAPIImplTest.java
@@ -32,6 +32,7 @@ import java.util.Collections;
import static org.hamcrest.CoreMatchers.is;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.junit.jupiter.api.Assertions.fail;
+import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
@@ -60,12 +61,16 @@ class ServerStatisticsAPIImplTest {
when(regCenter.getChildrenKeys("/")).thenReturn(Arrays.asList("test_job1",
"test_job2"));
when(regCenter.getChildrenKeys("/test_job1/servers")).thenReturn(Arrays.asList("ip1",
"ip2"));
when(regCenter.getChildrenKeys("/test_job2/servers")).thenReturn(Arrays.asList("ip1",
"ip2"));
- when(regCenter.get("/test_job1/servers/ip1")).thenReturn("DISABLED");
- when(regCenter.get("/test_job1/servers/ip2")).thenReturn("");
+
lenient().when(regCenter.get("/test_job1/servers/ip1")).thenReturn("ENABLED");
+
lenient().when(regCenter.get("/test_job1/servers/ip2")).thenReturn("ENABLED");
+
lenient().when(regCenter.get("/test_job2/servers/ip1")).thenReturn("ENABLED");
+
lenient().when(regCenter.get("/test_job2/servers/ip2")).thenReturn("ENABLED");
+
when(regCenter.getDirectly("/test_job1/servers/ip1")).thenReturn("DISABLED");
+ when(regCenter.getDirectly("/test_job1/servers/ip2")).thenReturn("");
when(regCenter.getChildrenKeys("/test_job1/instances")).thenReturn(Collections.singletonList("ip1@-@defaultInstance"));
- when(regCenter.get("/test_job2/servers/ip1")).thenReturn("DISABLED");
- when(regCenter.get("/test_job2/servers/ip2")).thenReturn("DISABLED");
+
when(regCenter.getDirectly("/test_job2/servers/ip1")).thenReturn("DISABLED");
+
when(regCenter.getDirectly("/test_job2/servers/ip2")).thenReturn("DISABLED");
when(regCenter.get("/test_job1/instances/ip1@-@defaultInstance")).thenReturn("jobInstanceId:
ip1@-@defaultInstance\nserverIp: ip1\n");
when(regCenter.get("/test_job2/instances/ip1@-@defaultInstance")).thenReturn("jobInstanceId:
ip1@-@defaultInstance\nserverIp: ip1\n");
when(regCenter.get("/test_job2/instances/ip2@-@defaultInstance2")).thenReturn("jobInstanceId:
ip2@-@defaultInstance2\nserverIp: ip2\n");