温馨提示×

温馨提示×

您好,登录后才能下订单哦!

密码登录×
登录注册×
其他方式登录
点击 登录注册 即表示同意《亿速云用户服务条款》

nacos RaftCore中signalPublish的原理和使用方法

发布时间:2021-06-24 10:33:48 来源:亿速云 阅读:224 作者:chen 栏目:大数据

这篇文章主要讲解了“nacos RaftCore中signalPublish的原理和使用方法”,文中的讲解内容简单清晰,易于学习与理解,下面请大家跟着小编的思路慢慢深入,一起来研究和学习“nacos RaftCore中signalPublish的原理和使用方法”吧!

本文主要研究一下nacos RaftCore的signalPublish

RaftCore

nacos-1.1.3/naming/src/main/java/com/alibaba/nacos/naming/consistency/persistent/raft/RaftCore.java

@Component public class RaftCore {     public static final String API_VOTE = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/vote";     public static final String API_BEAT = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/beat";     public static final String API_PUB = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/datum";     public static final String API_DEL = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/datum";     public static final String API_GET = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/datum";     public static final String API_ON_PUB = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/datum/commit";     public static final String API_ON_DEL = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/datum/commit";     public static final String API_GET_PEER = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/peer";     private ScheduledExecutorService executor = new ScheduledThreadPoolExecutor(1, new ThreadFactory() {         @Override         public Thread newThread(Runnable r) {             Thread t = new Thread(r);             t.setDaemon(true);             t.setName("com.alibaba.nacos.naming.raft.notifier");             return t;         }     });     public static final Lock OPERATE_LOCK = new ReentrantLock();     public static final int PUBLISH_TERM_INCREASE_COUNT = 100;     private volatile Map<String, List<RecordListener>> listeners = new ConcurrentHashMap<>();     private volatile ConcurrentMap<String, Datum> datums = new ConcurrentHashMap<>();     @Autowired     private RaftPeerSet peers;     @Autowired     private SwitchDomain switchDomain;     @Autowired     private GlobalConfig globalConfig;     @Autowired     private RaftProxy raftProxy;     @Autowired     private RaftStore raftStore;     public volatile Notifier notifier = new Notifier();     private boolean initialized = false;     @PostConstruct     public void init() throws Exception {         Loggers.RAFT.info("initializing Raft sub-system");         executor.submit(notifier);         long start = System.currentTimeMillis();         raftStore.loadDatums(notifier, datums);         setTerm(NumberUtils.toLong(raftStore.loadMeta().getProperty("term"), 0L));         Loggers.RAFT.info("cache loaded, datum count: {}, current term: {}", datums.size(), peers.getTerm());         while (true) {             if (notifier.tasks.size() <= 0) {                 break;             }             Thread.sleep(1000L);         }         initialized = true;         Loggers.RAFT.info("finish to load data from disk, cost: {} ms.", (System.currentTimeMillis() - start));         GlobalExecutor.registerMasterElection(new MasterElection());         GlobalExecutor.registerHeartbeat(new HeartBeat());         Loggers.RAFT.info("timer started: leader timeout ms: {}, heart-beat timeout ms: {}",             GlobalExecutor.LEADER_TIMEOUT_MS, GlobalExecutor.HEARTBEAT_INTERVAL_MS);     }     public Map<String, List<RecordListener>> getListeners() {         return listeners;     }     public void signalPublish(String key, Record value) throws Exception {         if (!isLeader()) {             JSONObject params = new JSONObject();             params.put("key", key);             params.put("value", value);             Map<String, String> parameters = new HashMap<>(1);             parameters.put("key", key);             raftProxy.proxyPostLarge(getLeader().ip, API_PUB, params.toJSONString(), parameters);             return;         }         try {             OPERATE_LOCK.lock();             long start = System.currentTimeMillis();             final Datum datum = new Datum();             datum.key = key;             datum.value = value;             if (getDatum(key) == null) {                 datum.timestamp.set(1L);             } else {                 datum.timestamp.set(getDatum(key).timestamp.incrementAndGet());             }             JSONObject json = new JSONObject();             json.put("datum", datum);             json.put("source", peers.local());             onPublish(datum, peers.local());             final String content = JSON.toJSONString(json);             final CountDownLatch latch = new CountDownLatch(peers.majorityCount());             for (final String server : peers.allServersIncludeMyself()) {                 if (isLeader(server)) {                     latch.countDown();                     continue;                 }                 final String url = buildURL(server, API_ON_PUB);                 HttpClient.asyncHttpPostLarge(url, Arrays.asList("key=" + key), content, new AsyncCompletionHandler<Integer>() {                     @Override                     public Integer onCompleted(Response response) throws Exception {                         if (response.getStatusCode() != HttpURLConnection.HTTP_OK) {                             Loggers.RAFT.warn("[RAFT] failed to publish data to peer, datumId={}, peer={}, http code={}",                                 datum.key, server, response.getStatusCode());                             return 1;                         }                         latch.countDown();                         return 0;                     }                     @Override                     public STATE onContentWriteCompleted() {                         return STATE.CONTINUE;                     }                 });             }             if (!latch.await(UtilsAndCommons.RAFT_PUBLISH_TIMEOUT, TimeUnit.MILLISECONDS)) {                 // only majority servers return success can we consider this update success                 Loggers.RAFT.error("data publish failed, caused failed to notify majority, key={}", key);                 throw new IllegalStateException("data publish failed, caused failed to notify majority, key=" + key);             }             long end = System.currentTimeMillis();             Loggers.RAFT.info("signalPublish cost {} ms, key: {}", (end - start), key);         } finally {             OPERATE_LOCK.unlock();         }     }     //...... }
  • signalPublish方法判断当前节点是否是leader,如果不是则转发publish到leader节点的/v1/ns/raft/datum接口

