From b3823240e5ccd944c58333a7827137d636237eae Mon Sep 17 00:00:00 2001 From: Ayan Alam Date: Mon, 28 Sep 2026 00:59:39 +0530 Subject: [PATCH 1/2] refactor(server): remove role election and legacy scheduler leftovers Nothing has started the role election state machine since #3082. Graphs started with a node id (tests, examples, and Gremlin scripts calling serverStarted(GlobalMasterInfo.master(...))) still built one, which wrote ~role_data schema and leaked an idle executor. TaskManager also kept the task-scheduler pool alive only to close transactions on it. Remove the masterelection classes other than GlobalMasterInfo, the HugeGraph.roleElectionStateMachine() accessor, the server.role_election and server.role.* options, the deprecation warnings added in #3082, and TaskManager's scheduler pool and role callbacks. Old config files that still set the removed keys only get HugeConfig's "redundant option" warning. HugeVertex keeps mapping the ~role_data label to HugeType.SERVER so vertices left in existing graphs still route to the same table. --- .../hugegraph/auth/HugeFactoryAuthProxy.java | 4 +- .../hugegraph/auth/HugeGraphAuthProxy.java | 7 - .../hugegraph/auth/StandardAuthenticator.java | 17 - .../hugegraph/config/ServerOptions.java | 9 - .../apache/hugegraph/core/GraphManager.java | 67 +--- .../java/org/apache/hugegraph/HugeGraph.java | 3 - .../apache/hugegraph/StandardHugeGraph.java | 37 -- .../hugegraph/masterelection/ClusterRole.java | 95 ----- .../masterelection/ClusterRoleStore.java | 27 -- .../hugegraph/masterelection/Config.java | 35 -- .../masterelection/RoleElectionConfig.java | 76 ---- .../masterelection/RoleElectionOptions.java | 95 ----- .../RoleElectionStateMachine.java | 25 -- .../masterelection/RoleListener.java | 33 -- .../StandardClusterRoleStore.java | 220 ----------- .../StandardRoleElectionStateMachine.java | 368 ------------------ .../masterelection/StandardRoleListener.java | 104 ----- .../masterelection/StateMachineContext.java | 46 --- .../hugegraph/structure/HugeVertex.java | 10 +- .../apache/hugegraph/task/TaskManager.java | 51 +-- .../apache/hugegraph/dist/RegisterUtil.java | 2 - .../hugegraph/core/MultiGraphsTest.java | 2 +- .../core/RoleElectionStateMachineTest.java | 330 ---------------- .../apache/hugegraph/unit/UnitTestSuite.java | 2 - .../core/TaskSchedulerServerInfoTest.java | 6 +- 25 files changed, 20 insertions(+), 1651 deletions(-) delete mode 100644 hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/ClusterRole.java delete mode 100644 hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/ClusterRoleStore.java delete mode 100644 hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/Config.java delete mode 100644 hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleElectionConfig.java delete mode 100644 hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleElectionOptions.java delete mode 100644 hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleElectionStateMachine.java delete mode 100644 hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleListener.java delete mode 100644 hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StandardClusterRoleStore.java delete mode 100644 hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StandardRoleElectionStateMachine.java delete mode 100644 hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StandardRoleListener.java delete mode 100644 hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StateMachineContext.java delete mode 100644 hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/core/RoleElectionStateMachineTest.java diff --git a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/HugeFactoryAuthProxy.java b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/HugeFactoryAuthProxy.java index 5ff2925e69..83f0a64f97 100644 --- a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/HugeFactoryAuthProxy.java +++ b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/HugeFactoryAuthProxy.java @@ -400,11 +400,11 @@ private static void registerPrivateActions() { "oneNumericField", "hasSameProperties"); Reflection.registerFieldsToFilter(TaskManager.class, "LOG", "SCHEDULE_PERIOD", "THREADS", "MANAGER", "schedulers", "taskExecutor", "taskDbExecutor", - "serverInfoDbExecutor", "schedulerExecutor", "contexts", + "serverInfoDbExecutor", "contexts", "$assertionsDisabled"); Reflection.registerMethodsToFilter(TaskManager.class, "lambda$0", "resetContext", "closeTaskTx", "setContext", "instance", - "closeSchedulerTx", "notifyNewTask", + "notifyNewTask", "scheduleOrExecuteJob", "scheduleOrExecuteJobForGraph"); Reflection.registerFieldsToFilter(StandardTaskScheduler.class, "LOG", "graph", "serverManager", "taskExecutor", "taskDbExecutor", diff --git a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/HugeGraphAuthProxy.java b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/HugeGraphAuthProxy.java index 510b838437..9ecda6b9f7 100644 --- a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/HugeGraphAuthProxy.java +++ b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/HugeGraphAuthProxy.java @@ -61,7 +61,6 @@ import org.apache.hugegraph.iterator.MapperIterator; import org.apache.hugegraph.kvstore.KvStore; import org.apache.hugegraph.masterelection.GlobalMasterInfo; -import org.apache.hugegraph.masterelection.RoleElectionStateMachine; import org.apache.hugegraph.rpc.RpcServiceConfig4Client; import org.apache.hugegraph.rpc.RpcServiceConfig4Server; import org.apache.hugegraph.schema.EdgeLabel; @@ -871,12 +870,6 @@ public AuthManager authManager() { return this.authManager; } - @Override - public RoleElectionStateMachine roleElectionStateMachine() { - this.verifyAdminPermission(); - return this.hugegraph.roleElectionStateMachine(); - } - @Override public void switchAuthManager(AuthManager authManager) { this.verifyAdminPermission(); diff --git a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/StandardAuthenticator.java b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/StandardAuthenticator.java index aecc8af282..f64ead4fb2 100644 --- a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/StandardAuthenticator.java +++ b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/StandardAuthenticator.java @@ -30,7 +30,6 @@ import org.apache.hugegraph.config.CoreOptions; import org.apache.hugegraph.config.HugeConfig; import org.apache.hugegraph.config.ServerOptions; -import org.apache.hugegraph.masterelection.RoleElectionOptions; import org.apache.hugegraph.rpc.RpcClientProviderWithAuth; import org.apache.hugegraph.util.ConfigUtil; import org.apache.hugegraph.util.E; @@ -132,7 +131,6 @@ public void setup(HugeConfig config) { String raftGroupPeers = config.get(ServerOptions.RAFT_GROUP_PEERS); graphConfig.addProperty(ServerOptions.RAFT_GROUP_PEERS.name(), raftGroupPeers); - this.transferRoleWorkerConfig(graphConfig, config); this.graph = (HugeGraph) GraphFactory.open(graphConfig); @@ -144,21 +142,6 @@ public void setup(HugeConfig config) { } } - private void transferRoleWorkerConfig(HugeConfig graphConfig, HugeConfig config) { - graphConfig.addProperty(RoleElectionOptions.NODE_EXTERNAL_URL.name(), - config.get(ServerOptions.REST_SERVER_URL)); - graphConfig.addProperty(RoleElectionOptions.BASE_TIMEOUT_MILLISECOND.name(), - config.get(RoleElectionOptions.BASE_TIMEOUT_MILLISECOND)); - graphConfig.addProperty(RoleElectionOptions.EXCEEDS_FAIL_COUNT.name(), - config.get(RoleElectionOptions.EXCEEDS_FAIL_COUNT)); - graphConfig.addProperty(RoleElectionOptions.RANDOM_TIMEOUT_MILLISECOND.name(), - config.get(RoleElectionOptions.RANDOM_TIMEOUT_MILLISECOND)); - graphConfig.addProperty(RoleElectionOptions.HEARTBEAT_INTERVAL_SECOND.name(), - config.get(RoleElectionOptions.HEARTBEAT_INTERVAL_SECOND)); - graphConfig.addProperty(RoleElectionOptions.MASTER_DEAD_TIMES.name(), - config.get(RoleElectionOptions.MASTER_DEAD_TIMES)); - } - /** * Verify if a user is legal * diff --git a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/config/ServerOptions.java b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/config/ServerOptions.java index d43b095b76..e6ed6954a1 100644 --- a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/config/ServerOptions.java +++ b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/config/ServerOptions.java @@ -42,15 +42,6 @@ public class ServerOptions extends OptionHolder { 1 ); - public static final ConfigOption ENABLE_SERVER_ROLE_ELECTION = - new ConfigOption<>( - "server.role_election", - "Whether to enable role election, if enabled, the server " + - "will elect a master node in the cluster.", - disallowEmpty(), - false - ); - public static final ConfigOption MAX_WORKER_THREADS = new ConfigOption<>( "restserver.max_worker_threads", diff --git a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/core/GraphManager.java b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/core/GraphManager.java index a386a0ca67..87ff5ca6f4 100644 --- a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/core/GraphManager.java +++ b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/core/GraphManager.java @@ -82,9 +82,6 @@ import org.apache.hugegraph.kvstore.KvStore; import org.apache.hugegraph.kvstore.KvStoreImpl; import org.apache.hugegraph.masterelection.GlobalMasterInfo; -import org.apache.hugegraph.masterelection.RoleElectionOptions; -import org.apache.hugegraph.masterelection.RoleElectionStateMachine; -import org.apache.hugegraph.masterelection.StandardRoleListener; import org.apache.hugegraph.meta.MetaDriver; import org.apache.hugegraph.meta.MetaManager; import org.apache.hugegraph.meta.PdMetaDriver; @@ -189,7 +186,6 @@ public final class GraphManager { private final Set serverUrlsToPd; private final Boolean serverDeployInK8s; private final HugeConfig config; - private RoleElectionStateMachine roleStateMachine; private K8sDriver.CA ca; private final boolean PDExist; @@ -226,7 +222,6 @@ public GraphManager(HugeConfig conf, EventHub hub) { this.rpcClient = new RpcClientProvider(conf); this.pdPeers = conf.get(ServerOptions.PD_PEERS); - this.roleStateMachine = null; this.globalNodeRoleInfo = new GlobalMasterInfo(); this.eventHub = hub; @@ -696,7 +691,7 @@ public void init() { this.waitGraphsReady(); this.checkBackendVersionOrExit(this.conf); - this.serverStarted(this.conf); + this.serverStarted(); this.addMetrics(this.conf); } @@ -1735,9 +1730,6 @@ public void close() { } this.destroyRpcServer(); this.unlistenChanges(); - if (this.roleStateMachine != null) { - this.roleStateMachine.shutdown(); - } } private void startRpcServer() { @@ -1837,8 +1829,6 @@ private void loadGraph(String name, String graphConfPath) { this.transferPdPeersConfig(config); - this.transferRoleWorkerConfig(config); - Graph graph = GraphFactory.open(config); this.graphs.put(defaultSpaceGraphName(name), graph); @@ -1868,21 +1858,6 @@ private void transferPdPeersConfig(HugeConfig config) { } } - private void transferRoleWorkerConfig(HugeConfig config) { - config.setProperty(RoleElectionOptions.NODE_EXTERNAL_URL.name(), - this.conf.get(ServerOptions.REST_SERVER_URL)); - config.setProperty(RoleElectionOptions.BASE_TIMEOUT_MILLISECOND.name(), - this.conf.get(RoleElectionOptions.BASE_TIMEOUT_MILLISECOND)); - config.setProperty(RoleElectionOptions.EXCEEDS_FAIL_COUNT.name(), - this.conf.get(RoleElectionOptions.EXCEEDS_FAIL_COUNT)); - config.setProperty(RoleElectionOptions.RANDOM_TIMEOUT_MILLISECOND.name(), - this.conf.get(RoleElectionOptions.RANDOM_TIMEOUT_MILLISECOND)); - config.setProperty(RoleElectionOptions.HEARTBEAT_INTERVAL_SECOND.name(), - this.conf.get(RoleElectionOptions.HEARTBEAT_INTERVAL_SECOND)); - config.setProperty(RoleElectionOptions.MASTER_DEAD_TIMES.name(), - this.conf.get(RoleElectionOptions.MASTER_DEAD_TIMES)); - } - private void waitGraphsReady() { if (!this.rpcServer.enabled()) { LOG.info("RpcServer is not enabled, skip wait graphs ready"); @@ -1928,16 +1903,6 @@ private void checkBackendVersionOrExit(HugeConfig config) { } private void initNodeRole() { - boolean enableRoleElection = config.get( - ServerOptions.ENABLE_SERVER_ROLE_ELECTION); - if (enableRoleElection) { - LOG.warn("The server.role_election option is deprecated and no " + - "longer supported (removed with server_info persistence). " + - "The configured server.role is still used for local node " + - "role initialization. Set server.role_election=false to " + - "suppress this warning."); - } - String role = config.get(ServerOptions.SERVER_ROLE); E.checkArgument(StringUtils.isNotEmpty(role), "The server role can't be null or empty"); @@ -1946,16 +1911,12 @@ private void initNodeRole() { this.globalNodeRoleInfo.initNodeRole(nodeRole); } - private void serverStarted(HugeConfig conf) { + private void serverStarted() { for (String graph : this.graphs()) { HugeGraph hugegraph = this.graph(graph); assert hugegraph != null; hugegraph.serverStarted(this.globalNodeRoleInfo); } - if (!this.globalNodeRoleInfo.nodeRole().computer() && this.supportRoleElection() && - config.get(ServerOptions.ENABLE_SERVER_ROLE_ELECTION)) { - LOG.info("Skip role state machine init (deprecated with server_info)"); - } } public SchemaTemplate schemaTemplate(String graphSpace, @@ -1964,30 +1925,6 @@ public SchemaTemplate schemaTemplate(String graphSpace, return this.metaManager.schemaTemplate(graphSpace, schemaTemplate); } - private void initRoleStateMachine() { - E.checkArgument(this.roleStateMachine == null, - "Repeated initialization of role state worker"); - this.globalNodeRoleInfo.supportElection(true); - this.roleStateMachine = this.authenticator().graph().roleElectionStateMachine(); - StandardRoleListener listener = new StandardRoleListener(TaskManager.instance(), - this.globalNodeRoleInfo); - this.roleStateMachine.start(listener); - } - - private boolean supportRoleElection() { - try { - if (!(this.authenticator() instanceof StandardAuthenticator)) { - LOG.info("{} authenticator does not support role election currently", - this.authenticator().getClass().getSimpleName()); - return false; - } - return true; - } catch (IllegalStateException e) { - LOG.info("{}, does not support role election currently", e.getMessage()); - return false; - } - } - private void addMetrics(HugeConfig config) { final MetricManager metric = MetricManager.INSTANCE; // Force to add a server reporter diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/HugeGraph.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/HugeGraph.java index 8fcd76d39b..927614488e 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/HugeGraph.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/HugeGraph.java @@ -37,7 +37,6 @@ import org.apache.hugegraph.config.TypedOption; import org.apache.hugegraph.kvstore.KvStore; import org.apache.hugegraph.masterelection.GlobalMasterInfo; -import org.apache.hugegraph.masterelection.RoleElectionStateMachine; import org.apache.hugegraph.rpc.RpcServiceConfig4Client; import org.apache.hugegraph.rpc.RpcServiceConfig4Server; import org.apache.hugegraph.schema.EdgeLabel; @@ -273,8 +272,6 @@ public interface HugeGraph extends Graph { AuthManager authManager(); - RoleElectionStateMachine roleElectionStateMachine(); - void switchAuthManager(AuthManager authManager); TaskScheduler taskScheduler(); diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/StandardHugeGraph.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/StandardHugeGraph.java index 8f2ff95d01..7f79c93bb7 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/StandardHugeGraph.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/StandardHugeGraph.java @@ -69,14 +69,7 @@ import org.apache.hugegraph.io.HugeGraphIoRegistry; import org.apache.hugegraph.job.EphemeralJob; import org.apache.hugegraph.kvstore.KvStore; -import org.apache.hugegraph.masterelection.ClusterRoleStore; -import org.apache.hugegraph.masterelection.Config; import org.apache.hugegraph.masterelection.GlobalMasterInfo; -import org.apache.hugegraph.masterelection.RoleElectionConfig; -import org.apache.hugegraph.masterelection.RoleElectionOptions; -import org.apache.hugegraph.masterelection.RoleElectionStateMachine; -import org.apache.hugegraph.masterelection.StandardClusterRoleStore; -import org.apache.hugegraph.masterelection.StandardRoleElectionStateMachine; import org.apache.hugegraph.memory.MemoryManager; import org.apache.hugegraph.memory.util.RoundUtil; import org.apache.hugegraph.meta.MetaManager; @@ -184,7 +177,6 @@ public class StandardHugeGraph implements HugeGraph { private volatile HugeVariables variables; private String graphSpace; private AuthManager authManager; - private RoleElectionStateMachine roleElectionStateMachine; private String nickname; private String creator; private Date createTime; @@ -226,13 +218,6 @@ public StandardHugeGraph(HugeConfig config) { this.taskManager = TaskManager.instance(); this.name = config.get(CoreOptions.STORE); - // Keep old config files upgrade-safe while ignoring the legacy scheduler. - if (config.containsKey("task.scheduler_type")) { - LOG.warn("Config key 'task.scheduler_type' is deprecated and " + - "ignored. The scheduler is auto-selected by backend " + - "type (hstore -> distributed, others -> local)."); - } - this.started = false; this.closed = false; this.mode = GraphMode.NONE; @@ -364,7 +349,6 @@ public void serverStarted(GlobalMasterInfo nodeInfo) { if (nodeInfo != null && nodeInfo.nodeId() != null) { this.serverInfoManager().initServerInfo(nodeInfo); - this.initRoleStateMachine(nodeInfo.nodeId()); } // TODO: check necessary? @@ -381,22 +365,6 @@ public void serverStarted(GlobalMasterInfo nodeInfo) { this.started = true; } - private void initRoleStateMachine(Id serverId) { - HugeConfig conf = this.configuration; - Config roleConfig = new RoleElectionConfig(serverId.toString(), - conf.get(RoleElectionOptions.NODE_EXTERNAL_URL), - conf.get(RoleElectionOptions.EXCEEDS_FAIL_COUNT), - conf.get( - RoleElectionOptions.RANDOM_TIMEOUT_MILLISECOND), - conf.get( - RoleElectionOptions.HEARTBEAT_INTERVAL_SECOND), - conf.get(RoleElectionOptions.MASTER_DEAD_TIMES), - conf.get( - RoleElectionOptions.BASE_TIMEOUT_MILLISECOND)); - ClusterRoleStore roleStore = new StandardClusterRoleStore(this.params); - this.roleElectionStateMachine = new StandardRoleElectionStateMachine(roleConfig, roleStore); - } - @Override public boolean started() { return this.started; @@ -1256,11 +1224,6 @@ public AuthManager authManager() { return this.authManager; } - @Override - public RoleElectionStateMachine roleElectionStateMachine() { - return this.roleElectionStateMachine; - } - @Override public void switchAuthManager(AuthManager authManager) { this.authManager = authManager; diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/ClusterRole.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/ClusterRole.java deleted file mode 100644 index f85f29cc88..0000000000 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/ClusterRole.java +++ /dev/null @@ -1,95 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.hugegraph.masterelection; - -import java.util.Objects; - -public class ClusterRole { - - private final String node; - private long clock; - private final int epoch; - private final String url; - - public ClusterRole(String node, String url, int epoch) { - this(node, url, epoch, 1); - } - - public ClusterRole(String node, String url, int epoch, long clock) { - this.node = node; - this.url = url; - this.epoch = epoch; - this.clock = clock; - } - - public void increaseClock() { - this.clock++; - } - - public boolean isMaster(String node) { - return Objects.equals(this.node, node); - } - - public int epoch() { - return this.epoch; - } - - public long clock() { - return this.clock; - } - - public void clock(long clock) { - this.clock = clock; - } - - public String node() { - return this.node; - } - - public String url() { - return this.url; - } - - @Override - public boolean equals(Object obj) { - if (this == obj) { - return true; - } - if (!(obj instanceof ClusterRole)) { - return false; - } - ClusterRole clusterRole = (ClusterRole) obj; - return clock == clusterRole.clock && - epoch == clusterRole.epoch && - Objects.equals(node, clusterRole.node); - } - - @Override - public int hashCode() { - return Objects.hash(node, clock, epoch); - } - - @Override - public String toString() { - return "RoleStateData{" + - "node='" + node + '\'' + - ", clock=" + clock + - ", epoch=" + epoch + - '}'; - } -} diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/ClusterRoleStore.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/ClusterRoleStore.java deleted file mode 100644 index f8bee5c0be..0000000000 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/ClusterRoleStore.java +++ /dev/null @@ -1,27 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.hugegraph.masterelection; - -import java.util.Optional; - -public interface ClusterRoleStore { - - boolean updateIfNodePresent(ClusterRole clusterRole); - - Optional query(); -} diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/Config.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/Config.java deleted file mode 100644 index 31de004f97..0000000000 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/Config.java +++ /dev/null @@ -1,35 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.hugegraph.masterelection; - -public interface Config { - - String node(); - - String url(); - - int exceedsFailCount(); - - long randomTimeoutMillisecond(); - - long heartBeatIntervalSecond(); - - int masterDeadTimes(); - - long baseTimeoutMillisecond(); -} diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleElectionConfig.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleElectionConfig.java deleted file mode 100644 index e389530949..0000000000 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleElectionConfig.java +++ /dev/null @@ -1,76 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.hugegraph.masterelection; - -public class RoleElectionConfig implements Config { - - private final String node; - private final String url; - private final int exceedsFailCount; - private final long randomTimeoutMillisecond; - private final long heartBeatIntervalSecond; - private final int masterDeadTimes; - private final long baseTimeoutMillisecond; - - public RoleElectionConfig(String node, String url, int exceedsFailCount, - long randomTimeoutMillisecond, long heartBeatIntervalSecond, - int masterDeadTimes, long baseTimeoutMillisecond) { - this.node = node; - this.url = url; - this.exceedsFailCount = exceedsFailCount; - this.randomTimeoutMillisecond = randomTimeoutMillisecond; - this.heartBeatIntervalSecond = heartBeatIntervalSecond; - this.masterDeadTimes = masterDeadTimes; - this.baseTimeoutMillisecond = baseTimeoutMillisecond; - } - - @Override - public String node() { - return this.node; - } - - @Override - public String url() { - return this.url; - } - - @Override - public int exceedsFailCount() { - return this.exceedsFailCount; - } - - @Override - public long randomTimeoutMillisecond() { - return this.randomTimeoutMillisecond; - } - - @Override - public long heartBeatIntervalSecond() { - return this.heartBeatIntervalSecond; - } - - @Override - public int masterDeadTimes() { - return this.masterDeadTimes; - } - - @Override - public long baseTimeoutMillisecond() { - return this.baseTimeoutMillisecond; - } -} diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleElectionOptions.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleElectionOptions.java deleted file mode 100644 index 748ec36bd0..0000000000 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleElectionOptions.java +++ /dev/null @@ -1,95 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.hugegraph.masterelection; - -import static org.apache.hugegraph.config.OptionChecker.disallowEmpty; -import static org.apache.hugegraph.config.OptionChecker.rangeInt; - -import org.apache.hugegraph.config.ConfigOption; -import org.apache.hugegraph.config.OptionHolder; - -public class RoleElectionOptions extends OptionHolder { - - private RoleElectionOptions() { - super(); - } - - private static volatile RoleElectionOptions instance; - - public static synchronized RoleElectionOptions instance() { - if (instance == null) { - instance = new RoleElectionOptions(); - // Should initialize all static members first, then register. - instance.registerOptions(); - } - return instance; - } - - public static final ConfigOption EXCEEDS_FAIL_COUNT = - new ConfigOption<>( - "server.role.fail_count", - "When the node failed count of update or query heartbeat " + - "is reaches this threshold, the node will become abdication state to guard" + - "safe property.", - rangeInt(0, Integer.MAX_VALUE), - 5 - ); - - public static final ConfigOption NODE_EXTERNAL_URL = - new ConfigOption<>( - "server.role.node_external_url", - "The url of external accessibility.", - disallowEmpty(), - "http://127.0.0.1:8080" - ); - - public static final ConfigOption RANDOM_TIMEOUT_MILLISECOND = - new ConfigOption<>( - "server.role.random_timeout", - "The random timeout in ms that be used when candidate node request " + - "to become master state to reduce competitive voting.", - rangeInt(0, Integer.MAX_VALUE), - 1000 - ); - - public static final ConfigOption HEARTBEAT_INTERVAL_SECOND = - new ConfigOption<>( - "server.role.heartbeat_interval", - "The role state machine heartbeat interval second time.", - rangeInt(0, Integer.MAX_VALUE), - 2 - ); - - public static final ConfigOption MASTER_DEAD_TIMES = - new ConfigOption<>( - "server.role.master_dead_times", - "When the worker node detects that the number of times " + - "the master node fails to update heartbeat reaches this threshold, " + - "the worker node will become to a candidate node.", - rangeInt(0, Integer.MAX_VALUE), - 10 - ); - - public static final ConfigOption BASE_TIMEOUT_MILLISECOND = - new ConfigOption<>( - "server.role.base_timeout", - "The role state machine candidate state base timeout time.", - rangeInt(0, Integer.MAX_VALUE), - 500 - ); -} diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleElectionStateMachine.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleElectionStateMachine.java deleted file mode 100644 index 4d60300d32..0000000000 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleElectionStateMachine.java +++ /dev/null @@ -1,25 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.hugegraph.masterelection; - -public interface RoleElectionStateMachine { - - void shutdown(); - - void start(RoleListener callback); -} diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleListener.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleListener.java deleted file mode 100644 index eba2aa1e5f..0000000000 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/RoleListener.java +++ /dev/null @@ -1,33 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.hugegraph.masterelection; - -public interface RoleListener { - - void onAsRoleMaster(StateMachineContext context); - - void onAsRoleWorker(StateMachineContext context); - - void onAsRoleCandidate(StateMachineContext context); - - void unknown(StateMachineContext context); - - void onAsRoleAbdication(StateMachineContext context); - - void error(StateMachineContext context, Throwable e); -} diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StandardClusterRoleStore.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StandardClusterRoleStore.java deleted file mode 100644 index fa1cc6a617..0000000000 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StandardClusterRoleStore.java +++ /dev/null @@ -1,220 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.hugegraph.masterelection; - -import java.util.ArrayList; -import java.util.Iterator; -import java.util.List; -import java.util.Objects; -import java.util.Optional; - -import org.apache.hugegraph.HugeGraphParams; -import org.apache.hugegraph.auth.SchemaDefine; -import org.apache.hugegraph.backend.query.Condition; -import org.apache.hugegraph.backend.query.ConditionQuery; -import org.apache.hugegraph.backend.store.BackendEntry; -import org.apache.hugegraph.backend.tx.GraphTransaction; -import org.apache.hugegraph.schema.VertexLabel; -import org.apache.hugegraph.structure.HugeVertex; -import org.apache.hugegraph.type.HugeType; -import org.apache.hugegraph.type.define.DataType; -import org.apache.hugegraph.type.define.HugeKeys; -import org.apache.hugegraph.util.Log; -import org.apache.tinkerpop.gremlin.structure.Graph; -import org.apache.tinkerpop.gremlin.structure.T; -import org.apache.tinkerpop.gremlin.structure.Vertex; -import org.slf4j.Logger; - -public class StandardClusterRoleStore implements ClusterRoleStore { - - private static final Logger LOG = Log.logger(StandardClusterRoleStore.class); - private static final int RETRY_QUERY_TIMEOUT = 200; - - private final HugeGraphParams graph; - - private boolean firstTime; - - public StandardClusterRoleStore(HugeGraphParams graph) { - this.graph = graph; - Schema schema = new Schema(graph); - schema.initSchemaIfNeeded(); - this.firstTime = true; - } - - @Override - public boolean updateIfNodePresent(ClusterRole clusterRole) { - // if epoch increase, update and return true - // if epoch equal, ignore different node, return false - Optional oldClusterRoleOpt = this.queryVertex(); - if (oldClusterRoleOpt.isPresent()) { - ClusterRole oldClusterRole = this.from(oldClusterRoleOpt.get()); - if (clusterRole.epoch() < oldClusterRole.epoch()) { - return false; - } - - if (clusterRole.epoch() == oldClusterRole.epoch() && - !Objects.equals(clusterRole.node(), oldClusterRole.node())) { - return false; - } - LOG.trace("Server {} epoch {} begin remove data old epoch {}, ", - clusterRole.node(), clusterRole.epoch(), oldClusterRole.epoch()); - this.graph.systemTransaction().removeVertex((HugeVertex) oldClusterRoleOpt.get()); - this.graph.systemTransaction().commitOrRollback(); - LOG.trace("Server {} epoch {} success remove data old epoch {}, ", - clusterRole.node(), clusterRole.epoch(), oldClusterRole.epoch()); - } - try { - GraphTransaction tx = this.graph.systemTransaction(); - tx.doUpdateIfAbsent(this.constructEntry(clusterRole)); - tx.commitOrRollback(); - LOG.trace("Server {} epoch {} success update data", - clusterRole.node(), clusterRole.epoch()); - } catch (Throwable ignore) { - LOG.trace("Server {} epoch {} fail update data", - clusterRole.node(), clusterRole.epoch()); - return false; - } - - return true; - } - - private BackendEntry constructEntry(ClusterRole clusterRole) { - List list = new ArrayList<>(8); - list.add(T.label); - list.add(P.ROLE_DATA); - - list.add(P.NODE); - list.add(clusterRole.node()); - - list.add(P.URL); - list.add(clusterRole.url()); - - list.add(P.CLOCK); - list.add(clusterRole.clock()); - - list.add(P.EPOCH); - list.add(clusterRole.epoch()); - - list.add(P.TYPE); - list.add("default"); - - HugeVertex vertex = this.graph.systemTransaction() - .constructVertex(false, list.toArray()); - - return this.graph.serializer().writeVertex(vertex); - } - - @Override - public Optional query() { - Optional vertex = this.queryVertex(); - if (!vertex.isPresent() && !this.firstTime) { - // If query nothing, retry once - try { - Thread.sleep(RETRY_QUERY_TIMEOUT); - } catch (InterruptedException ignored) { - } - - vertex = this.queryVertex(); - } - this.firstTime = false; - return vertex.map(this::from); - } - - private ClusterRole from(Vertex vertex) { - String node = (String) vertex.property(P.NODE).value(); - String url = (String) vertex.property(P.URL).value(); - Long clock = (Long) vertex.property(P.CLOCK).value(); - Integer epoch = (Integer) vertex.property(P.EPOCH).value(); - - return new ClusterRole(node, url, epoch, clock); - } - - private Optional queryVertex() { - GraphTransaction tx = this.graph.systemTransaction(); - ConditionQuery query; - if (this.graph.backendStoreFeatures().supportsTaskAndServerVertex()) { - query = new ConditionQuery(HugeType.SERVER); - } else { - query = new ConditionQuery(HugeType.VERTEX); - } - VertexLabel vl = this.graph.graph().vertexLabel(P.ROLE_DATA); - query.eq(HugeKeys.LABEL, vl.id()); - query.query(Condition.eq(vl.primaryKeys().get(0), "default")); - query.showHidden(true); - Iterator vertexIterator = tx.queryVertices(query); - if (vertexIterator.hasNext()) { - return Optional.of(vertexIterator.next()); - } - - return Optional.empty(); - } - - public static final class P { - - public static final String ROLE_DATA = Graph.Hidden.hide("role_data"); - - public static final String LABEL = T.label.getAccessor(); - - public static final String NODE = Graph.Hidden.hide("role_node"); - - public static final String CLOCK = Graph.Hidden.hide("role_clock"); - - public static final String EPOCH = Graph.Hidden.hide("role_epoch"); - - public static final String URL = Graph.Hidden.hide("role_url"); - - public static final String TYPE = Graph.Hidden.hide("role_type"); - } - - public static final class Schema extends SchemaDefine { - - public Schema(HugeGraphParams graph) { - super(graph, P.ROLE_DATA); - } - - @Override - public void initSchemaIfNeeded() { - if (this.existVertexLabel(this.label)) { - return; - } - - String[] properties = this.initProperties(); - - VertexLabel label = this.schema() - .vertexLabel(this.label) - .enableLabelIndex(true) - .usePrimaryKeyId() - .primaryKeys(P.TYPE) - .properties(properties) - .build(); - this.graph.schemaTransaction().addVertexLabel(label); - } - - private String[] initProperties() { - List props = new ArrayList<>(); - - props.add(createPropertyKey(P.NODE, DataType.TEXT)); - props.add(createPropertyKey(P.URL, DataType.TEXT)); - props.add(createPropertyKey(P.CLOCK, DataType.LONG)); - props.add(createPropertyKey(P.EPOCH, DataType.INT)); - props.add(createPropertyKey(P.TYPE, DataType.TEXT)); - - return super.initProperties(props); - } - } -} diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StandardRoleElectionStateMachine.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StandardRoleElectionStateMachine.java deleted file mode 100644 index e813f9e6c3..0000000000 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StandardRoleElectionStateMachine.java +++ /dev/null @@ -1,368 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.hugegraph.masterelection; - -import java.security.SecureRandom; -import java.util.Optional; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.locks.LockSupport; - -import org.apache.hugegraph.util.E; -import org.apache.hugegraph.util.Log; -import org.slf4j.Logger; - -public class StandardRoleElectionStateMachine implements RoleElectionStateMachine { - - private static final Logger LOG = Log.logger(StandardRoleElectionStateMachine.class); - - private final Config config; - private final ClusterRoleStore roleStore; - private final ExecutorService applyThread; - - private volatile boolean shutdown; - private volatile RoleState state; - - public StandardRoleElectionStateMachine(Config config, ClusterRoleStore roleStore) { - this.config = config; - this.roleStore = roleStore; - this.applyThread = Executors.newSingleThreadExecutor(); - this.state = new UnknownState(null); - this.shutdown = false; - } - - @Override - public void shutdown() { - if (this.shutdown) { - return; - } - this.shutdown = true; - this.applyThread.shutdown(); - } - - @Override - public void start(RoleListener stateMachineCallback) { - this.applyThread.execute(() -> this.apply(stateMachineCallback)); - } - - private void apply(RoleListener stateMachineCallback) { - int failCount = 0; - StateMachineContextImpl context = new StateMachineContextImpl(this); - while (!this.shutdown) { - E.checkArgumentNotNull(this.state, "State don't be null"); - try { - RoleState pre = this.state; - this.state = state.transform(context); - LOG.trace("server {} epoch {} role state change {} to {}", - context.node(), context.epoch(), pre.getClass().getSimpleName(), - this.state.getClass().getSimpleName()); - Callback runnable = this.state.callback(stateMachineCallback); - runnable.call(context); - failCount = 0; - } catch (Throwable e) { - stateMachineCallback.error(context, e); - failCount++; - if (failCount >= this.config.exceedsFailCount()) { - this.state = new AbdicationState(context.epoch()); - Callback runnable = this.state.callback(stateMachineCallback); - runnable.call(context); - } - } - } - } - - protected ClusterRoleStore roleStore() { - return this.roleStore; - } - - private interface RoleState { - - SecureRandom SECURE_RANDOM = new SecureRandom(); - - RoleState transform(StateMachineContext context); - - Callback callback(RoleListener callback); - - static void heartBeatPark(StateMachineContext context) { - long heartBeatIntervalSecond = context.config().heartBeatIntervalSecond(); - LockSupport.parkNanos(heartBeatIntervalSecond * 1_000_000_000); - } - - static void randomPark(StateMachineContext context) { - long randomTimeout = context.config().randomTimeoutMillisecond(); - long baseTime = context.config().baseTimeoutMillisecond(); - long timeout = (long) (baseTime + (randomTimeout / 10.0 * SECURE_RANDOM.nextInt(11))); - LockSupport.parkNanos(timeout * 1_000_000); - } - } - - @FunctionalInterface - private interface Callback { - - void call(StateMachineContext context); - } - - private static class UnknownState implements RoleState { - - final Integer epoch; - - public UnknownState(Integer epoch) { - this.epoch = epoch; - } - - @Override - public RoleState transform(StateMachineContext context) { - ClusterRoleStore adapter = context.roleStore(); - Optional clusterRoleOpt = adapter.query(); - if (!clusterRoleOpt.isPresent()) { - context.reset(); - Integer nextEpoch = this.epoch == null ? 1 : this.epoch + 1; - context.epoch(nextEpoch); - return new CandidateState(nextEpoch); - } - - ClusterRole clusterRole = clusterRoleOpt.get(); - if (this.epoch != null && clusterRole.epoch() < this.epoch) { - context.reset(); - Integer nextEpoch = this.epoch + 1; - context.epoch(nextEpoch); - return new CandidateState(nextEpoch); - } - - context.epoch(clusterRole.epoch()); - context.master(new MasterServerInfoImpl(clusterRole.node(), clusterRole.url())); - if (clusterRole.isMaster(context.node())) { - return new MasterState(clusterRole); - } else { - return new WorkerState(clusterRole); - } - } - - @Override - public Callback callback(RoleListener callback) { - return callback::unknown; - } - } - - private static class AbdicationState implements RoleState { - - private final Integer epoch; - - public AbdicationState(Integer epoch) { - this.epoch = epoch; - } - - @Override - public RoleState transform(StateMachineContext context) { - context.master(null); - RoleState.heartBeatPark(context); - return new UnknownState(this.epoch).transform(context); - } - - @Override - public Callback callback(RoleListener callback) { - return callback::onAsRoleAbdication; - } - } - - private static class MasterState implements RoleState { - - private final ClusterRole clusterRole; - - public MasterState(ClusterRole clusterRole) { - this.clusterRole = clusterRole; - } - - @Override - public RoleState transform(StateMachineContext context) { - this.clusterRole.increaseClock(); - RoleState.heartBeatPark(context); - if (context.roleStore().updateIfNodePresent(this.clusterRole)) { - return this; - } - context.reset(); - context.epoch(this.clusterRole.epoch()); - return new UnknownState(this.clusterRole.epoch()).transform(context); - } - - @Override - public Callback callback(RoleListener callback) { - return callback::onAsRoleMaster; - } - } - - private static class WorkerState implements RoleState { - - private ClusterRole clusterRole; - private int clock; - - public WorkerState(ClusterRole clusterRole) { - this.clusterRole = clusterRole; - this.clock = 0; - } - - @Override - public RoleState transform(StateMachineContext context) { - RoleState.heartBeatPark(context); - RoleState nextState = new UnknownState(this.clusterRole.epoch()).transform(context); - if (nextState instanceof WorkerState) { - this.merge((WorkerState) nextState); - if (this.clock > context.config().masterDeadTimes()) { - return new CandidateState(this.clusterRole.epoch() + 1); - } else { - return this; - } - } else { - return nextState; - } - } - - @Override - public Callback callback(RoleListener callback) { - return callback::onAsRoleWorker; - } - - public void merge(WorkerState state) { - if (state.clusterRole.epoch() > this.clusterRole.epoch()) { - this.clock = 0; - this.clusterRole = state.clusterRole; - } else if (state.clusterRole.epoch() < this.clusterRole.epoch()) { - throw new IllegalStateException("Epoch must increase"); - } else if (state.clusterRole.epoch() == this.clusterRole.epoch() && - state.clusterRole.clock() < this.clusterRole.clock()) { - throw new IllegalStateException("Clock must increase"); - } else if (state.clusterRole.epoch() == this.clusterRole.epoch() && - state.clusterRole.clock() > this.clusterRole.clock()) { - this.clock = 0; - this.clusterRole = state.clusterRole; - } else { - this.clock++; - } - } - } - - private static class CandidateState implements RoleState { - - private final Integer epoch; - - public CandidateState(Integer epoch) { - this.epoch = epoch; - } - - @Override - public RoleState transform(StateMachineContext context) { - RoleState.randomPark(context); - int epoch = this.epoch == null ? 1 : this.epoch; - ClusterRole clusterRole = new ClusterRole(context.config().node(), - context.config().url(), epoch); - // The master failover completed - context.epoch(clusterRole.epoch()); - if (context.roleStore().updateIfNodePresent(clusterRole)) { - context.master(new MasterServerInfoImpl(clusterRole.node(), clusterRole.url())); - return new MasterState(clusterRole); - } else { - return new UnknownState(epoch).transform(context); - } - } - - @Override - public Callback callback(RoleListener callback) { - return callback::onAsRoleCandidate; - } - } - - private static class StateMachineContextImpl implements StateMachineContext { - - private Integer epoch; - private final String node; - private final StandardRoleElectionStateMachine machine; - - private MasterServerInfo masterServerInfo; - - public StateMachineContextImpl(StandardRoleElectionStateMachine machine) { - this.node = machine.config.node(); - this.machine = machine; - } - - @Override - public void master(MasterServerInfo info) { - this.masterServerInfo = info; - } - - @Override - public Integer epoch() { - return this.epoch; - } - - @Override - public String node() { - return this.node; - } - - @Override - public void epoch(Integer epoch) { - this.epoch = epoch; - } - - @Override - public ClusterRoleStore roleStore() { - return this.machine.roleStore(); - } - - @Override - public Config config() { - return this.machine.config; - } - - @Override - public MasterServerInfo master() { - return this.masterServerInfo; - } - - @Override - public RoleElectionStateMachine stateMachine() { - return this.machine; - } - - @Override - public void reset() { - this.epoch = null; - } - } - - private static class MasterServerInfoImpl implements StateMachineContext.MasterServerInfo { - - private final String node; - private final String url; - - public MasterServerInfoImpl(String node, String url) { - this.node = node; - this.url = url; - } - - @Override - public String url() { - return this.url; - } - - @Override - public String node() { - return this.node; - } - } -} diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StandardRoleListener.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StandardRoleListener.java deleted file mode 100644 index ce91ca2d5f..0000000000 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StandardRoleListener.java +++ /dev/null @@ -1,104 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.hugegraph.masterelection; - -import java.util.Objects; - -import org.apache.hugegraph.task.TaskManager; -import org.apache.hugegraph.util.Log; -import org.slf4j.Logger; - -public class StandardRoleListener implements RoleListener { - - private static final Logger LOG = Log.logger(StandardRoleListener.class); - - private final TaskManager taskManager; - - private final GlobalMasterInfo roleInfo; - - private volatile boolean selfIsMaster; - - public StandardRoleListener(TaskManager taskManager, - GlobalMasterInfo roleInfo) { - this.taskManager = taskManager; - this.roleInfo = roleInfo; - this.selfIsMaster = false; - } - - @Override - public void onAsRoleMaster(StateMachineContext context) { - if (!selfIsMaster) { - this.taskManager.onAsRoleMaster(); - LOG.info("Server {} change to master role", context.config().node()); - } - this.updateMasterInfo(context); - this.selfIsMaster = true; - } - - @Override - public void onAsRoleWorker(StateMachineContext context) { - if (this.selfIsMaster) { - this.taskManager.onAsRoleWorker(); - LOG.info("Server {} change to worker role", context.config().node()); - } - this.updateMasterInfo(context); - this.selfIsMaster = false; - } - - @Override - public void onAsRoleCandidate(StateMachineContext context) { - // pass - } - - @Override - public void onAsRoleAbdication(StateMachineContext context) { - if (this.selfIsMaster) { - this.taskManager.onAsRoleWorker(); - LOG.info("Server {} change to worker role", context.config().node()); - } - this.updateMasterInfo(context); - this.selfIsMaster = false; - } - - @Override - public void error(StateMachineContext context, Throwable e) { - LOG.error("Server {} exception occurred", context.config().node(), e); - } - - @Override - public void unknown(StateMachineContext context) { - if (this.selfIsMaster) { - this.taskManager.onAsRoleWorker(); - LOG.info("Server {} change to worker role", context.config().node()); - } - this.updateMasterInfo(context); - - this.selfIsMaster = false; - } - - public void updateMasterInfo(StateMachineContext context) { - StateMachineContext.MasterServerInfo master = context.master(); - if (master == null) { - this.roleInfo.resetMasterInfo(); - return; - } - - boolean isMaster = Objects.equals(context.node(), master.node()); - this.roleInfo.masterInfo(isMaster, master.url()); - } -} diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StateMachineContext.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StateMachineContext.java deleted file mode 100644 index 9c0789a41e..0000000000 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/masterelection/StateMachineContext.java +++ /dev/null @@ -1,46 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.hugegraph.masterelection; - -public interface StateMachineContext { - - Integer epoch(); - - String node(); - - RoleElectionStateMachine stateMachine(); - - void epoch(Integer epoch); - - Config config(); - - MasterServerInfo master(); - - void master(MasterServerInfo info); - - ClusterRoleStore roleStore(); - - void reset(); - - interface MasterServerInfo { - - String url(); - - String node(); - } -} diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/structure/HugeVertex.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/structure/HugeVertex.java index 389b0f0e8c..ae1a45b845 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/structure/HugeVertex.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/structure/HugeVertex.java @@ -39,7 +39,6 @@ import org.apache.hugegraph.backend.serializer.BytesBuffer; import org.apache.hugegraph.backend.tx.GraphTransaction; import org.apache.hugegraph.config.CoreOptions; -import org.apache.hugegraph.masterelection.StandardClusterRoleStore; import org.apache.hugegraph.perf.PerfUtil.Watched; import org.apache.hugegraph.schema.EdgeLabel; import org.apache.hugegraph.schema.PropertyKey; @@ -58,6 +57,7 @@ import org.apache.logging.log4j.util.Strings; import org.apache.tinkerpop.gremlin.structure.Direction; import org.apache.tinkerpop.gremlin.structure.Edge; +import org.apache.tinkerpop.gremlin.structure.Graph; import org.apache.tinkerpop.gremlin.structure.T; import org.apache.tinkerpop.gremlin.structure.Vertex; import org.apache.tinkerpop.gremlin.structure.VertexProperty; @@ -71,6 +71,12 @@ public class HugeVertex extends HugeElement implements Vertex, Cloneable { private static final List EMPTY_LIST = ImmutableList.of(); + /* + * Labels of the removed server info and role election vertices, graphs + * created by older versions may still store them + */ + private static final String LEGACY_ROLE_DATA_LABEL = Graph.Hidden.hide("role_data"); + private Id id; private VertexLabel label; protected Collection edges; @@ -101,7 +107,7 @@ public HugeType type() { } if (label != null && (label.name().equals(HugeServerInfo.P.SERVER) || - label.name().equals(StandardClusterRoleStore.P.ROLE_DATA))) { + label.name().equals(LEGACY_ROLE_DATA_LABEL))) { return HugeType.SERVER; } return HugeType.VERTEX; diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/TaskManager.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/TaskManager.java index ec063754d8..135fb47a42 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/TaskManager.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/TaskManager.java @@ -36,10 +36,7 @@ /** * Central task management system that coordinates task scheduling and execution. - * Manages task schedulers for different graphs and handles role-based execution. - *

