hzh0425 commented on code in PR #4195:
URL: https://github.com/apache/rocketmq/pull/4195#discussion_r857059424
##########
namesrv/src/main/java/org/apache/rocketmq/namesrv/NamesrvController.java:
##########
@@ -82,24 +85,36 @@ public class NamesrvController {
private ExecutorService defaultExecutor;
private ExecutorService clientRequestExecutor;
+ private ExecutorService controllerRequestExecutor;
private BlockingQueue<Runnable> defaultThreadPoolQueue;
private BlockingQueue<Runnable> clientRequestThreadPoolQueue;
+ private BlockingQueue<Runnable> controllerRequestThreadPoolQueue;
private Configuration configuration;
private FileWatchService fileWatchService;
+ private Controller controller;
+
public NamesrvController(NamesrvConfig namesrvConfig, NettyServerConfig
nettyServerConfig) {
this(namesrvConfig, nettyServerConfig, new NettyClientConfig());
}
- public NamesrvController(NamesrvConfig namesrvConfig, NettyServerConfig
nettyServerConfig, NettyClientConfig nettyClientConfig) {
+ public NamesrvController(NamesrvConfig namesrvConfig, NettyServerConfig
nettyServerConfig,
+ NettyClientConfig nettyClientConfig) {
this.namesrvConfig = namesrvConfig;
this.nettyServerConfig = nettyServerConfig;
this.nettyClientConfig = nettyClientConfig;
this.kvConfigManager = new KVConfigManager(this);
this.brokerHousekeepingService = new BrokerHousekeepingService(this);
- this.routeInfoManager = new RouteInfoManager(namesrvConfig, this);
+ if (namesrvConfig.isStartupController()) {
+ final DLedgerConfig config = new DLedgerConfig();
+ config.setGroup(namesrvConfig.getControllerDLegerGroup());
+ config.setPeers(namesrvConfig.getControllerDLegerPeers());
+ config.setSelfId(namesrvConfig.getControllerDLegerSelfId());
+ this.controller = new DledgerController(config,
namesrvConfig.isEnableElectUncleanMaster());
Review Comment:
讨论方案: 新增一个 DledgerControllerConfig
##########
namesrv/src/main/java/org/apache/rocketmq/namesrv/controller/impl/DledgerController.java:
##########
@@ -0,0 +1,353 @@
+/*
+ * 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.rocketmq.namesrv.controller.impl;
+
+import io.openmessaging.storage.dledger.AppendFuture;
+import io.openmessaging.storage.dledger.DLedgerConfig;
+import io.openmessaging.storage.dledger.DLedgerLeaderElector;
+import io.openmessaging.storage.dledger.DLedgerServer;
+import io.openmessaging.storage.dledger.MemberState;
+import io.openmessaging.storage.dledger.protocol.AppendEntryRequest;
+import io.openmessaging.storage.dledger.protocol.AppendEntryResponse;
+import io.openmessaging.storage.dledger.protocol.BatchAppendEntryRequest;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.TimeUnit;
+import java.util.function.Supplier;
+import org.apache.rocketmq.common.constant.LoggerName;
+import
org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
+import
org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetResponseHeader;
+import
org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterRequestHeader;
+import
org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterResponseHeader;
+import
org.apache.rocketmq.common.protocol.header.namesrv.controller.GetMetaDataResponseHeader;
+import
org.apache.rocketmq.common.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
+import
org.apache.rocketmq.common.protocol.header.namesrv.controller.GetReplicaInfoResponseHeader;
+import
org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerRequestHeader;
+import
org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerResponseHeader;
+import org.apache.rocketmq.logging.InternalLogger;
+import org.apache.rocketmq.logging.InternalLoggerFactory;
+import org.apache.rocketmq.namesrv.controller.Controller;
+import org.apache.rocketmq.namesrv.controller.manager.ReplicasInfoManager;
+import org.apache.rocketmq.namesrv.controller.manager.event.ControllerResult;
+import org.apache.rocketmq.namesrv.controller.manager.event.EventMessage;
+import org.apache.rocketmq.namesrv.controller.manager.event.EventSerializer;
+import org.apache.rocketmq.remoting.common.ServiceThread;
+
+/**
+ * The implementation of controller, based on dledger (raft).
+ */
+public class DledgerController implements Controller {
+
+ private static final InternalLogger log =
InternalLoggerFactory.getLogger(LoggerName.CONTROLLER_LOGGER_NAME);
+ private final DLedgerServer dLedgerServer;
+ private final DLedgerConfig dLedgerConfig;
+ private final ReplicasInfoManager replicasInfoManager;
+ private final EventScheduler scheduler;
+ private final EventSerializer eventSerializer;
+ private final RoleChangeHandler roleHandler;
+ private final DledgerControllerStateMachine statemachine;
+ private volatile boolean isScheduling = false;
+
+ public DledgerController(final DLedgerConfig dLedgerConfig, final boolean
isEnableElectUncleanMaster) {
+ this.dLedgerConfig = dLedgerConfig;
+
+ this.eventSerializer = new EventSerializer();
+
+ this.scheduler = new EventScheduler();
+ this.roleHandler = new RoleChangeHandler(dLedgerConfig.getSelfId());
+ this.replicasInfoManager = new
ReplicasInfoManager(isEnableElectUncleanMaster);
+ this.statemachine = new
DledgerControllerStateMachine(replicasInfoManager, this.eventSerializer,
dLedgerConfig.getSelfId());
+
+ // Register statemachine and role handler.
+ this.dLedgerServer = new DLedgerServer(dLedgerConfig);
+ this.dLedgerServer.registerStateMachine(this.statemachine);
+
this.dLedgerServer.getdLedgerLeaderElector().addRoleChangeHandler(this.roleHandler);
+ }
+
+ @Override
+ public void startup() {
+ this.dLedgerServer.startup();
+ }
+
+ @Override
+ public void shutdown() {
+ this.dLedgerServer.shutdown();
+ }
+
+ @Override
+ public void startScheduling() {
+ if (!this.isScheduling) {
+ log.info("Start scheduling controller events");
+ this.isScheduling = true;
+ this.scheduler.start();
+ }
+ }
+
+ @Override
+ public void stopScheduling() {
+ if (this.isScheduling) {
+ log.info("Stop scheduling controller events");
+ this.isScheduling = false;
+ this.scheduler.shutdown(true);
+ }
+ }
+
+ @Override
+ public CompletableFuture<AlterSyncStateSetResponseHeader>
alterSyncStateSet(
+ AlterSyncStateSetRequestHeader request) {
+ if (!this.roleHandler.isLeaderState()) {
+ log.warn("Current controller {} is not leader, reject
alterSyncStateSet request", this.dLedgerConfig.getSelfId());
+ return null;
+ }
+ return this.scheduler.appendEvent("alterSyncStateSet",
+ () -> this.replicasInfoManager.alterSyncStateSet(request), true);
+ }
+
+ @Override
+ public CompletableFuture<ElectMasterResponseHeader> electMaster(final
ElectMasterRequestHeader request) {
+ if (!this.roleHandler.isLeaderState()) {
+ log.warn("Current controller {} is not leader, reject electMaster
request", this.dLedgerConfig.getSelfId());
+ return null;
+ }
+ return this.scheduler.appendEvent("electMaster",
+ () -> this.replicasInfoManager.electMaster(request), true);
+ }
+
+ @Override
+ public CompletableFuture<RegisterBrokerResponseHeader>
registerBroker(RegisterBrokerRequestHeader request) {
+ if (!this.roleHandler.isLeaderState()) {
+ log.warn("Current controller {} is not leader, reject
registerBroker request", this.dLedgerConfig.getSelfId());
+ return null;
+ }
+ return this.scheduler.appendEvent("registerBroker",
+ () -> this.replicasInfoManager.registerBroker(request), true);
+ }
+
+ @Override
+ public CompletableFuture<GetReplicaInfoResponseHeader>
getReplicaInfo(final GetReplicaInfoRequestHeader request) {
+ if (!this.roleHandler.isLeaderState()) {
+ log.warn("Current controller {} is not leader, reject
getReplicaInfo request", this.dLedgerConfig.getSelfId());
+ return null;
+ }
+ return this.scheduler.appendEvent("getReplicaInfo",
+ () -> this.replicasInfoManager.getReplicaInfo(request), false);
+ }
+
+ @Override
+ public GetMetaDataResponseHeader getControllerMetadata() {
+ final MemberState state = getMemberState();
+ return new GetMetaDataResponseHeader(state.getLeaderId(),
state.getLeaderAddr());
+ }
+
+ /**
+ * Event scheduler, schedule event handler from event queue
+ */
+ class EventScheduler extends ServiceThread {
+ private final BlockingQueue<EventHandler> eventQueue;
+
+ public EventScheduler() {
+ this.eventQueue = new LinkedBlockingQueue<>(1024);
+ }
+
+ @Override
+ public String getServiceName() {
+ return EventScheduler.class.getName();
+ }
+
+ @Override
+ public void run() {
+ log.info("Start event scheduler.");
+ while (!isStopped()) {
+ EventHandler handler;
+ try {
+ handler = this.eventQueue.poll(5, TimeUnit.SECONDS);
+ } catch (final InterruptedException e) {
+ continue;
+ }
+ try {
+ if (handler != null) {
+ handler.run();
+ }
+ } catch (final Throwable e) {
+ handler.handleException(e);
+ }
+ }
+
+ }
+
+ public <T> CompletableFuture<T> appendEvent(final String name, final
Supplier<ControllerResult<T>> supplier,
+ boolean isWriteEvent) {
+ if (isStopped()) {
+ return null;
+ }
+ final EventHandler<T> event = new ControllerEventHandler<>(name,
supplier, isWriteEvent);
+ int tryTimes = 0;
+ while (true) {
+ try {
+ if (!this.eventQueue.offer(event, 5, TimeUnit.SECONDS)) {
+ continue;
+ }
+ return event.future();
+ } catch (final InterruptedException e) {
+ log.error("Error happen in EventScheduler when append
event", e);
+ tryTimes++;
+ if (tryTimes > 3) {
+ return null;
+ }
+ }
+ }
+ }
+ }
+
+ /**
+ * Event handler, get events from supplier, and append events to dledger
+ */
+ class ControllerEventHandler<T> implements EventHandler<T> {
+ private final String name;
+ private final Supplier<ControllerResult<T>> supplier;
+ private final CompletableFuture<T> future;
+ private final boolean isWriteEvent;
+
+ ControllerEventHandler(final String name, final
Supplier<ControllerResult<T>> supplier,
+ final boolean isWriteEvent) {
+ this.name = name;
+ this.supplier = supplier;
+ this.future = new CompletableFuture<>();
+ this.isWriteEvent = isWriteEvent;
+ }
+
+ @Override
+ public void run() throws Throwable {
+ final ControllerResult<T> result = this.supplier.get();
+ log.info("Event queue run event {}, get the result {}", this.name,
result);
+ boolean appendSuccess = true;
+ if (this.isWriteEvent) {
+ final List<EventMessage> events = result.getEvents();
+ final List<byte[]> eventBytes = new ArrayList<>(events.size());
+ for (final EventMessage event : events) {
+ if (event != null) {
+ final byte[] data =
DledgerController.this.eventSerializer.serialize(event);
+ if (data != null && data.length > 0) {
+ eventBytes.add(data);
+ }
+ }
+ }
+ // Append events to dledger
+ if (!eventBytes.isEmpty()) {
+ final BatchAppendEntryRequest request = new
BatchAppendEntryRequest();
+ request.setBatchMsgs(eventBytes);
+ appendSuccess = appendToDledgerAndWait(request);
+ }
+ } else {
+ // Now the dledger don't have the function of Read-Index or
Lease-Read,
+ // So we still need to propose an empty request to dledger.
+ final AppendEntryRequest request = new AppendEntryRequest();
+ request.setBody(new byte[0]);
+ appendSuccess = appendToDledgerAndWait(request);
+ }
Review Comment:
讨论方案: 先开一个配置(isProcessReadEvet = {ReadEvent 是否要走日志}). 后续可以为 Dledger 新增
Read-Index
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]