扫二维码与项目经理沟通
我们在微信上24小时期待你的声音
解答本文疑问/技术咨询/运营咨询/技术建议/互联网交流
这篇文章主要讲解了“nacos RaftCore中signalPublish的原理和使用方法”,文中的讲解内容简单清晰,易于学习与理解,下面请大家跟着小编的思路慢慢深入,一起来研究和学习“nacos RaftCore中signalPublish的原理和使用方法”吧!
十多年建站经验, 网站制作、网站建设客户的见证与正确选择。成都创新互联公司提供完善的营销型网页建站明细报价表。后期开发更加便捷高效,我们致力于追求更美、更快、更规范。
本文主要研究一下nacos RaftCore的signalPublish
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> listeners = new ConcurrentHashMap<>(); private volatile ConcurrentMap 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 > 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 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 () { @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
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>() { }); } else if (KeyBuilder.matchSwitchKey(datumJson.getString(key))) { datum = JSON.parseObject(jsonObject.getString("datum"), new TypeReference >() { }); } else if (KeyBuilder.matchServiceMetaKey(datumJson.getString(key))) { datum = JSON.parseObject(jsonObject.getString("datum"), new TypeReference >() { }); } raftConsistencyService.onPut(datum, source); return "ok"; } //...... }
onPublish方法主要是执行raftConsistencyService.onPut(datum, source)
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)
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> listeners = new ConcurrentHashMap<>(); private volatile ConcurrentMap 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的原理和使用方法这一问题有了更深刻的体会,具体使用情况还需要大家实践验证。这里是创新互联,小编将为大家推送更多相关知识点的文章,欢迎关注!
我们在微信上24小时期待你的声音
解答本文疑问/技术咨询/运营咨询/技术建议/互联网交流