- * Note: The local master-worker mechanism will be deprecated in version 1.7 - * (configuration has been removed from config files). + * Manages the task schedulers of each graph and the executors they share. */ public final class TaskManager { @@ -49,7 +46,6 @@ public final class TaskManager { public static final String TASK_WORKER = TASK_WORKER_PREFIX + "-%d"; public static final String TASK_DB_WORKER = "task-db-worker-%d"; public static final String SERVER_INFO_DB_WORKER = "server-info-db-worker-%d"; - public static final String TASK_SCHEDULER = "task-scheduler-%d"; public static final String OLAP_TASK_WORKER = "olap-task-worker-%d"; public static final String SCHEMA_TASK_WORKER = "schema-task-worker-%d"; @@ -66,7 +62,6 @@ public final class TaskManager { private final ExecutorService taskExecutor; private final ExecutorService taskDbExecutor; private final ExecutorService serverInfoDbExecutor; - private final PausableScheduledThreadPool schedulerExecutor; private final ExecutorService schemaTaskExecutor; private final ExecutorService olapTaskExecutor; @@ -93,10 +88,6 @@ private TaskManager(int pool) { this.ephemeralTaskExecutor = ExecutorUtil.newFixedThreadPool(pool, EPHEMERAL_TASK_WORKER); this.distributedSchedulerExecutor = ExecutorUtil.newPausableScheduledThreadPool(1, DISTRIBUTED_TASK_SCHEDULER); - - // For a schedule task to run, just one thread is ok - this.schedulerExecutor = ExecutorUtil.newPausableScheduledThreadPool( - 1, TASK_SCHEDULER); } public void addScheduler(HugeGraphParams graph) { @@ -160,10 +151,6 @@ public void closeScheduler(HugeGraphParams graph) { this.closeTaskTx(graph); } - if (!this.schedulerExecutor.isTerminated()) { - this.closeSchedulerTx(graph); - } - if (!this.distributedSchedulerExecutor.isTerminated()) { this.closeDistributedSchedulerTx(graph); } @@ -190,21 +177,6 @@ private void closeTaskTx(HugeGraphParams graph) { } } - private void closeSchedulerTx(HugeGraphParams graph) { - final Callable closeTx = () -> { - // Do close-tx for the current thread - graph.closeTx(); - // Let other threads run - Thread.yield(); - return null; - }; - try { - this.schedulerExecutor.submit(closeTx).get(); - } catch (Exception e) { - throw new HugeException("Exception when closing scheduler tx", e); - } - } - private void closeDistributedSchedulerTx(HugeGraphParams graph) { final Callable closeTx = () -> { // Do close-tx for the current thread @@ -236,19 +208,10 @@ public void shutdown(long timeout) { assert this.schedulers.isEmpty() : this.schedulers.size(); Throwable ex = null; - boolean terminated = this.schedulerExecutor.isTerminated(); + boolean terminated = this.distributedSchedulerExecutor.isTerminated(); final TimeUnit unit = TimeUnit.SECONDS; - if (!this.schedulerExecutor.isShutdown()) { - this.schedulerExecutor.shutdown(); - try { - terminated = this.schedulerExecutor.awaitTermination(timeout, unit); - } catch (Throwable e) { - ex = e; - } - } - - if (terminated && !this.distributedSchedulerExecutor.isShutdown()) { + if (!this.distributedSchedulerExecutor.isShutdown()) { this.distributedSchedulerExecutor.shutdown(); try { terminated = this.distributedSchedulerExecutor.awaitTermination(timeout, unit); @@ -331,14 +294,6 @@ public int pendingTasks() { return size; } - public void onAsRoleMaster() { - // ServerInfo based role propagation is deprecated. - } - - public void onAsRoleWorker() { - // ServerInfo based role propagation is deprecated. - } - private static final ThreadLocal CONTEXTS = new ThreadLocal<>(); public static void setContext(String context) { diff --git a/hugegraph-server/hugegraph-dist/src/main/java/org/apache/hugegraph/dist/RegisterUtil.java b/hugegraph-server/hugegraph-dist/src/main/java/org/apache/hugegraph/dist/RegisterUtil.java index 44c074c6a1..380dd996ad 100644 --- a/hugegraph-server/hugegraph-dist/src/main/java/org/apache/hugegraph/dist/RegisterUtil.java +++ b/hugegraph-server/hugegraph-dist/src/main/java/org/apache/hugegraph/dist/RegisterUtil.java @@ -30,7 +30,6 @@ import org.apache.hugegraph.config.CoreOptions; import org.apache.hugegraph.config.HugeConfig; import org.apache.hugegraph.config.OptionSpace; -import org.apache.hugegraph.masterelection.RoleElectionOptions; import org.apache.hugegraph.plugin.HugeGraphPlugin; import org.apache.hugegraph.util.E; import org.apache.hugegraph.util.Log; @@ -51,7 +50,6 @@ public class RegisterUtil { static { OptionSpace.register("core", CoreOptions.instance()); OptionSpace.register("dist", DistOptions.instance()); - OptionSpace.register("masterElection", RoleElectionOptions.instance()); } public static void registerBackends() { diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/core/MultiGraphsTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/core/MultiGraphsTest.java index 7ed172acd9..bb45b47e3b 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/core/MultiGraphsTest.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/core/MultiGraphsTest.java @@ -373,7 +373,7 @@ public void testCreateGraphsWithDifferentNameDifferentBackends() { } @Test - public void testOpenGraphWithDeprecatedTaskSchedulerType() { + public void testOpenGraphWithRemovedTaskSchedulerType() { HugeGraph graph = openGraphWithBackend("legacySchedulerType", "rocksdb", "binary", "task.scheduler_type", "local"); diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/core/RoleElectionStateMachineTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/core/RoleElectionStateMachineTest.java deleted file mode 100644 index eab1e6b791..0000000000 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/core/RoleElectionStateMachineTest.java +++ /dev/null @@ -1,330 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.hugegraph.core; - -import java.util.ArrayList; -import java.util.Collections; -import java.util.HashMap; -import java.util.HashSet; -import java.util.List; -import java.util.Map; -import java.util.Objects; -import java.util.Optional; -import java.util.Set; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.locks.LockSupport; - -import org.apache.hugegraph.masterelection.ClusterRole; -import org.apache.hugegraph.masterelection.ClusterRoleStore; -import org.apache.hugegraph.masterelection.Config; -import org.apache.hugegraph.masterelection.RoleElectionStateMachine; -import org.apache.hugegraph.masterelection.RoleListener; -import org.apache.hugegraph.masterelection.StandardRoleElectionStateMachine; -import org.apache.hugegraph.masterelection.StateMachineContext; -import org.apache.hugegraph.testutil.Assert; -import org.apache.hugegraph.testutil.Utils; -import org.junit.Test; - -public class RoleElectionStateMachineTest { - - private static class LogEntry { - - private final Integer epoch; - - private final String node; - - private final Role role; - - enum Role { - master, - worker, - candidate, - abdication, - unknown - } - - public LogEntry(Integer epoch, String node, Role role) { - this.epoch = epoch; - this.node = node; - this.role = role; - } - - @Override - public boolean equals(Object obj) { - if (this == obj) { - return true; - } - if (!(obj instanceof LogEntry)) { - return false; - } - LogEntry logEntry = (LogEntry) obj; - return Objects.equals(this.epoch, logEntry.epoch) && - Objects.equals(this.node, logEntry.node) && - this.role == logEntry.role; - } - - @Override - public int hashCode() { - return Objects.hash(this.epoch, this.node, this.role); - } - - @Override - public String toString() { - return "LogEntry{" + - "epoch=" + this.epoch + - ", node='" + this.node + '\'' + - ", role=" + this.role + - '}'; - } - } - - private static class TestConfig implements Config { - - private final String node; - - public TestConfig(String node) { - this.node = node; - } - - @Override - public String node() { - return this.node; - } - - @Override - public String url() { - return "http://127.0.0.1:8080"; - } - - @Override - public int exceedsFailCount() { - return 2; - } - - @Override - public long randomTimeoutMillisecond() { - return 400; - } - - @Override - public long heartBeatIntervalSecond() { - return 1; - } - - @Override - public int masterDeadTimes() { - return 5; - } - - @Override - public long baseTimeoutMillisecond() { - return 100; - } - } - - @Test - public void testStateMachine() throws InterruptedException { - final int MAX_COUNT = 200; - CountDownLatch stop = new CountDownLatch(4); - List logRecords = Collections.synchronizedList(new ArrayList<>(MAX_COUNT)); - List masterNodes = Collections.synchronizedList(new ArrayList<>(MAX_COUNT)); - RoleListener callback = new RoleListener() { - - @Override - public void onAsRoleMaster(StateMachineContext context) { - Integer epochId = context.epoch(); - String node = context.node(); - logRecords.add(new LogEntry(epochId, node, LogEntry.Role.master)); - if (logRecords.size() > MAX_COUNT) { - context.stateMachine().shutdown(); - } - Utils.println("master node: " + node); - masterNodes.add(node); - } - - @Override - public void onAsRoleWorker(StateMachineContext context) { - Integer epochId = context.epoch(); - String node = context.node(); - logRecords.add(new LogEntry(epochId, node, LogEntry.Role.worker)); - if (logRecords.size() > MAX_COUNT) { - context.stateMachine().shutdown(); - } - } - - @Override - public void onAsRoleCandidate(StateMachineContext context) { - Integer epochId = context.epoch(); - String node = context.node(); - logRecords.add(new LogEntry(epochId, node, LogEntry.Role.candidate)); - if (logRecords.size() > MAX_COUNT) { - context.stateMachine().shutdown(); - } - } - - @Override - public void unknown(StateMachineContext context) { - Integer epochId = context.epoch(); - String node = context.node(); - logRecords.add(new LogEntry(epochId, node, LogEntry.Role.unknown)); - if (logRecords.size() > MAX_COUNT) { - context.stateMachine().shutdown(); - } - } - - @Override - public void onAsRoleAbdication(StateMachineContext context) { - Integer epochId = context.epoch(); - String node = context.node(); - logRecords.add(new LogEntry(epochId, node, LogEntry.Role.abdication)); - if (logRecords.size() > MAX_COUNT) { - context.stateMachine().shutdown(); - } - } - - @Override - public void error(StateMachineContext context, Throwable e) { - Utils.println("state machine error: node " + - context.node() + " message " + e.getMessage()); - } - }; - - final List clusterRoleLogs = Collections.synchronizedList( - new ArrayList<>(100)); - - final ClusterRoleStore clusterRoleStore = new ClusterRoleStore() { - - volatile int epoch = 0; - - final Map data = new ConcurrentHashMap<>(); - - ClusterRole copy(ClusterRole clusterRole) { - if (clusterRole == null) { - return null; - } - return new ClusterRole(clusterRole.node(), clusterRole.url(), - clusterRole.epoch(), clusterRole.clock()); - } - - @Override - public boolean updateIfNodePresent(ClusterRole clusterRole) { - if (clusterRole.epoch() < this.epoch) { - return false; - } - - ClusterRole copy = this.copy(clusterRole); - ClusterRole newClusterRole = this.data.compute(copy.epoch(), (key, value) -> { - if (copy.epoch() > this.epoch) { - this.epoch = copy.epoch(); - Assert.assertNull(value); - clusterRoleLogs.add(copy); - Utils.println("The node " + copy + " become new master:"); - return copy; - } - - Assert.assertEquals(value.epoch(), copy.epoch()); - if (Objects.equals(value.node(), copy.node()) && - value.clock() <= copy.clock()) { - Utils.println("The master node " + copy + " keep heartbeat"); - clusterRoleLogs.add(copy); - if (value.clock() == copy.clock()) { - Assert.fail("Clock must increase when same epoch and node id"); - } - return copy; - } - return value; - - }); - return Objects.equals(newClusterRole, copy); - } - - @Override - public Optional query() { - return Optional.ofNullable(this.copy(this.data.get(this.epoch))); - } - }; - - RoleElectionStateMachine[] machines = new RoleElectionStateMachine[4]; - Thread node1 = new Thread(() -> { - Config config = new TestConfig("1"); - RoleElectionStateMachine stateMachine = - new StandardRoleElectionStateMachine(config, clusterRoleStore); - machines[1] = stateMachine; - stateMachine.start(callback); - stop.countDown(); - }); - - Thread node2 = new Thread(() -> { - Config config = new TestConfig("2"); - RoleElectionStateMachine stateMachine = - new StandardRoleElectionStateMachine(config, clusterRoleStore); - machines[2] = stateMachine; - stateMachine.start(callback); - stop.countDown(); - }); - - Thread node3 = new Thread(() -> { - Config config = new TestConfig("3"); - RoleElectionStateMachine stateMachine = - new StandardRoleElectionStateMachine(config, clusterRoleStore); - machines[3] = stateMachine; - stateMachine.start(callback); - stop.countDown(); - }); - - node1.start(); - node2.start(); - node3.start(); - - Thread randomShutdown = new Thread(() -> { - Set dropNodes = new HashSet<>(); - while (dropNodes.size() < 3) { - LockSupport.parkNanos(5_000_000_000L); - int size = masterNodes.size(); - if (size < 1) { - continue; - } - String node = masterNodes.get(size - 1); - if (dropNodes.contains(node)) { - continue; - } - machines[Integer.parseInt(node)].shutdown(); - dropNodes.add(node); - Utils.println("----shutdown machine " + node); - } - stop.countDown(); - }); - - randomShutdown.start(); - stop.await(); - - Assert.assertGt(0, logRecords.size()); - Map masters = new HashMap<>(); - for (LogEntry entry : logRecords) { - if (entry.role == LogEntry.Role.master) { - String lastNode = masters.putIfAbsent(entry.epoch, entry.node); - if (lastNode != null) { - Assert.assertEquals(lastNode, entry.node); - } - } - } - - Assert.assertGt(0, masters.size()); - } -} diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/UnitTestSuite.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/UnitTestSuite.java index 0e010ae5f7..4189c1692d 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/UnitTestSuite.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/UnitTestSuite.java @@ -24,7 +24,6 @@ import org.apache.hugegraph.backend.page.QueryListTest; import org.apache.hugegraph.backend.tx.GraphIndexTransactionTest; import org.apache.hugegraph.backend.tx.GraphTransactionTest; -import org.apache.hugegraph.core.RoleElectionStateMachineTest; import org.apache.hugegraph.meta.EtcdMetaDriverTest; import org.apache.hugegraph.meta.MetaManagerSchemaCacheClearEventTest; import org.apache.hugegraph.meta.managers.AuthMetaManagerTest; @@ -181,7 +180,6 @@ SystemSchemaStoreTest.class, ServerInfoManagerTest.class, TaskSchedulerServerInfoTest.class, - RoleElectionStateMachineTest.class, HugeGraphAuthProxyTest.class, SchemaElementTest.class, ShortestPathTraverserTest.class, diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/TaskSchedulerServerInfoTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/TaskSchedulerServerInfoTest.java index 758b8c1f2d..5acd19d177 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/TaskSchedulerServerInfoTest.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/TaskSchedulerServerInfoTest.java @@ -156,9 +156,11 @@ public void testGraphManagerDoesNotGenerateServerIdWhenElectionDisabled() { } @Test - public void testGraphManagerWarnsOnRoleElection() { + public void testGraphManagerIgnoresRemovedRoleElectionOptions() { + // Old rest-server.properties files may still set these keys PropertiesConfiguration conf = new PropertiesConfiguration(); - conf.setProperty(ServerOptions.ENABLE_SERVER_ROLE_ELECTION.name(), true); + conf.setProperty("server.role_election", true); + conf.setProperty("server.role.fail_count", 5); HugeConfig config = new HugeConfig(conf); GraphManager manager = new GraphManager(config, new EventHub("test")); From 5117958f0be8c12bac0eeb1360de2c593c8b03b1 Mon Sep 17 00:00:00 2001 From: Ayan Alam Date: Mon, 28 Sep 2026 01:00:40 +0530 Subject: [PATCH 2/2] refactor(server): drop HugeServerInfo and the server info db executor ServerInfoManager.init() and heartbeat() have been no-ops since #3082, and tx()/call() had no callers, so the server-info-db-worker pool that backed them never ran anything. Remove those methods, the executor, the constructor parameter threaded through the task schedulers, and the HugeServerInfo vertex class. ServerInfoManager itself stays: StandardHugeGraph still hands it the node info, and TaskScheduler.serverManager() is part of the interface. --- .../hugegraph/auth/HugeFactoryAuthProxy.java | 3 +- .../hugegraph/structure/HugeVertex.java | 8 +- .../task/DistributedTaskScheduler.java | 5 +- .../apache/hugegraph/task/HugeServerInfo.java | 317 ------------------ .../hugegraph/task/ServerInfoManager.java | 40 +-- .../hugegraph/task/StandardTaskScheduler.java | 5 +- .../task/TaskAndResultScheduler.java | 7 +- .../apache/hugegraph/task/TaskManager.java | 19 +- .../task/TaskAndResultSchedulerTest.java | 10 +- .../unit/core/ServerInfoManagerTest.java | 22 +- .../core/TaskSchedulerServerInfoTest.java | 39 ++- 11 files changed, 47 insertions(+), 428 deletions(-) delete mode 100644 hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/HugeServerInfo.java diff --git a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/HugeFactoryAuthProxy.java b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/HugeFactoryAuthProxy.java index 83f0a64f97..bb954446db 100644 --- a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/HugeFactoryAuthProxy.java +++ b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/auth/HugeFactoryAuthProxy.java @@ -400,8 +400,7 @@ private static void registerPrivateActions() { "oneNumericField", "hasSameProperties"); Reflection.registerFieldsToFilter(TaskManager.class, "LOG", "SCHEDULE_PERIOD", "THREADS", "MANAGER", "schedulers", "taskExecutor", "taskDbExecutor", - "serverInfoDbExecutor", "contexts", - "$assertionsDisabled"); + "contexts", "$assertionsDisabled"); Reflection.registerMethodsToFilter(TaskManager.class, "lambda$0", "resetContext", "closeTaskTx", "setContext", "instance", "notifyNewTask", diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/structure/HugeVertex.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/structure/HugeVertex.java index ae1a45b845..573e2f948e 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/structure/HugeVertex.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/structure/HugeVertex.java @@ -43,7 +43,6 @@ import org.apache.hugegraph.schema.EdgeLabel; import org.apache.hugegraph.schema.PropertyKey; import org.apache.hugegraph.schema.VertexLabel; -import org.apache.hugegraph.task.HugeServerInfo; import org.apache.hugegraph.task.HugeTask; import org.apache.hugegraph.task.HugeTaskResult; import org.apache.hugegraph.type.HugeType; @@ -72,9 +71,10 @@ public class HugeVertex extends HugeElement implements Vertex, Cloneable { private static final List EMPTY_LIST = ImmutableList.of(); /* - * Labels of the removed server info and role election vertices, graphs - * created by older versions may still store them + * Labels of the removed server info and role election vertices. Graphs + * created by older versions may still store them. */ + private static final String LEGACY_SERVER_LABEL = Graph.Hidden.hide("server"); private static final String LEGACY_ROLE_DATA_LABEL = Graph.Hidden.hide("role_data"); private Id id; @@ -106,7 +106,7 @@ public HugeType type() { return HugeType.TASK; } if (label != null && - (label.name().equals(HugeServerInfo.P.SERVER) || + (label.name().equals(LEGACY_SERVER_LABEL) || label.name().equals(LEGACY_ROLE_DATA_LABEL))) { return HugeType.SERVER; } diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/DistributedTaskScheduler.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/DistributedTaskScheduler.java index 57885a9d99..4e9e0b13a4 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/DistributedTaskScheduler.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/DistributedTaskScheduler.java @@ -75,9 +75,8 @@ public DistributedTaskScheduler(HugeGraphParams graph, ExecutorService schemaTaskExecutor, ExecutorService olapTaskExecutor, ExecutorService gremlinTaskExecutor, - ExecutorService ephemeralTaskExecutor, - ExecutorService serverInfoDbExecutor) { - super(graph, serverInfoDbExecutor); + ExecutorService ephemeralTaskExecutor) { + super(graph); this.taskDbExecutor = taskDbExecutor; this.schemaTaskExecutor = schemaTaskExecutor; diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/HugeServerInfo.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/HugeServerInfo.java deleted file mode 100644 index f0485f6656..0000000000 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/HugeServerInfo.java +++ /dev/null @@ -1,317 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.hugegraph.task; - -import java.util.ArrayList; -import java.util.Date; -import java.util.HashMap; -import java.util.List; -import java.util.Map; - -import org.apache.hugegraph.HugeGraph; -import org.apache.hugegraph.HugeGraphParams; -import org.apache.hugegraph.backend.id.Id; -import org.apache.hugegraph.backend.id.IdGenerator; -import org.apache.hugegraph.schema.IndexLabel; -import org.apache.hugegraph.schema.PropertyKey; -import org.apache.hugegraph.schema.SchemaManager; -import org.apache.hugegraph.schema.VertexLabel; -import org.apache.hugegraph.type.HugeType; -import org.apache.hugegraph.type.define.Cardinality; -import org.apache.hugegraph.type.define.DataType; -import org.apache.hugegraph.type.define.NodeRole; -import org.apache.hugegraph.type.define.SerialEnum; -import org.apache.hugegraph.util.DateUtil; -import org.apache.hugegraph.util.E; -import org.apache.tinkerpop.gremlin.structure.Graph; -import org.apache.tinkerpop.gremlin.structure.T; -import org.apache.tinkerpop.gremlin.structure.Vertex; -import org.apache.tinkerpop.gremlin.structure.VertexProperty; - -public class HugeServerInfo { - - // Unit millisecond - private static final long EXPIRED_INTERVAL = TaskManager.SCHEDULE_PERIOD * 10; - - private NodeRole role; - private Date updateTime; - private int maxLoad; - private int load; - private final Id id; - - private transient boolean updated = false; - - public HugeServerInfo(String name, NodeRole role) { - this(IdGenerator.of(name), role); - } - - public HugeServerInfo(Id id) { - this.id = id; - this.role = NodeRole.WORKER; - this.maxLoad = 0; - this.load = 0; - this.updateTime = DateUtil.now(); - } - - public HugeServerInfo(Id id, NodeRole role) { - this.id = id; - this.load = 0; - this.role = role; - this.updateTime = DateUtil.now(); - } - - public Id id() { - return this.id; - } - - public String name() { - return this.id.asString(); - } - - public NodeRole role() { - return this.role; - } - - public void role(NodeRole role) { - this.role = role; - } - - public int maxLoad() { - return this.maxLoad; - } - - public void maxLoad(int maxLoad) { - this.maxLoad = maxLoad; - } - - public int load() { - return this.load; - } - - public void load(int load) { - this.load = load; - } - - public void increaseLoad(int delta) { - this.load += delta; - this.updated = true; - } - - public long expireTime() { - return this.updateTime.getTime() + EXPIRED_INTERVAL; - } - - public Date updateTime() { - return this.updateTime; - } - - public void updateTime(Date updateTime) { - this.updateTime = updateTime; - } - - public boolean alive() { - long now = DateUtil.now().getTime(); - return this.updateTime != null && - this.updateTime.getTime() + EXPIRED_INTERVAL > now; - } - - public boolean updated() { - return this.updated; - } - - @Override - public String toString() { - return String.format("HugeServerInfo(%s)%s", this.id, this.asMap()); - } - - protected boolean property(String key, Object value) { - switch (key) { - case P.ROLE: - this.role = SerialEnum.fromCode(NodeRole.class, (byte) value); - break; - case P.MAX_LOAD: - this.maxLoad = (int) value; - break; - case P.LOAD: - this.load = (int) value; - break; - case P.UPDATE_TIME: - this.updateTime = (Date) value; - break; - default: - throw new AssertionError("Unsupported key: " + key); - } - return true; - } - - protected Object[] asArray() { - E.checkState(this.id != null, "Server id can't be null"); - - List list = new ArrayList<>(12); - - list.add(T.label); - list.add(P.SERVER); - - list.add(T.id); - list.add(this.id); - - list.add(P.ROLE); - list.add(this.role.code()); - - list.add(P.MAX_LOAD); - list.add(this.maxLoad); - - list.add(P.LOAD); - list.add(this.load); - - list.add(P.UPDATE_TIME); - list.add(this.updateTime); - - return list.toArray(); - } - - public Map asMap() { - E.checkState(this.id != null, "Server id can't be null"); - - Map map = new HashMap<>(); - - map.put(Graph.Hidden.unHide(P.ID), this.id); - map.put(Graph.Hidden.unHide(P.LABEL), P.SERVER); - map.put(Graph.Hidden.unHide(P.ROLE), this.role); - map.put(Graph.Hidden.unHide(P.MAX_LOAD), this.maxLoad); - map.put(Graph.Hidden.unHide(P.LOAD), this.load); - map.put(Graph.Hidden.unHide(P.UPDATE_TIME), this.updateTime); - - return map; - } - - public static HugeServerInfo fromVertex(Vertex vertex) { - HugeServerInfo serverInfo = new HugeServerInfo((Id) vertex.id()); - for (var iter = vertex.properties(); iter.hasNext(); ) { - VertexProperty prop = iter.next(); - serverInfo.property(prop.key(), prop.value()); - } - return serverInfo; - } - - public static Schema schema(HugeGraphParams graph) { - return new Schema(graph); - } - - public static final class P { - - public static final String SERVER = Graph.Hidden.hide("server"); - - public static final String ID = T.id.getAccessor(); - public static final String LABEL = T.label.getAccessor(); - - public static final String NAME = "~server_name"; - public static final String ROLE = "~server_role"; - public static final String LOAD = "~server_load"; - public static final String MAX_LOAD = "~server_max_load"; - public static final String UPDATE_TIME = "~server_update_time"; - - public static String unhide(String key) { - final String prefix = Graph.Hidden.hide("server_"); - if (key.startsWith(prefix)) { - return key.substring(prefix.length()); - } - return key; - } - } - - public static final class Schema { - - public static final String SERVER = P.SERVER; - - private final HugeGraphParams graph; - - public Schema(HugeGraphParams graph) { - this.graph = graph; - } - - public void initSchemaIfNeeded() { - if (this.existVertexLabel(SERVER)) { - return; - } - - HugeGraph graph = this.graph.graph(); - String[] properties = this.initProperties(); - - // Create vertex label '~server' - VertexLabel label = graph.schema().vertexLabel(SERVER) - .properties(properties) - .useCustomizeStringId() - .nullableKeys(P.ROLE, P.MAX_LOAD, P.LOAD, P.UPDATE_TIME) - .enableLabelIndex(true) - .build(); - this.graph.schemaTransaction().addVertexLabel(label); - } - - private String[] initProperties() { - List props = new ArrayList<>(); - props.add(createPropertyKey(P.ROLE, DataType.BYTE)); - props.add(createPropertyKey(P.MAX_LOAD, DataType.INT)); - props.add(createPropertyKey(P.LOAD, DataType.INT)); - props.add(createPropertyKey(P.UPDATE_TIME, DataType.DATE)); - - return props.toArray(new String[0]); - } - - public boolean existVertexLabel(String label) { - return this.graph.schemaTransaction().getVertexLabel(label) != null; - } - - @SuppressWarnings("unused") - private String createPropertyKey(String name) { - return this.createPropertyKey(name, DataType.TEXT); - } - - private String createPropertyKey(String name, DataType dataType) { - return this.createPropertyKey(name, dataType, Cardinality.SINGLE); - } - - private String createPropertyKey(String name, DataType dataType, Cardinality cardinality) { - SchemaManager schema = this.graph.graph().schema(); - PropertyKey propertyKey = schema.propertyKey(name) - .dataType(dataType) - .cardinality(cardinality) - .build(); - this.graph.schemaTransaction().addPropertyKey(propertyKey); - return name; - } - - @SuppressWarnings("unused") - private IndexLabel createIndexLabel(VertexLabel label, String field) { - SchemaManager schema = this.graph.graph().schema(); - String name = Graph.Hidden.hide("server-index-by-" + field); - IndexLabel indexLabel = schema.indexLabel(name) - .on(HugeType.VERTEX_LABEL, SERVER) - .by(field) - .build(); - this.graph.schemaTransaction().addIndexLabel(label, indexLabel); - return indexLabel; - } - - @SuppressWarnings("unused") - private IndexLabel indexLabel(String field) { - String name = Graph.Hidden.hide("server-index-by-" + field); - return this.graph.graph().indexLabel(name); - } - } -} diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/ServerInfoManager.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/ServerInfoManager.java index dac78f2fcc..3d3585caf4 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/ServerInfoManager.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/ServerInfoManager.java @@ -17,15 +17,9 @@ package org.apache.hugegraph.task; -import java.util.concurrent.Callable; -import java.util.concurrent.ExecutorService; - -import org.apache.hugegraph.HugeException; import org.apache.hugegraph.HugeGraphParams; import org.apache.hugegraph.backend.id.Id; import org.apache.hugegraph.backend.id.IdGenerator; -import org.apache.hugegraph.backend.tx.GraphTransaction; -import org.apache.hugegraph.exception.ConnectionException; import org.apache.hugegraph.masterelection.GlobalMasterInfo; import org.apache.hugegraph.type.define.NodeRole; import org.apache.hugegraph.util.E; @@ -33,31 +27,22 @@ public class ServerInfoManager { private final HugeGraphParams graph; - private final ExecutorService dbExecutor; private volatile GlobalMasterInfo globalNodeInfo; private volatile boolean closed; - public ServerInfoManager(HugeGraphParams graph, ExecutorService dbExecutor) { + public ServerInfoManager(HugeGraphParams graph) { E.checkNotNull(graph, "graph"); - E.checkNotNull(dbExecutor, "db executor"); this.graph = graph; - this.dbExecutor = dbExecutor; this.globalNodeInfo = null; this.closed = false; } - public void init() { - // ServerInfo is soft-disabled; keep this method for compatibility. - } - public synchronized boolean close() { - // ServerInfo persistence is soft-deprecated; init() and heartbeat() - // are no-ops, so there's nothing to clean up in close(). this.closed = true; return true; } @@ -103,27 +88,4 @@ public NodeRole selfNodeRole() { public boolean selfIsMaster() { return this.selfNodeRole() != null && this.selfNodeRole().master(); } - - public synchronized void heartbeat() { - // ServerInfo heartbeat is deprecated for local scheduling. - } - - private GraphTransaction tx() { - assert Thread.currentThread().getName().contains("server-info-db-worker"); - return this.graph.systemTransaction(); - } - - private V call(Callable callable) { - assert !Thread.currentThread().getName().startsWith( - "server-info-db-worker") : "can't call by itself"; - try { - // Pass context for db thread - callable = new TaskManager.ContextCallable<>(callable); - // Ensure all db operations are executed in dbExecutor thread(s) - return this.dbExecutor.submit(callable).get(); - } catch (Throwable e) { - throw new HugeException("Failed to update/query server info: %s", - e, e.toString()); - } - } } diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/StandardTaskScheduler.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/StandardTaskScheduler.java index 269094bc7e..3367572580 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/StandardTaskScheduler.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/StandardTaskScheduler.java @@ -74,8 +74,7 @@ public class StandardTaskScheduler implements TaskScheduler { public StandardTaskScheduler(HugeGraphParams graph, ExecutorService taskExecutor, - ExecutorService taskDbExecutor, - ExecutorService serverInfoDbExecutor) { + ExecutorService taskDbExecutor) { E.checkNotNull(graph, "graph"); E.checkNotNull(taskExecutor, "taskExecutor"); E.checkNotNull(taskDbExecutor, "dbExecutor"); @@ -84,7 +83,7 @@ public StandardTaskScheduler(HugeGraphParams graph, this.taskExecutor = taskExecutor; this.taskDbExecutor = taskDbExecutor; - this.serverManager = new ServerInfoManager(graph, serverInfoDbExecutor); + this.serverManager = new ServerInfoManager(graph); this.tasks = new ConcurrentHashMap<>(); this.taskTx = null; diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/TaskAndResultScheduler.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/TaskAndResultScheduler.java index 70d4d2f571..ffb54c71b6 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/TaskAndResultScheduler.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/TaskAndResultScheduler.java @@ -21,7 +21,6 @@ import java.util.Iterator; import java.util.List; import java.util.Map; -import java.util.concurrent.ExecutorService; import org.apache.hugegraph.HugeGraphParams; import org.apache.hugegraph.backend.id.Id; @@ -61,16 +60,14 @@ public abstract class TaskAndResultScheduler implements TaskScheduler { private final ServerInfoManager serverManager; - public TaskAndResultScheduler( - HugeGraphParams graph, - ExecutorService serverInfoDbExecutor) { + public TaskAndResultScheduler(HugeGraphParams graph) { E.checkNotNull(graph, "graph"); this.graph = graph; this.graphSpace = graph.graph().graphSpace(); this.graphName = graph.name(); - this.serverManager = new ServerInfoManager(graph, serverInfoDbExecutor); + this.serverManager = new ServerInfoManager(graph); } @Override diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/TaskManager.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/TaskManager.java index 135fb47a42..cbdc1948c2 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/TaskManager.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/task/TaskManager.java @@ -45,7 +45,6 @@ public final class TaskManager { public static final String TASK_WORKER_PREFIX = "task-worker"; public static final String TASK_WORKER = TASK_WORKER_PREFIX + "-%d"; public static final String TASK_DB_WORKER = "task-db-worker-%d"; - public static final String SERVER_INFO_DB_WORKER = "server-info-db-worker-%d"; public static final String OLAP_TASK_WORKER = "olap-task-worker-%d"; public static final String SCHEMA_TASK_WORKER = "schema-task-worker-%d"; @@ -61,7 +60,6 @@ public final class TaskManager { private final ExecutorService taskExecutor; private final ExecutorService taskDbExecutor; - private final ExecutorService serverInfoDbExecutor; private final ExecutorService schemaTaskExecutor; private final ExecutorService olapTaskExecutor; @@ -80,8 +78,6 @@ private TaskManager(int pool) { // For save/query task state, just one thread is ok this.taskDbExecutor = ExecutorUtil.newFixedThreadPool( 1, TASK_DB_WORKER); - this.serverInfoDbExecutor = ExecutorUtil.newFixedThreadPool( - 1, SERVER_INFO_DB_WORKER); this.schemaTaskExecutor = ExecutorUtil.newFixedThreadPool(pool, SCHEMA_TASK_WORKER); this.olapTaskExecutor = ExecutorUtil.newFixedThreadPool(pool, OLAP_TASK_WORKER); @@ -107,8 +103,7 @@ public void addScheduler(HugeGraphParams graph) { schemaTaskExecutor, olapTaskExecutor, taskExecutor, /* gremlinTaskExecutor */ - ephemeralTaskExecutor, - serverInfoDbExecutor); + ephemeralTaskExecutor); this.schedulers.put(graph, scheduler); break; } @@ -118,8 +113,7 @@ public void addScheduler(HugeGraphParams graph) { new StandardTaskScheduler( graph, this.taskExecutor, - this.taskDbExecutor, - this.serverInfoDbExecutor); + this.taskDbExecutor); this.schedulers.put(graph, scheduler); break; } @@ -229,15 +223,6 @@ public void shutdown(long timeout) { } } - if (terminated && !this.serverInfoDbExecutor.isShutdown()) { - this.serverInfoDbExecutor.shutdown(); - try { - terminated = this.serverInfoDbExecutor.awaitTermination(timeout, unit); - } catch (Throwable e) { - ex = e; - } - } - if (terminated && !this.taskDbExecutor.isShutdown()) { this.taskDbExecutor.shutdown(); try { diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/task/TaskAndResultSchedulerTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/task/TaskAndResultSchedulerTest.java index 15fabf1302..a31c70e8d1 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/task/TaskAndResultSchedulerTest.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/task/TaskAndResultSchedulerTest.java @@ -209,7 +209,6 @@ private static class TestDistributedTaskScheduler private final ExecutorService olapTaskExecutor; private final ExecutorService gremlinTaskExecutor; private final ExecutorService ephemeralTaskExecutor; - private final ExecutorService serverInfoDbExecutor; private final AtomicInteger resultReadCount; private volatile boolean forbidResultRead; private volatile boolean failNextTaskResultDelete; @@ -220,7 +219,6 @@ private static class TestDistributedTaskScheduler Executors.newSingleThreadExecutor(), Executors.newSingleThreadExecutor(), Executors.newSingleThreadExecutor(), - Executors.newSingleThreadExecutor(), Executors.newSingleThreadExecutor()); } @@ -231,18 +229,15 @@ private TestDistributedTaskScheduler( ExecutorService schemaTaskExecutor, ExecutorService olapTaskExecutor, ExecutorService gremlinTaskExecutor, - ExecutorService ephemeralTaskExecutor, - ExecutorService serverInfoDbExecutor) { + ExecutorService ephemeralTaskExecutor) { super(graph, schedulerExecutor, taskDbExecutor, schemaTaskExecutor, - olapTaskExecutor, gremlinTaskExecutor, ephemeralTaskExecutor, - serverInfoDbExecutor); + olapTaskExecutor, gremlinTaskExecutor, ephemeralTaskExecutor); this.schedulerExecutor = schedulerExecutor; this.taskDbExecutor = taskDbExecutor; this.schemaTaskExecutor = schemaTaskExecutor; this.olapTaskExecutor = olapTaskExecutor; this.gremlinTaskExecutor = gremlinTaskExecutor; this.ephemeralTaskExecutor = ephemeralTaskExecutor; - this.serverInfoDbExecutor = serverInfoDbExecutor; this.resultReadCount = new AtomicInteger(); this.forbidResultRead = false; this.failNextTaskResultDelete = false; @@ -327,7 +322,6 @@ public void closeAndShutdown() { this.olapTaskExecutor.shutdownNow(); this.gremlinTaskExecutor.shutdownNow(); this.ephemeralTaskExecutor.shutdownNow(); - this.serverInfoDbExecutor.shutdownNow(); } } diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/ServerInfoManagerTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/ServerInfoManagerTest.java index 607d81de91..3784cef086 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/ServerInfoManagerTest.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/ServerInfoManagerTest.java @@ -17,8 +17,6 @@ package org.apache.hugegraph.unit.core; -import java.util.concurrent.ExecutorService; - import org.apache.hugegraph.HugeGraphParams; import org.apache.hugegraph.backend.id.Id; import org.apache.hugegraph.masterelection.GlobalMasterInfo; @@ -44,24 +42,8 @@ public void setup() { Mockito.when(hugegraphParams.spaceGraphName()) .thenReturn("DEFAULT-hugegraph"); - ExecutorService executor = Mockito.mock(ExecutorService.class); - - this.sysGraphManager = new ServerInfoManager(sysGraphParams, executor); - this.hugegraphManager = new ServerInfoManager(hugegraphParams, executor); - } - - @Test - public void testInitDoesNotAccessBackendStore() { - HugeGraphParams graphParams = Mockito.mock(HugeGraphParams.class); - ExecutorService executor = Mockito.mock(ExecutorService.class); - ServerInfoManager manager = new ServerInfoManager(graphParams, executor); - - manager.init(); - - Mockito.verify(graphParams, Mockito.never()).systemTransaction(); - Mockito.verify(graphParams, Mockito.never()).backendStoreFeatures(); - Mockito.verify(graphParams, Mockito.never()).graph(); - Mockito.verify(graphParams, Mockito.never()).closeTx(); + this.sysGraphManager = new ServerInfoManager(sysGraphParams); + this.hugegraphManager = new ServerInfoManager(hugegraphParams); } @Test diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/TaskSchedulerServerInfoTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/TaskSchedulerServerInfoTest.java index 5acd19d177..f7dd3e370a 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/TaskSchedulerServerInfoTest.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/TaskSchedulerServerInfoTest.java @@ -32,12 +32,18 @@ import org.apache.hugegraph.core.GraphManager; import org.apache.hugegraph.event.EventHub; import org.apache.hugegraph.exception.NotFoundException; +import org.apache.hugegraph.schema.VertexLabel; +import org.apache.hugegraph.structure.HugeVertex; import org.apache.hugegraph.task.DistributedTaskScheduler; import org.apache.hugegraph.task.HugeTask; import org.apache.hugegraph.task.TaskCallable; import org.apache.hugegraph.task.TaskStatus; import org.apache.hugegraph.testutil.Assert; import org.apache.hugegraph.testutil.Whitebox; +import org.apache.hugegraph.type.HugeType; +import org.apache.hugegraph.type.define.IdStrategy; +import org.apache.hugegraph.type.define.NodeRole; +import org.apache.hugegraph.unit.FakeObjects; import org.apache.hugegraph.util.ExecutorUtil; import org.junit.Test; import org.mockito.Mockito; @@ -63,13 +69,11 @@ public void testDistributedCheckRequirementDoesNotNeedServerInfo() { ExecutorService olapTaskExecutor = Executors.newSingleThreadExecutor(); ExecutorService gremlinTaskExecutor = Executors.newSingleThreadExecutor(); ExecutorService ephemeralTaskExecutor = Executors.newSingleThreadExecutor(); - ExecutorService serverInfoDbExecutor = Executors.newSingleThreadExecutor(); try { DistributedTaskScheduler scheduler = new DistributedTaskScheduler( params, schedulerExecutor, taskDbExecutor, schemaTaskExecutor, - olapTaskExecutor, gremlinTaskExecutor, ephemeralTaskExecutor, - serverInfoDbExecutor); + olapTaskExecutor, gremlinTaskExecutor, ephemeralTaskExecutor); scheduler.checkRequirement("schedule"); } finally { schedulerExecutor.shutdownNow(); @@ -78,7 +82,6 @@ public void testDistributedCheckRequirementDoesNotNeedServerInfo() { olapTaskExecutor.shutdownNow(); gremlinTaskExecutor.shutdownNow(); ephemeralTaskExecutor.shutdownNow(); - serverInfoDbExecutor.shutdownNow(); } } @@ -101,13 +104,11 @@ public void testDistributedCancelTreatsMissingRetryTaskAsGone() { ExecutorService olapTaskExecutor = Executors.newSingleThreadExecutor(); ExecutorService gremlinTaskExecutor = Executors.newSingleThreadExecutor(); ExecutorService ephemeralTaskExecutor = Executors.newSingleThreadExecutor(); - ExecutorService serverInfoDbExecutor = Executors.newSingleThreadExecutor(); try { DistributedTaskScheduler scheduler = new DistributedTaskScheduler( params, schedulerExecutor, taskDbExecutor, schemaTaskExecutor, - olapTaskExecutor, gremlinTaskExecutor, ephemeralTaskExecutor, - serverInfoDbExecutor) { + olapTaskExecutor, gremlinTaskExecutor, ephemeralTaskExecutor) { @Override protected boolean updateStatus(Id id, TaskStatus prestatus, @@ -138,12 +139,11 @@ public Object call() { olapTaskExecutor.shutdownNow(); gremlinTaskExecutor.shutdownNow(); ephemeralTaskExecutor.shutdownNow(); - serverInfoDbExecutor.shutdownNow(); } } @Test - public void testGraphManagerDoesNotGenerateServerIdWhenElectionDisabled() { + public void testGraphManagerDoesNotGenerateServerId() { HugeConfig config = newConfig(); GraphManager manager = new GraphManager(config, new EventHub("test")); @@ -161,16 +161,35 @@ public void testGraphManagerIgnoresRemovedRoleElectionOptions() { PropertiesConfiguration conf = new PropertiesConfiguration(); conf.setProperty("server.role_election", true); conf.setProperty("server.role.fail_count", 5); + conf.setProperty(ServerOptions.SERVER_ROLE.name(), "worker"); HugeConfig config = new HugeConfig(conf); GraphManager manager = new GraphManager(config, new EventHub("test")); try { - Assert.assertNotNull(manager); + Assert.assertEquals(NodeRole.WORKER, + manager.globalNodeRoleInfo().nodeRole()); } finally { manager.close(); } } + @Test + public void testLegacyServerLabelsKeepServerType() { + // Graphs created by older versions may still store these vertices + FakeObjects fakeObject = new FakeObjects(); + for (String name : new String[]{"~server", "~role_data"}) { + VertexLabel label = fakeObject.newVertexLabel(IdGenerator.of(name), name, + IdStrategy.CUSTOMIZE_STRING); + HugeVertex vertex = new HugeVertex(fakeObject.graph(), null, label); + Assert.assertEquals(HugeType.SERVER, vertex.type()); + } + + VertexLabel person = fakeObject.newVertexLabel(IdGenerator.of(1), "person", + IdStrategy.CUSTOMIZE_STRING); + HugeVertex vertex = new HugeVertex(fakeObject.graph(), null, person); + Assert.assertEquals(HugeType.VERTEX, vertex.type()); + } + @Test public void testGraphManagerDoesNotInjectPdPeersForStandaloneRocksDB() { HugeConfig serverConfig = newConfig();