  • 如果是leader则构造datum以及peers.majorityCount()大小的CountDownLatch,然后遍历peers.allServersIncludeMyself(),对于leader节点直接latch.countDown(),对于非leader节点则发送异步请求,请求/v1/ns/raft/datum/commit接口,在onCompleted的时候,如果请求成功执行latch.countDown()

  • 最后对于CountDownLatch未能在RAFT_PUBLISH_TIMEOUT返回的,抛出IllegalStateException

RaftController

nacos-1.1.3/naming/src/main/java/com/alibaba/nacos/naming/controllers/RaftController.java

@RestController @RequestMapping({UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft",     UtilsAndCommons.NACOS_SERVER_CONTEXT + UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft"}) public class RaftController {     @Autowired     private RaftConsistencyServiceImpl raftConsistencyService;     @Autowired     private ServiceManager serviceManager;     @Autowired     private RaftCore raftCore;     //......     @NeedAuth     @RequestMapping(value = "/datum/commit", method = RequestMethod.POST)     public String onPublish(HttpServletRequest request, HttpServletResponse response) throws Exception {         response.setHeader("Content-Type", "application/json; charset=" + getAcceptEncoding(request));         response.setHeader("Cache-Control", "no-cache");         response.setHeader("Content-Encode", "gzip");         String entity = IOUtils.toString(request.getInputStream(), "UTF-8");         String value = URLDecoder.decode(entity, "UTF-8");         JSONObject jsonObject = JSON.parseObject(value);         String key = "key";         RaftPeer source = JSON.parseObject(jsonObject.getString("source"), RaftPeer.class);         JSONObject datumJson = jsonObject.getJSONObject("datum");         Datum datum = null;         if (KeyBuilder.matchInstanceListKey(datumJson.getString(key))) {             datum = JSON.parseObject(jsonObject.getString("datum"), new TypeReference<Datum<Instances>>() {             });         } else if (KeyBuilder.matchSwitchKey(datumJson.getString(key))) {             datum = JSON.parseObject(jsonObject.getString("datum"), new TypeReference<Datum<SwitchDomain>>() {             });         } else if (KeyBuilder.matchServiceMetaKey(datumJson.getString(key))) {             datum = JSON.parseObject(jsonObject.getString("datum"), new TypeReference<Datum<Service>>() {             });         }         raftConsistencyService.onPut(datum, source);         return "ok";     }     //...... }
  • onPublish方法主要是执行raftConsistencyService.onPut(datum, source)

RaftConsistencyServiceImpl

nacos-1.1.3/naming/src/main/java/com/alibaba/nacos/naming/consistency/persistent/raft/RaftConsistencyServiceImpl.java

@Service public class RaftConsistencyServiceImpl implements PersistentConsistencyService {     @Autowired     private RaftCore raftCore;     @Autowired     private RaftPeerSet peers;     @Autowired     private SwitchDomain switchDomain;     //......     public void onPut(Datum datum, RaftPeer source) throws NacosException {         try {             raftCore.onPublish(datum, source);         } catch (Exception e) {             Loggers.RAFT.error("Raft onPut failed.", e);             throw new NacosException(NacosException.SERVER_ERROR, "Raft onPut failed, datum:" + datum + ", source: " + source, e);         }     }     //...... }
  • onPut方法执行的是raftCore.onPublish(datum, source)

RaftCore

nacos-1.1.3/naming/src/main/java/com/alibaba/nacos/naming/consistency/persistent/raft/RaftCore.java

@Component public class RaftCore {     public static final String API_VOTE = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/vote";     public static final String API_BEAT = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/beat";     public static final String API_PUB = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/datum";     public static final String API_DEL = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/datum";     public static final String API_GET = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/datum";     public static final String API_ON_PUB = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/datum/commit";     public static final String API_ON_DEL = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/datum/commit";     public static final String API_GET_PEER = UtilsAndCommons.NACOS_NAMING_CONTEXT + "/raft/peer";     private ScheduledExecutorService executor = new ScheduledThreadPoolExecutor(1, new ThreadFactory() {         @Override         public Thread newThread(Runnable r) {             Thread t = new Thread(r);             t.setDaemon(true);             t.setName("com.alibaba.nacos.naming.raft.notifier");             return t;         }     });     public static final Lock OPERATE_LOCK = new ReentrantLock();     public static final int PUBLISH_TERM_INCREASE_COUNT = 100;     private volatile Map<String, List<RecordListener>> listeners = new ConcurrentHashMap<>();     private volatile ConcurrentMap<String, Datum> datums = new ConcurrentHashMap<>();     @Autowired     private RaftPeerSet peers;     @Autowired     private SwitchDomain switchDomain;     @Autowired     private GlobalConfig globalConfig;     @Autowired     private RaftProxy raftProxy;     @Autowired     private RaftStore raftStore;     public volatile Notifier notifier = new Notifier();     private boolean initialized = false;     //......     public void onPublish(Datum datum, RaftPeer source) throws Exception {         RaftPeer local = peers.local();         if (datum.value == null) {             Loggers.RAFT.warn("received empty datum");             throw new IllegalStateException("received empty datum");         }         if (!peers.isLeader(source.ip)) {             Loggers.RAFT.warn("peer {} tried to publish data but wasn't leader, leader: {}",                 JSON.toJSONString(source), JSON.toJSONString(getLeader()));             throw new IllegalStateException("peer(" + source.ip + ") tried to publish " +                 "data but wasn't leader");         }         if (source.term.get() < local.term.get()) {             Loggers.RAFT.warn("out of date publish, pub-term: {}, cur-term: {}",                 JSON.toJSONString(source), JSON.toJSONString(local));             throw new IllegalStateException("out of date publish, pub-term:"                 + source.term.get() + ", cur-term: " + local.term.get());         }         local.resetLeaderDue();         // if data should be persistent, usually this is always true:         if (KeyBuilder.matchPersistentKey(datum.key)) {             raftStore.write(datum);         }         datums.put(datum.key, datum);         if (isLeader()) {             local.term.addAndGet(PUBLISH_TERM_INCREASE_COUNT);         } else {             if (local.term.get() + PUBLISH_TERM_INCREASE_COUNT > source.term.get()) {                 //set leader term:                 getLeader().term.set(source.term.get());                 local.term.set(getLeader().term.get());             } else {                 local.term.addAndGet(PUBLISH_TERM_INCREASE_COUNT);             }         }         raftStore.updateTerm(local.term.get());         notifier.addTask(datum.key, ApplyAction.CHANGE);         Loggers.RAFT.info("data added/updated, key={}, term={}", datum.key, local.term);     }     //...... }
  • onPublish方法首先判断请求的节点是否是leader,不是则抛出IllegalStateException;对于source.term小于local.term的抛出IllegalStateException

