Skip to content
Merged
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
16 changes: 16 additions & 0 deletions hugegraph-store/docs/operations-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -498,6 +498,22 @@ scp backup-store1-*.tar.gz backup-server:/backups/
- Deploy new Store node with same configuration
- PD automatically assigns partitions to new node
- Wait for data replication (may take hours)
- If the new node reuses the failed node's raft address with an empty data directory (for
example a Kubernetes StatefulSet Pod rebuilt after its volume was lost), it registers under
a new store ID while the old ID still holds its partitions. Retire the old ID on the PD
leader so the replicas move to the new node:
```bash
# Two entries share the address: the new ID and the old one
curl http://192.168.1.10:8620/v1/stores
curl -X POST -H 'Content-Type: application/json' -d '{"storeState":"Tombstone"}' \
http://192.168.1.10:8620/v1/store/<oldStoreId>
curl http://192.168.1.10:8620/v1/task/patrolPartitions
```
Every shard group in `/v1/shardGroups` should then list the new ID, and the new node's
`http://<store>:8520/v1/partition/<partitionId>` should answer for each group. Then remove
the old record with `curl -X DELETE http://192.168.1.10:8620/v1/store/<oldStoreId>`. Store
versions without the fix for apache/hugegraph#3227 never finish this: the groups keep the
old ID and the new node stays empty.

