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));