  • 之后执行local.resetLeaderDue(),以及raftStore.write(datum),datums.put(datum.key, datum);对于leader节点执行local.term.addAndGet(PUBLISH_TERM_INCREASE_COUNT),非leader节点则更新leader term以及local.term

  • 最后执行raftStore.updateTerm(local.term.get())以及notifier.addTask(datum.key, ApplyAction.CHANGE)

小结

  • signalPublish方法判断当前节点是否是leader,如果不是则转发publish到leader节点的/v1/ns/raft/datum接口

  • 如果是leader则构造datum以及peers.majorityCount()大小的CountDownLatch,然后遍历peers.allServersIncludeMyself(),对于leader节点直接latch.countDown(),对于非leader节点则发送异步请求,请求/v1/ns/raft/datum/commit接口,在onCompleted的时候,如果请求成功执行latch.countDown()

  • 最后对于CountDownLatch未能在RAFT_PUBLISH_TIMEOUT返回的,抛出IllegalStateException

感谢各位的阅读,以上就是“nacos RaftCore中signalPublish的原理和使用方法”的内容了,经过本文的学习后,相信大家对nacos RaftCore中signalPublish的原理和使用方法这一问题有了更深刻的体会,具体使用情况还需要大家实践验证。这里是亿速云,小编将为大家推送更多相关知识点的文章,欢迎关注!

向AI问一下细节

免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。

AI