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 d949032c7 Handle missing instances during lifecycle traversal (#2534)
d949032c7 is described below

commit d949032c7bf6c26e3b3b7e1285d744d2bd63cba2
Author: Hongsheng Zhong <[email protected]>
AuthorDate: Wed Jul 29 11:18:16 2026 +0800

    Handle missing instances during lifecycle traversal (#2534)
    
    What: Read instance metadata directly in lifecycle traversal, skip 
instances that disappear from statistics, and treat concurrent shutdown 
deletion as idempotent.
    
    Why: Cached metadata can outlive an ephemeral instance after child 
enumeration, causing ghost counts, stale identities, or null-related failures.
    
    Risk: This does not provide an atomic snapshot across assignment, metadata, 
and sharding marker nodes.
---
 .../lifecycle/internal/operate/JobOperateAPIImpl.java      | 11 +++++++----
 .../internal/statistics/JobStatisticsAPIImpl.java          |  6 +++++-
 .../internal/statistics/ServerStatisticsAPIImpl.java       |  6 +++++-
 .../internal/statistics/ShardingStatisticsAPIImpl.java     |  7 ++++---
 .../lifecycle/internal/operate/JobOperateAPIImplTest.java  | 14 ++++++++++++--
 .../internal/statistics/JobStatisticsAPIImplTest.java      | 12 +++++++++++-
 .../internal/statistics/ServerStatisticsAPIImplTest.java   | 11 +++++++----
 .../internal/statistics/ShardingStatisticsAPIImplTest.java |  8 ++++++--
 8 files changed, 57 insertions(+), 18 deletions(-)

diff --git 
a/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/operate/JobOperateAPIImpl.java
 
b/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/operate/JobOperateAPIImpl.java
index 87f750377..a61ed7ab4 100644
--- 
a/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/operate/JobOperateAPIImpl.java
+++ 
b/lifecycle/src/main/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/operate/JobOperateAPIImpl.java
@@ -96,8 +96,7 @@ public final class JobOperateAPIImpl implements JobOperateAPI 
{
         if (null != jobName && null != serverIp) {
             JobNodePath jobNodePath = new JobNodePath(jobName);
             for (String each : 
regCenter.getChildrenKeys(jobNodePath.getInstancesNodePath())) {
-                JobInstance jobInstance = 
YamlEngine.unmarshal(regCenter.get(jobNodePath.getInstanceNodePath(each)), 
JobInstance.class);
-                if (serverIp.equals(jobInstance.getServerIp())) {
+                if (isInstanceOnServer(jobNodePath, each, serverIp)) {
                     regCenter.remove(jobNodePath.getInstanceNodePath(each));
                 }
             }
@@ -112,8 +111,7 @@ public final class JobOperateAPIImpl implements 
JobOperateAPI {
                 JobNodePath jobNodePath = new JobNodePath(job);
                 List<String> instances = 
regCenter.getChildrenKeys(jobNodePath.getInstancesNodePath());
                 for (String each : instances) {
-                    JobInstance jobInstance = 
YamlEngine.unmarshal(regCenter.get(jobNodePath.getInstanceNodePath(each)), 
JobInstance.class);
-                    if (serverIp.equals(jobInstance.getServerIp())) {
+                    if (isInstanceOnServer(jobNodePath, each, serverIp)) {
                         
regCenter.remove(jobNodePath.getInstanceNodePath(each));
                     }
                 }
@@ -121,6 +119,11 @@ public final class JobOperateAPIImpl implements 
JobOperateAPI {
         }
     }
     
+    private boolean isInstanceOnServer(final JobNodePath jobNodePath, final 
String instanceId, final String serverIp) {
+        String instanceData = 
regCenter.getDirectly(jobNodePath.getInstanceNodePath(instanceId));
+        return null != instanceData && 
serverIp.equals(YamlEngine.unmarshal(instanceData, 
JobInstance.class).getServerIp());
+    }
+    
     @Override
     public void remove(final String jobName, final String serverIp) {
         shutdown(jobName, serverIp);
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 3d42af8b8..a46f5cb96 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
@@ -160,7 +160,11 @@ public final class JobStatisticsAPIImpl implements 
JobStatisticsAPI {
         JobNodePath jobNodePath = new JobNodePath(jobName);
         List<String> instances = 
regCenter.getChildrenKeys(jobNodePath.getInstancesNodePath());
         for (String each : instances) {
-            JobInstance jobInstance = 
YamlEngine.unmarshal(regCenter.get(jobNodePath.getInstanceNodePath(each)), 
JobInstance.class);
+            String instanceData = 
regCenter.getDirectly(jobNodePath.getInstanceNodePath(each));
+            if (null == instanceData) {
+                continue;
+            }
+            JobInstance jobInstance = YamlEngine.unmarshal(instanceData, 
JobInstance.class);
             if (ip.equals(jobInstance.getServerIp())) {
                 result++;
             }
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 a1da756ba..a2ecc5c99 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
@@ -67,7 +67,11 @@ public final class ServerStatisticsAPIImpl implements 
ServerStatisticsAPI {
             }
             List<String> instances = 
regCenter.getChildrenKeys(jobNodePath.getInstancesNodePath());
             for (String each : instances) {
-                JobInstance jobInstance = 
YamlEngine.unmarshal(regCenter.get(jobNodePath.getInstanceNodePath(each)), 
JobInstance.class);
+                String instanceData = 
regCenter.getDirectly(jobNodePath.getInstanceNodePath(each));
+                if (null == instanceData) {
+                    continue;
+                }
+                JobInstance jobInstance = YamlEngine.unmarshal(instanceData, 
JobInstance.class);
                 if (null != jobInstance) {
                     ServerBriefInfo serverInfo = 
servers.get(jobInstance.getServerIp());
                     if (null != serverInfo) {
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 071263736..757b8b7d4 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
@@ -57,11 +57,12 @@ public final class ShardingStatisticsAPIImpl implements 
ShardingStatisticsAPI {
         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));
+        String instanceData = null == instanceId ? null : 
regCenter.getDirectly(jobNodePath.getInstanceNodePath(instanceId));
+        boolean shardingError = null == instanceData;
         
result.setStatus(ShardingInfo.ShardingStatus.getShardingStatus(disabled, 
running, shardingError));
         
result.setFailover(regCenter.isExisted(jobNodePath.getShardingNodePath(item, 
"failover")));
-        if (null != instanceId) {
-            JobInstance jobInstance = 
YamlEngine.unmarshal(regCenter.get(jobNodePath.getInstanceNodePath(instanceId)),
 JobInstance.class);
+        if (null != instanceData) {
+            JobInstance jobInstance = YamlEngine.unmarshal(instanceData, 
JobInstance.class);
             result.setServerIp(jobInstance.getServerIp());
             result.setInstanceId(jobInstance.getJobInstanceId());
         }
diff --git 
a/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/operate/JobOperateAPIImplTest.java
 
b/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/operate/JobOperateAPIImplTest.java
index a295f521d..c655aebc8 100644
--- 
a/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/operate/JobOperateAPIImplTest.java
+++ 
b/lifecycle/src/test/java/org/apache/shardingsphere/elasticjob/lifecycle/internal/operate/JobOperateAPIImplTest.java
@@ -30,6 +30,7 @@ import java.util.Collections;
 
 import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
@@ -114,11 +115,19 @@ class JobOperateAPIImplTest {
     @Test
     void assertShutdownWithJobNameAndServerIp() {
         
when(regCenter.getChildrenKeys("/test_job/instances")).thenReturn(Collections.singletonList("localhost@-@defaultInstance"));
-        
when(regCenter.get("/test_job/instances/localhost@-@defaultInstance")).thenReturn("jobInstanceId:
 localhost@-@defaultInstance\nserverIp: localhost\n");
+        
when(regCenter.getDirectly("/test_job/instances/localhost@-@defaultInstance")).thenReturn("jobInstanceId:
 localhost@-@defaultInstance\nserverIp: localhost\n");
         jobOperateAPI.shutdown("test_job", "localhost");
         
verify(regCenter).remove("/test_job/instances/localhost@-@defaultInstance");
     }
     
+    @Test
+    void assertShutdownWithJobNameAndServerIpWhenInstanceDisappears() {
+        
when(regCenter.getChildrenKeys("/test_job/instances")).thenReturn(Collections.singletonList("localhost@-@defaultInstance"));
+        
when(regCenter.get("/test_job/instances/localhost@-@defaultInstance")).thenReturn("jobInstanceId:
 localhost@-@defaultInstance\nserverIp: localhost\n");
+        jobOperateAPI.shutdown("test_job", "localhost");
+        verify(regCenter, 
never()).remove("/test_job/instances/localhost@-@defaultInstance");
+    }
+    
     @Test
     void assertShutdownWithJobName() {
         
when(regCenter.getChildrenKeys("/test_job/instances")).thenReturn(Arrays.asList("ip1@-@defaultInstance",
 "ip2@-@defaultInstance"));
@@ -134,9 +143,10 @@ class JobOperateAPIImplTest {
         
when(regCenter.getChildrenKeys("/test_job2/instances")).thenReturn(Collections.singletonList("localhost@-@defaultInstance"));
         
when(regCenter.get("/test_job1/instances/localhost@-@defaultInstance")).thenReturn("jobInstanceId:
 localhost@-@defaultInstance\nserverIp: localhost\n");
         
when(regCenter.get("/test_job2/instances/localhost@-@defaultInstance")).thenReturn("jobInstanceId:
 localhost@-@defaultInstance\nserverIp: localhost\n");
+        
when(regCenter.getDirectly("/test_job2/instances/localhost@-@defaultInstance")).thenReturn("jobInstanceId:
 localhost@-@defaultInstance\nserverIp: localhost\n");
         jobOperateAPI.shutdown(null, "localhost");
         verify(regCenter).getChildrenKeys("/");
-        
verify(regCenter).remove("/test_job1/instances/localhost@-@defaultInstance");
+        verify(regCenter, 
never()).remove("/test_job1/instances/localhost@-@defaultInstance");
         
verify(regCenter).remove("/test_job2/instances/localhost@-@defaultInstance");
     }
     
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 a736d3517..83fbce092 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
@@ -161,7 +161,7 @@ class JobStatisticsAPIImplTest {
         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");
+        
when(regCenter.getDirectly("/test_job_1/instances/ip1@-@defaultInstance")).thenReturn("jobInstanceId:
 ip1@-@defaultInstance\nserverIp: ip1\n");
         int i = 0;
         for (JobBriefInfo each : jobStatisticsAPI.getJobsBriefInfo("ip1")) {
             assertThat(each.getJobName(), is("test_job_" + ++i));
@@ -174,4 +174,14 @@ class JobStatisticsAPIImplTest {
             }
         }
     }
+    
+    @Test
+    void assertGetJobsBriefInfoByIpWhenInstanceDisappears() {
+        
when(regCenter.getChildrenKeys("/")).thenReturn(Collections.singletonList("test_job"));
+        when(regCenter.isExisted("/test_job/servers/ip1")).thenReturn(true);
+        
when(regCenter.getChildrenKeys("/test_job/instances")).thenReturn(Collections.singletonList("ip1@-@defaultInstance"));
+        
when(regCenter.get("/test_job/instances/ip1@-@defaultInstance")).thenReturn("jobInstanceId:
 ip1@-@defaultInstance\nserverIp: ip1\n");
+        JobBriefInfo actual = 
jobStatisticsAPI.getJobsBriefInfo("ip1").iterator().next();
+        assertThat(actual.getInstanceCount(), is(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 1d7b40a42..4dc6074e8 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
@@ -67,13 +67,16 @@ class ServerStatisticsAPIImplTest {
         
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.getChildrenKeys("/test_job1/instances")).thenReturn(Collections.singletonList("ip1@-@oldInstance"));
         
         
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");
+        
lenient().when(regCenter.get("/test_job1/instances/ip1@-@oldInstance")).thenReturn("jobInstanceId:
 ip1@-@oldInstance\nserverIp: ip1\n");
+        
when(regCenter.getDirectly("/test_job1/instances/ip1@-@oldInstance")).thenReturn(null);
+        
lenient().when(regCenter.get("/test_job2/instances/ip1@-@defaultInstance")).thenReturn("jobInstanceId:
 ip1@-@defaultInstance\nserverIp: ip1\n");
+        
lenient().when(regCenter.get("/test_job2/instances/ip2@-@defaultInstance2")).thenReturn("jobInstanceId:
 ip2@-@defaultInstance2\nserverIp: ip2\n");
+        
when(regCenter.getDirectly("/test_job2/instances/ip1@-@defaultInstance")).thenReturn("jobInstanceId:
 ip1@-@defaultInstance\nserverIp: ip1\n");
+        
when(regCenter.getDirectly("/test_job2/instances/ip2@-@defaultInstance2")).thenReturn("jobInstanceId:
 ip2@-@defaultInstance2\nserverIp: ip2\n");
         
when(regCenter.getChildrenKeys("/test_job2/instances")).thenReturn(Arrays.asList("ip1@-@defaultInstance",
 "ip2@-@defaultInstance2"));
         
         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 2c0dbb62b..f7157fb68 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
@@ -30,6 +30,7 @@ import java.util.Arrays;
 
 import static org.hamcrest.CoreMatchers.is;
 import static org.hamcrest.MatcherAssert.assertThat;
+import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.Mockito.when;
 
@@ -63,6 +64,9 @@ class ShardingStatisticsAPIImplTest {
         
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");
         
when(regCenter.get("/test_job/instances/ip4@-@4123")).thenReturn("jobInstanceId:
 ip4@-@4123\nserverIp: ip4\n");
+        
when(regCenter.getDirectly("/test_job/instances/ip1@-@1234")).thenReturn("jobInstanceId:
 ip1@-@1234\nserverIp: ip1\n");
+        
when(regCenter.getDirectly("/test_job/instances/ip3@-@3412")).thenReturn("jobInstanceId:
 ip3@-@3412\nserverIp: ip3\n");
+        
when(regCenter.getDirectly("/test_job/instances/ip4@-@4123")).thenReturn("jobInstanceId:
 ip4@-@4123\nserverIp: ip4\n");
         
when(regCenter.isExisted("/test_job/instances/ip4@-@4123")).thenReturn(true);
         
when(regCenter.isExisted("/test_job/sharding/0/running")).thenReturn(true);
         
when(regCenter.isExisted("/test_job/sharding/1/running")).thenReturn(false);
@@ -84,8 +88,8 @@ class ShardingStatisticsAPIImplTest {
                 case 2:
                     assertTrue(each.isFailover());
                     assertThat(each.getStatus(), 
is(ShardingInfo.ShardingStatus.SHARDING_FLAG));
-                    assertThat(each.getServerIp(), is("ip2"));
-                    assertThat(each.getInstanceId(), is("ip2@-@2341"));
+                    assertNull(each.getServerIp());
+                    assertNull(each.getInstanceId());
                     break;
                 case 3:
                     assertThat(each.getStatus(), 
is(ShardingInfo.ShardingStatus.DISABLED));

Reply via email to