diff --git a/hugegraph-store/docs/operations-guide.md b/hugegraph-store/docs/operations-guide.md index 8835dc1e48..6eb80ce136 100644 --- a/hugegraph-store/docs/operations-guide.md +++ b/hugegraph-store/docs/operations-guide.md @@ -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/ + 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://:8520/v1/partition/` should answer for each group. Then remove + the old record with `curl -X DELETE http://192.168.1.10:8620/v1/store/`. 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 diff --git a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/PartitionEngine.java b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/PartitionEngine.java index a70f17465f..6850267299 100644 --- a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/PartitionEngine.java +++ b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/PartitionEngine.java @@ -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; @@ -339,31 +340,8 @@ public Status changePeers(List 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()); @@ -433,6 +411,140 @@ public Status changePeers(List 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 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 shards) { + Map pdIds = partitionManager.shardIdsByEndpoint(shards); + Map localIds = + partitionManager.shardIdsByEndpoint(shardGroup.getMetaPbShard()); + String self = raftNode.getNodeId().getPeerId().getEndpoint().toString(); + List 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 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 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 changedShardEndpoints(Map pdIds, Map localIds, + List confEndpoints, String self) { + List changed = new ArrayList<>(); + Set 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")); @@ -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 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 { @@ -647,7 +760,7 @@ public void onConfigurationCommitted(Configuration conf) { } List 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 { @@ -663,7 +776,6 @@ public void onConfigurationCommitted(Configuration conf) { // partitionManager.changeShards(partition, shardGroup.getMetaPbShard()); // }); try { - var pdGroup = storeEngine.getPdProvider().getShardGroupDirect(getGroupId()); List peers = partitionManager.shards2Peers(pdGroup.getShardsList()); Long leaderStoreId = null; @@ -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); } @@ -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 diff --git a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/meta/PartitionManager.java b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/meta/PartitionManager.java index cc66893ec2..d078466c68 100644 --- a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/meta/PartitionManager.java +++ b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/meta/PartitionManager.java @@ -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 shardIdsByEndpoint(List shards) { + Map 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 shards = group.getShards(); for (Shard shard : shards) { diff --git a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/StoreIdChangeTest.java b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/StoreIdChangeTest.java new file mode 100644 index 0000000000..823c0f0d92 --- /dev/null +++ b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/StoreIdChangeTest.java @@ -0,0 +1,162 @@ +/* + * 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.hugegraph.store; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +import java.nio.file.Files; +import java.util.List; +import java.util.Map; + +import org.apache.hugegraph.pd.grpc.Metapb; +import org.apache.hugegraph.store.meta.PartitionManager; +import org.apache.hugegraph.store.meta.Shard; +import org.apache.hugegraph.store.meta.ShardGroup; +import org.apache.hugegraph.store.meta.Store; +import org.apache.hugegraph.store.options.HgStoreEngineOptions; +import org.apache.hugegraph.store.pd.PdProvider; +import org.junit.Before; +import org.junit.Test; +import org.mockito.Mockito; + +/** + * A Store rebuilt with an empty disk at the same raft address registers under a new store id. + */ +public class StoreIdChangeTest { + + private static final String ADDR0 = "store-0:8510"; + private static final String ADDR1 = "store-1:8510"; + private static final String ADDR2 = "store-2:8510"; + private static final long S0 = 10L; + private static final long S1 = 11L; + // old and new id of the Store at ADDR2 + private static final long X = 12L; + private static final long Y = 13L; + + private PdProvider pdProvider; + private PartitionManager partitionManager; + + @Before + public void setUp() throws Exception { + pdProvider = Mockito.mock(PdProvider.class); + mockStore(S0, ADDR0); + mockStore(S1, ADDR1); + mockStore(X, ADDR2); + mockStore(Y, ADDR2); + String path = Files.createTempDirectory("store-id-change").toString(); + HgStoreEngineOptions options = new HgStoreEngineOptions(); + options.setDataPath(path); + options.setRaftPath(path); + partitionManager = new PartitionManager(pdProvider, options); + } + + private void mockStore(long id, String raftAddress) { + Store store = new Store(Metapb.Store.newBuilder().setId(id) + .setRaftAddress(raftAddress).build()); + Mockito.when(pdProvider.getStoreByID(id)).thenReturn(store); + } + + private static ShardGroup localGroup(long... storeIds) { + ShardGroup group = new ShardGroup(); + group.setId(1); + for (long id : storeIds) { + Shard shard = new Shard(); + shard.setStoreId(id); + shard.setRole(Metapb.ShardRole.Follower); + group.getShards().add(shard); + } + return group; + } + + private static Metapb.ShardGroup pdGroup(long... storeIds) { + Metapb.ShardGroup.Builder builder = Metapb.ShardGroup.newBuilder().setId(1); + for (long id : storeIds) { + builder.addShards(Metapb.Shard.newBuilder().setStoreId(id) + .setRole(Metapb.ShardRole.Follower)); + } + return builder.build(); + } + + @Test + public void testEndpointKeepsLocalIdStillInPdGroup() { + Store store = partitionManager.getStoreByRaftEndpoint(localGroup(S0, S1, X), + pdGroup(S0, S1, X), ADDR2); + assertEquals(X, store.getId()); + } + + @Test + public void testEndpointTakesPdIdWhenLocalIdRetired() { + Store store = partitionManager.getStoreByRaftEndpoint(localGroup(S0, S1, X), + pdGroup(S0, S1, Y), ADDR2); + assertEquals(Y, store.getId()); + } + + @Test + public void testEndpointTakesPdIdWhenLocalIdDeleted() { + Mockito.when(pdProvider.getStoreByID(X)).thenReturn(null); + Store store = partitionManager.getStoreByRaftEndpoint(localGroup(S0, S1, X), + pdGroup(S0, S1, Y), ADDR2); + assertEquals(Y, store.getId()); + } + + @Test + public void testEndpointWithoutPdGroupUsesLocalGroup() { + Store store = partitionManager.getStoreByRaftEndpoint(localGroup(S0, S1, X), null, + ADDR2); + assertEquals(X, store.getId()); + // unknown endpoint: id 0, the caller asks the endpoint itself + store = partitionManager.getStoreByRaftEndpoint(localGroup(S0, S1, X), + pdGroup(S0, S1, X), "store-3:8510"); + assertEquals(0L, store.getId()); + } + + @Test + public void testChangedEndpointsFindsNewIdAtSameAddress() { + Map pdIds = partitionManager.shardIdsByEndpoint( + pdGroup(S0, S1, Y).getShardsList()); + Map localIds = partitionManager.shardIdsByEndpoint( + localGroup(S0, S1, X).getMetaPbShard()); + List changed = PartitionEngine.changedShardEndpoints( + pdIds, localIds, List.of(ADDR0, ADDR1, ADDR2), ADDR0); + assertEquals(List.of(ADDR2), changed); + } + + @Test + public void testChangedEndpointsIgnoresMembershipChange() { + Map pdIds = Map.of(ADDR0, S0, ADDR1, S1, "store-3:8510", Y); + Map localIds = Map.of(ADDR0, S0, ADDR1, S1, ADDR2, X); + assertTrue(PartitionEngine.changedShardEndpoints( + pdIds, localIds, List.of(ADDR0, ADDR1, ADDR2), ADDR0).isEmpty()); + } + + @Test + public void testChangedEndpointsNeverIncludesSelf() { + Map pdIds = Map.of(ADDR0, Y, ADDR1, S1, ADDR2, X); + Map localIds = Map.of(ADDR0, S0, ADDR1, S1, ADDR2, X); + assertTrue(PartitionEngine.changedShardEndpoints( + pdIds, localIds, List.of(ADDR0, ADDR1, ADDR2), ADDR0).isEmpty()); + } + + @Test + public void testChangedEndpointsEmptyWhenIdsMatch() { + Map ids = Map.of(ADDR0, S0, ADDR1, S1, ADDR2, Y); + assertTrue(PartitionEngine.changedShardEndpoints( + ids, ids, List.of(ADDR0, ADDR1, ADDR2), ADDR0).isEmpty()); + } +} diff --git a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/CoreSuiteTest.java b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/CoreSuiteTest.java index 6afd046e18..830c348244 100644 --- a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/CoreSuiteTest.java +++ b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/CoreSuiteTest.java @@ -17,6 +17,7 @@ package org.apache.hugegraph.store.core; +import org.apache.hugegraph.store.StoreIdChangeTest; import org.apache.hugegraph.store.core.snapshot.HgSnapshotHandlerTest; import org.junit.runner.RunWith; import org.junit.runners.Suite; @@ -44,7 +45,8 @@ // HgBusinessImplTest.class @RunWith(Suite.class) @Suite.SuiteClasses({ - HgSnapshotHandlerTest.class + HgSnapshotHandlerTest.class, + StoreIdChangeTest.class }) @Slf4j public class CoreSuiteTest {