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 d8d6f3e5f Use direct reads for lifecycle sharding assignments (#2532)
d8d6f3e5f is described below
commit d8d6f3e5fe6a32d547a50c6310eb58a9919df6b4
Author: Hongsheng Zhong <[email protected]>
AuthorDate: Wed Jul 29 10:44:26 2026 +0800
Use direct reads for lifecycle sharding assignments (#2532)
What: Read sharding assignment values directly in lifecycle statistics and
cover stale cached values in focused tests.
Why: Lifecycle management queries should follow authoritative assignments
after resharding instead of an asynchronously updated local cache.
---
.../internal/statistics/JobStatisticsAPIImpl.java | 2 +-
.../statistics/ShardingStatisticsAPIImpl.java | 2 +-
.../statistics/JobStatisticsAPIImplTest.java | 25 +++++++++++-----------
.../statistics/ShardingStatisticsAPIImplTest.java | 7 +++++-
4 files changed, 21 insertions(+), 15 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 7c82fb517..3d42af8b8 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
@@ -108,7 +108,7 @@ public final class JobStatisticsAPIImpl implements
JobStatisticsAPI {
private boolean isHasShardingFlag(final JobNodePath jobNodePath, final
List<String> instances) {
Set<String> shardingInstances = new HashSet<>();
for (String each :
regCenter.getChildrenKeys(jobNodePath.getShardingNodePath())) {
- String instanceId =
regCenter.get(jobNodePath.getShardingNodePath(each, "instance"));
+ String instanceId =
regCenter.getDirectly(jobNodePath.getShardingNodePath(each, "instance"));
if (null != instanceId && !instanceId.isEmpty()) {
shardingInstances.add(instanceId);
}
diff --git
a/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ShardingStatisticsAPIImpl.java
b/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ShardingStatisticsAPIImpl.java
index 678343b3f..071263736 100644
---
a/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ShardingStatisticsAPIImpl.java
+++
b/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ShardingStatisticsAPIImpl.java
@@ -54,7 +54,7 @@ public final class ShardingStatisticsAPIImpl implements
ShardingStatisticsAPI {
ShardingInfo result = new ShardingInfo();
result.setItem(Integer.parseInt(item));
JobNodePath jobNodePath = new JobNodePath(jobName);
- String instanceId =
regCenter.get(jobNodePath.getShardingNodePath(item, "instance"));
+ String instanceId =
regCenter.getDirectly(jobNodePath.getShardingNodePath(item, "instance"));
boolean disabled =
regCenter.isExisted(jobNodePath.getShardingNodePath(item, "disabled"));
boolean running =
regCenter.isExisted(jobNodePath.getShardingNodePath(item, "running"));
boolean shardingError =
!regCenter.isExisted(jobNodePath.getInstanceNodePath(instanceId));
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 8099ac923..a736d3517 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
@@ -60,9 +60,10 @@ class JobStatisticsAPIImplTest {
when(regCenter.getChildrenKeys("/test_job/servers")).thenReturn(Arrays.asList("ip1",
"ip2"));
when(regCenter.getChildrenKeys("/test_job/instances")).thenReturn(Arrays.asList("ip1@-@defaultInstance",
"ip2@-@defaultInstance"));
when(regCenter.getChildrenKeys("/test_job/sharding")).thenReturn(Arrays.asList("0",
"1", "2"));
-
when(regCenter.get("/test_job/sharding/0/instance")).thenReturn("ip1@-@defaultInstance");
-
when(regCenter.get("/test_job/sharding/1/instance")).thenReturn("ip1@-@defaultInstance");
-
when(regCenter.get("/test_job/sharding/2/instance")).thenReturn("ip2@-@defaultInstance");
+
when(regCenter.get("/test_job/sharding/0/instance")).thenReturn("ip3@-@oldInstance");
+
when(regCenter.getDirectly("/test_job/sharding/0/instance")).thenReturn("ip1@-@defaultInstance");
+
when(regCenter.getDirectly("/test_job/sharding/1/instance")).thenReturn("ip1@-@defaultInstance");
+
when(regCenter.getDirectly("/test_job/sharding/2/instance")).thenReturn("ip2@-@defaultInstance");
when(regCenter.getChildrenKeys("/test_job/instances")).thenReturn(Arrays.asList("ip1@-@defaultInstance",
"ip2@-@defaultInstance"));
JobBriefInfo jobBrief = jobStatisticsAPI.getJobBriefInfo("test_job");
assertThat(jobBrief.getJobName(), is("test_job"));
@@ -80,8 +81,8 @@ class JobStatisticsAPIImplTest {
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");
-
when(regCenter.get("/test_job/sharding/1/instance")).thenReturn("ip2@-@defaultInstance");
+
when(regCenter.getDirectly("/test_job/sharding/0/instance")).thenReturn("ip1@-@defaultInstance");
+
when(regCenter.getDirectly("/test_job/sharding/1/instance")).thenReturn("ip2@-@defaultInstance");
JobBriefInfo jobBrief = jobStatisticsAPI.getJobBriefInfo("test_job");
assertThat(jobBrief.getStatus(), is(JobBriefInfo.JobStatus.OK));
}
@@ -105,9 +106,9 @@ class JobStatisticsAPIImplTest {
when(regCenter.getChildrenKeys("/test_job/servers")).thenReturn(Arrays.asList("ip1",
"ip2"));
when(regCenter.getChildrenKeys("/test_job/instances")).thenReturn(Arrays.asList("ip1@-@defaultInstance",
"ip2@-@defaultInstance"));
when(regCenter.getChildrenKeys("/test_job/sharding")).thenReturn(Arrays.asList("0",
"1", "2"));
-
when(regCenter.get("/test_job/sharding/0/instance")).thenReturn("ip1@-@defaultInstance");
-
when(regCenter.get("/test_job/sharding/1/instance")).thenReturn("ip2@-@defaultInstance");
-
when(regCenter.get("/test_job/sharding/2/instance")).thenReturn("ip3@-@defaultInstance");
+
when(regCenter.getDirectly("/test_job/sharding/0/instance")).thenReturn("ip1@-@defaultInstance");
+
when(regCenter.getDirectly("/test_job/sharding/1/instance")).thenReturn("ip2@-@defaultInstance");
+
when(regCenter.getDirectly("/test_job/sharding/2/instance")).thenReturn("ip3@-@defaultInstance");
JobBriefInfo jobBrief = jobStatisticsAPI.getJobBriefInfo("test_job");
assertThat(jobBrief.getStatus(),
is(JobBriefInfo.JobStatus.SHARDING_FLAG));
}
@@ -133,11 +134,11 @@ class JobStatisticsAPIImplTest {
when(regCenter.getChildrenKeys("/test_job_1/servers")).thenReturn(Arrays.asList("ip1",
"ip2"));
when(regCenter.getChildrenKeys("/test_job_2/servers")).thenReturn(Arrays.asList("ip3",
"ip4"));
when(regCenter.getChildrenKeys("/test_job_1/sharding")).thenReturn(Arrays.asList("0",
"1"));
-
when(regCenter.get("/test_job_1/sharding/0/instance")).thenReturn("ip1@-@defaultInstance");
-
when(regCenter.get("/test_job_1/sharding/1/instance")).thenReturn("ip2@-@defaultInstance");
+
when(regCenter.getDirectly("/test_job_1/sharding/0/instance")).thenReturn("ip1@-@defaultInstance");
+
when(regCenter.getDirectly("/test_job_1/sharding/1/instance")).thenReturn("ip2@-@defaultInstance");
when(regCenter.getChildrenKeys("/test_job_2/sharding")).thenReturn(Arrays.asList("0",
"1"));
-
when(regCenter.get("/test_job_2/sharding/0/instance")).thenReturn("ip3@-@defaultInstance");
-
when(regCenter.get("/test_job_2/sharding/1/instance")).thenReturn("ip4@-@defaultInstance");
+
when(regCenter.getDirectly("/test_job_2/sharding/0/instance")).thenReturn("ip3@-@defaultInstance");
+
when(regCenter.getDirectly("/test_job_2/sharding/1/instance")).thenReturn("ip4@-@defaultInstance");
when(regCenter.getChildrenKeys("/test_job_1/instances")).thenReturn(Arrays.asList("ip1@-@defaultInstance",
"ip2@-@defaultInstance"));
when(regCenter.getChildrenKeys("/test_job_2/instances")).thenReturn(Arrays.asList("ip3@-@defaultInstance",
"ip4@-@defaultInstance"));
int i = 0;
diff --git
a/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ShardingStatisticsAPIImplTest.java
b/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ShardingStatisticsAPIImplTest.java
index 863bf8bd0..2c0dbb62b 100644
---
a/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ShardingStatisticsAPIImplTest.java
+++
b/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/statistics/ShardingStatisticsAPIImplTest.java
@@ -50,10 +50,15 @@ class ShardingStatisticsAPIImplTest {
@Test
void assertGetShardingInfo() {
when(regCenter.getChildrenKeys("/test_job/sharding")).thenReturn(Arrays.asList("0",
"1", "2", "3"));
-
when(regCenter.get("/test_job/sharding/0/instance")).thenReturn("ip1@-@1234");
+
when(regCenter.get("/test_job/sharding/0/instance")).thenReturn("ip0@-@0000");
when(regCenter.get("/test_job/sharding/1/instance")).thenReturn("ip2@-@2341");
when(regCenter.get("/test_job/sharding/2/instance")).thenReturn("ip3@-@3412");
when(regCenter.get("/test_job/sharding/3/instance")).thenReturn("ip4@-@4123");
+
when(regCenter.getDirectly("/test_job/sharding/0/instance")).thenReturn("ip1@-@1234");
+
when(regCenter.getDirectly("/test_job/sharding/1/instance")).thenReturn("ip2@-@2341");
+
when(regCenter.getDirectly("/test_job/sharding/2/instance")).thenReturn("ip3@-@3412");
+
when(regCenter.getDirectly("/test_job/sharding/3/instance")).thenReturn("ip4@-@4123");
+
when(regCenter.get("/test_job/instances/ip0@-@0000")).thenReturn("jobInstanceId:
ip0@-@0000\nserverIp: ip0\n");
when(regCenter.get("/test_job/instances/ip1@-@1234")).thenReturn("jobInstanceId:
ip1@-@1234\nserverIp: ip1\n");
when(regCenter.get("/test_job/instances/ip2@-@2341")).thenReturn("jobInstanceId:
ip2@-@2341\nserverIp: ip2\n");
when(regCenter.get("/test_job/instances/ip3@-@3412")).thenReturn("jobInstanceId:
ip3@-@3412\nserverIp: ip3\n");