Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -400,11 +400,10 @@ private static void registerPrivateActions() {
"oneNumericField", "hasSameProperties");
Reflection.registerFieldsToFilter(TaskManager.class, "LOG", "SCHEDULE_PERIOD", "THREADS",
"MANAGER", "schedulers", "taskExecutor", "taskDbExecutor",
"serverInfoDbExecutor", "schedulerExecutor", "contexts",
"$assertionsDisabled");
"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",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);

Expand All @@ -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
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,15 +42,6 @@ public class ServerOptions extends OptionHolder {
1
);

public static final ConfigOption<Boolean> 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<Integer> MAX_WORKER_THREADS =
new ConfigOption<>(
"restserver.max_worker_threads",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -189,7 +186,6 @@ public final class GraphManager {
private final Set<String> serverUrlsToPd;
private final Boolean serverDeployInK8s;
private final HugeConfig config;
private RoleElectionStateMachine roleStateMachine;
private K8sDriver.CA ca;
private final boolean PDExist;

Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -696,7 +691,7 @@ public void init() {
this.waitGraphsReady();

this.checkBackendVersionOrExit(this.conf);
this.serverStarted(this.conf);
this.serverStarted();

this.addMetrics(this.conf);
}
Expand Down Expand Up @@ -1735,9 +1730,6 @@ public void close() {
}
this.destroyRpcServer();
this.unlistenChanges();
if (this.roleStateMachine != null) {
this.roleStateMachine.shutdown();
}
}

private void startRpcServer() {
Expand Down Expand Up @@ -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);

Expand Down Expand Up @@ -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");
Expand Down Expand Up @@ -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");
Expand All @@ -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,
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -273,8 +272,6 @@ public interface HugeGraph extends Graph {

AuthManager authManager();

RoleElectionStateMachine roleElectionStateMachine();

void switchAuthManager(AuthManager authManager);
Comment thread
byteayan marked this conversation as resolved.

TaskScheduler taskScheduler();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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?
Expand All @@ -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;
Expand Down Expand Up @@ -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;
Expand Down
Loading