4. **Verify**: Check partition distribution
```bash
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.atomic.AtomicBoolean;
Expand Down Expand Up @@ -339,31 +340,8 @@ public Status changePeers(List<String> peers, final Closure done) {
}

// 2.2 Waiting learner to synchronize snapshot (check added learner)
//todo Each learner will wait for 1s, if another one is not sync.Consider using
// countdownLatch
boolean allLearnerSnapshotOk = false;
long current = System.currentTimeMillis();
while (!allLearnerSnapshotOk) {
boolean snapshotOk = true;
for (var peerId : addPeers) {
var state = getReplicatorState(JRaftUtils.getPeerId(peerId));
log.info("Raft {}, peer:{}, replicate state:{}", getGroupId(), peerId, state);
if (state != Replicator.State.Replicate) {
snapshotOk = false;
}
}
allLearnerSnapshotOk = snapshotOk;

if (!allLearnerSnapshotOk) {
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
log.warn("Raft {} sleep when check learner snapshot", getGroupId());
}
}
if (System.currentTimeMillis() - current > 600 * 1000) {
return HgRaftError.TASK_CONTINUE.toStatus();
}
if (!waitForReplicate(addPeers)) {
return HgRaftError.TASK_CONTINUE.toStatus();
}

log.info("Raft {} replicate status is OK", getGroupId());
Expand Down Expand Up @@ -433,6 +411,140 @@ public Status changePeers(List<String> peers, final Closure done) {
return removeSelf ? HgRaftError.TASK_CONTINUE.toStatus() : HgRaftError.OK.toStatus();
}

/**
* Wait until the replicator of every peer is in Replicate state (snapshot installed).
*
* @return false if the peers did not catch up within 600s
*/
private boolean waitForReplicate(List<String> peers) {
//todo Each learner will wait for 1s, if another one is not sync.Consider using
// countdownLatch
boolean allLearnerSnapshotOk = false;
long current = System.currentTimeMillis();
while (!allLearnerSnapshotOk) {
boolean snapshotOk = true;
for (var peerId : peers) {
var state = getReplicatorState(JRaftUtils.getPeerId(peerId));
log.info("Raft {}, peer:{}, replicate state:{}", getGroupId(), peerId, state);
if (state != Replicator.State.Replicate) {
snapshotOk = false;
}
}
allLearnerSnapshotOk = snapshotOk;

if (!allLearnerSnapshotOk) {
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
log.warn("Raft {} sleep when check learner snapshot", getGroupId());
}
}
if (System.currentTimeMillis() - current > 600 * 1000) {
return false;
}
}
return true;
}

/**
* PD's shard list and the raft configuration have the same endpoints, but PD names another
* store id at one of them: a Store rebuilt with an empty disk at the same raft address
* registers under a new id. The raft configuration has nothing to change, so:
* 1. Create the raft node on that endpoint, the leader then installs a snapshot into it.
* 2. Take PD's store ids into the local shard group and report it to PD, otherwise the
* partition heartbeat writes the old id back to PD.
* 3. Wait for the peer to catch up.
*
* @param shards shard list from PD
* @return OK when nothing changed or the peers caught up, the createRaftNode status when an
* endpoint is unreachable (retried when its replicator comes online), TASK_ERROR when a peer
* has not caught up in time (its replicator keeps installing the snapshot)
*/
private Status syncShardIdentities(List<Metapb.Shard> shards) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

0c2a180 adds the procedure to hugegraph-store/docs/operations-guide.md (Single Store Node Failure, step 3): retire the old id with Tombstone, run patrolPartitions, verify, then delete it, with the note that Store versions without this fix never finish it. The website only lists these PD endpoints (quickstart hugegraph-pd), with no replacement procedure to correct. The Helm chart README is not in master yet; it lives in #3218, and its "do not delete a Store volume" warning changes there with a version caveat once this is merged and rerun on master images.

Map<String, Long> pdIds = partitionManager.shardIdsByEndpoint(shards);
Map<String, Long> localIds =
partitionManager.shardIdsByEndpoint(shardGroup.getMetaPbShard());
String self = raftNode.getNodeId().getPeerId().getEndpoint().toString();
List<String> changed = changedShardEndpoints(pdIds, localIds,
RaftUtils.getAllEndpoints(raftNode), self);
if (changed.isEmpty()) {
return HgRaftError.OK.toStatus();
}
log.info("Raft {} store id changed at raft address {}, local {}, pd {}",
getGroupId(), changed, localIds, pdIds);

for (String peer : changed) {
FutureClosure closure = new FutureClosure();
storeEngine.getHgCmdClient().createRaftNode(peer, partitionManager.getPartitionList(
getGroupId()), getCurrentConf(), closure);
Status status = closure.get();
if (!status.isOk()) {
log.info("Raft {} createRaftNode, peer:{}, reason:{}", getGroupId(), peer,
status.getErrorMsg());
return status;
}
}
doSnapshot(status -> log.info("Raft {} snapshot after create raft node, result:{}",
getGroupId(), status));

List<Long> peerIds = new ArrayList<>();
for (String peer : RaftUtils.getPeerEndpoints(raftNode)) {
Long id = pdIds.getOrDefault(peer.toLowerCase(), localIds.get(peer.toLowerCase()));
if (id != null) {
peerIds.add(id);
}
}
List<Long> learners = new ArrayList<>();
for (String learner : RaftUtils.getLearnerEndpoints(raftNode)) {
Long id = pdIds.getOrDefault(learner.toLowerCase(),
localIds.get(learner.toLowerCase()));
if (id != null) {
learners.add(id);
}
}
shardGroup.changeShardList(peerIds, learners, partitionManager.getStore().getId());
partitionManager.updateShardGroup(shardGroup);
try {
partitionManager.getPdProvider().updateShardGroup(shardGroup.getProtoObj());
} catch (PDException e) {
log.warn("Raft {} update shard group to pd failed, {}", getGroupId(), e.getMessage());
}
log.info("Raft {} shard group after store id change {}", getGroupId(),
shardGroup.getMetaPbShard());

// The store ids are taken before the wait: a rebuilt Store that restarts before PD names
// it would exit in loadPartitions. A timeout loses nothing, the replicator keeps going.
if (!waitForReplicate(changed)) {
log.warn("Raft {} peers {} not caught up in time, replication continues",
getGroupId(), changed);
return HgRaftError.TASK_ERROR.toStatus();
}
return HgRaftError.OK.toStatus();
}

/**
* Endpoints of the raft configuration, except self, where PD names a different store id
* than the local shard group. Empty if PD's endpoints differ from the configuration, which
* is a membership change for changePeers.
*/
static List<String> changedShardEndpoints(Map<String, Long> pdIds, Map<String, Long> localIds,
List<String> confEndpoints, String self) {
List<String> changed = new ArrayList<>();
Set<String> endpoints = new HashSet<>();
confEndpoints.forEach(endpoint -> endpoints.add(endpoint.toLowerCase()));
if (!endpoints.equals(pdIds.keySet())) {
return changed;
}
for (String endpoint : confEndpoints) {
String key = endpoint.toLowerCase();
if (!key.equals(self.toLowerCase()) &&
!Objects.equals(pdIds.get(key), localIds.get(key))) {
changed.add(endpoint);
}
}
return changed;
}

public void addRaftTask(RaftOperation operation, RaftClosure closure) {
if (!isLeader()) {
closure.run(new Status(HgRaftError.NOT_LEADER.getNumber(), "Not leader"));
Expand Down Expand Up @@ -635,10 +747,11 @@ public void onConfigurationCommitted(Configuration conf) {
try {
// Update shardlist
log.info("Raft {} onConfigurationCommitted, conf is {}", getGroupId(), conf.toString());
var pdGroup = storeEngine.getPdProvider().getShardGroupDirect(getGroupId());
// According to raft endpoint find storeId
List<Long> peerIds = new ArrayList<>();
for (String peer : RaftUtils.getPeerEndpoints(conf)) {
Store store = getStoreByEndpoint(peer);
Store store = getStoreByEndpoint(pdGroup, peer);
if (store != null) {
peerIds.add(store.getId());
} else {
Expand All @@ -647,7 +760,7 @@ public void onConfigurationCommitted(Configuration conf) {
}
List<Long> learners = new ArrayList<>();
for (String learner : RaftUtils.getLearnerEndpoints(conf)) {
Store store = getStoreByEndpoint(learner);
Store store = getStoreByEndpoint(pdGroup, learner);
if (store != null) {
learners.add(store.getId());
} else {
Expand All @@ -663,7 +776,6 @@ public void onConfigurationCommitted(Configuration conf) {
// partitionManager.changeShards(partition, shardGroup.getMetaPbShard());
// });
try {
var pdGroup = storeEngine.getPdProvider().getShardGroupDirect(getGroupId());
List<String> peers = partitionManager.shards2Peers(pdGroup.getShardsList());

Long leaderStoreId = null;
Expand Down Expand Up @@ -693,8 +805,8 @@ public void onConfigurationCommitted(Configuration conf) {

}

private Store getStoreByEndpoint(String endpoint) {
Store store = partitionManager.getStoreByRaftEndpoint(getShardGroup(), endpoint);
private Store getStoreByEndpoint(Metapb.ShardGroup pdGroup, String endpoint) {
Store store = partitionManager.getStoreByRaftEndpoint(getShardGroup(), pdGroup, endpoint);
if (store == null || store.getId() == 0) {
store = this.storeEngine.getHgCmdClient().getStoreInfo(endpoint);
}
Expand Down Expand Up @@ -770,6 +882,9 @@ public void doChangeShard(final MetaTask.Task task, Closure done) {
return;
}
Status result = changePeers(peers, null);
if (result.isOk()) {
result = syncShardIdentities(task.getChangeShard().getShardList());
}

if (result.getCode() == HgRaftError.TASK_CONTINUE.getNumber()) {
// Need to resend a request
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -881,6 +881,41 @@ public Store getStoreByRaftEndpoint(ShardGroup group, String endpoint) {
return result[0];
}

/**
* According to raft address to find Store, checked against PD's current shard group.
* A Store rebuilt with an empty disk at the same raft address registers under a new id,
* so the id held locally for that address may no longer be in PD's group; use PD's id then.
*/
public Store getStoreByRaftEndpoint(ShardGroup group, Metapb.ShardGroup pdGroup,
String endpoint) {
Store store = getStoreByRaftEndpoint(group, endpoint);
if (pdGroup == null ||
pdGroup.getShardsList().stream().anyMatch(s -> s.getStoreId() == store.getId())) {
return store;
}
Long pdStoreId = shardIdsByEndpoint(pdGroup.getShardsList()).get(endpoint.toLowerCase());
if (pdStoreId == null) {
return store;
}
log.info("Raft {} endpoint {} is store {} locally but {} in PD, using PD",
pdGroup.getId(), endpoint, store.getId(), pdStoreId);
return getStore(pdStoreId);
}

/**
* Raft address (lower case) to store id, for the shards whose store PD can still resolve
*/
public Map<String, Long> shardIdsByEndpoint(List<Metapb.Shard> shards) {
Map<String, Long> result = new HashMap<>();
for (Metapb.Shard shard : shards) {
Store store = getStore(shard.getStoreId());
if (store != null && !store.getRaftAddress().isEmpty()) {
result.put(store.getRaftAddress().toLowerCase(), shard.getStoreId());
}
}
return result;
}

public Shard getShardByEndpoint(ShardGroup group, String endpoint) {
List<Shard> shards = group.getShards();
for (Shard shard : shards) {
Expand Down
Loading
Loading