(boolean isMove)
| 564 | |
| 565 | // 开始分桶流程有两个线程需要访问:timer & raft.UserThreadExecutor |
| 566 | private void startSplit(boolean isMove) throws Exception { |
| 567 | lock(); |
| 568 | try { |
| 569 | if (!getRaft().isLeader()) |
| 570 | return; |
| 571 | |
| 572 | // 后面需要在lambda中传递系列号作为上下文,使用成员变量是不是会跟随变化? |
| 573 | var serialNo = manager.atomicSerialNo.incrementAndGet(); |
| 574 | splitSerialNo = serialNo; |
| 575 | var bucket = stateMachine.getBucket(); |
| 576 | RocksIterator it = null; |
| 577 | try { |
| 578 | var splitting = bucket.getSplittingMeta(); // 对于timer,这个会调用两次。 |
| 579 | if (null == splitting) { |
| 580 | // 第一次开始分桶,准备阶段。 |
| 581 | // 这个阶段在timer回调中执行,可以同步调用一些网络接口。 |
| 582 | // 先去manager查一下可用的manager是否够,简单判断,不原子化。 |
| 583 | if (manager.getMasterAgent().checkFreeManager() < dbh2Config.getRaftClusterCount()) { |
| 584 | logger.warn("splitting not enough free manager. isMove={}", isMove); |
| 585 | return; |
| 586 | } |
| 587 | // 上一次分桶结束的deleteRange可能还没compact,此时keyNumbers不准确,这里总是执行一次。 |
| 588 | bucket.getData().compact(); |
| 589 | |
| 590 | it = isMove ? locateFirst() : locateMiddle(); |
| 591 | if (null == it) { |
| 592 | logger.info("splitting break start: it is null. isMove={}", isMove); |
| 593 | return; // empty?不需要执行后续操作。break progress. |
| 594 | } |
| 595 | var newMeta = stateMachine.getBucket().getBucketMeta().copy(); |
| 596 | newMeta.setRaftConfig(""); |
| 597 | if (!isMove) |
| 598 | newMeta.setKeyFirst(new Binary(it.key())); |
| 599 | splitting = manager.getMasterAgent().createSplitBucket(newMeta); |
| 600 | |
| 601 | // 设置分桶进行中的标记到raft集群中。 |
| 602 | getRaft().appendLog(new LogSetSplittingMeta(splitting)); |
| 603 | // 创建到分桶目标的客户端。 |
| 604 | logger.info("splitting start... isMove={} {}->{}", |
| 605 | isMove, formatMeta(bucket.getBucketMeta()), formatMeta(splitting)); |
| 606 | } |
| 607 | |
| 608 | // 重启的时候,需要重建到分桶的连接。 |
| 609 | if (null == dbh2Splitting) { |
| 610 | dbh2Splitting = new Dbh2Agent(splitting.getRaftConfig(), RaftAgentNetClient::new); |
| 611 | dbh2Splitting.getRaftAgent().setPendingLimit(Integer.MAX_VALUE); |
| 612 | } |
| 613 | |
| 614 | var server = (Dbh2RaftServer)getRaft().getServer(); |
| 615 | performPrepareQueue(server.takePrepareQueue()); |
| 616 | |
| 617 | if (null == it) { |
| 618 | // 重新开始分桶时走这个分支,根据上次找到的middle,定位it。 |
| 619 | it = isMove ? locateFirst() : locateMiddle(splitting.getKeyFirst()); |
| 620 | if (null == it) { |
| 621 | logger.info("splitting break restart: it is null. isMove={}", isMove); |
| 622 | return; |
| 623 | } |
no test coverage detected