-
Notifications
You must be signed in to change notification settings - Fork 640
fix(server): converge schema caches across HStore servers #3237
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
a4f7ae1
be93242
070a961
7f5672a
be867b2
c26d961
2d3dd41
1d4335f
8d6e186
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -46,6 +46,7 @@ | |
| import com.alipay.sofa.jraft.Node; | ||
| import com.alipay.sofa.jraft.RaftGroupService; | ||
| import com.alipay.sofa.jraft.Status; | ||
| import com.alipay.sofa.jraft.closure.ReadIndexClosure; | ||
| import com.alipay.sofa.jraft.conf.Configuration; | ||
| import com.alipay.sofa.jraft.core.Replicator; | ||
| import com.alipay.sofa.jraft.core.State; | ||
|
|
@@ -74,6 +75,15 @@ public class RaftEngine { | |
| */ | ||
| private static final long ALIVE_PEERS_REFRESH_MS = 1000L; | ||
|
|
||
| private static final int READ_INDEX_RETRIES = 5; | ||
| private static final long READ_INDEX_RETRY_DELAY_MS = 20L; | ||
| private static final ScheduledExecutorService READ_INDEX_RETRY = | ||
| Executors.newSingleThreadScheduledExecutor(runnable -> { | ||
| Thread thread = new Thread(runnable, "pd-raft-read-index-retry"); | ||
| thread.setDaemon(true); | ||
| return thread; | ||
| }); | ||
|
|
||
| private volatile static RaftEngine instance = new RaftEngine(); | ||
| private RaftStateMachine stateMachine; | ||
| private String groupId = "pd_raft"; | ||
|
|
@@ -133,9 +143,7 @@ public synchronized boolean init(PDConfig.Raft config) { | |
| // Snapshot interval | ||
| nodeOptions.setSnapshotIntervalSecs(config.getSnapshotInterval()); | ||
|
|
||
| nodeOptions.setRpcConnectTimeoutMs(config.getRpcTimeout()); | ||
| nodeOptions.setRpcDefaultTimeout(config.getRpcTimeout()); | ||
| nodeOptions.setRpcInstallSnapshotTimeout(config.getRpcTimeout()); | ||
| setRpcTimeouts(nodeOptions, config); | ||
| // TODO: tune RaftOptions for PD (see hugegraph-store PartitionEngine for reference) | ||
|
|
||
| final PeerId serverId = JRaftUtils.getPeerId(config.getAddress()); | ||
|
|
@@ -152,6 +160,20 @@ public synchronized boolean init(PDConfig.Raft config) { | |
| return this.raftNode != null; | ||
| } | ||
|
|
||
| /** | ||
| * Wire the three raft rpc timeouts from their own options. jraft pings a peer to open a | ||
| * connection, and a candidate does so for every peer while it holds the node lock, so a | ||
| * peer that accepts the connection but never answers (a stopped process, a lost host) | ||
| * stalls the election, and the answers to the other peers' votes, for the whole connect | ||
| * timeout. Installing a snapshot sends the whole store and needs far longer than a | ||
| * normal request. | ||
| */ | ||
| static void setRpcTimeouts(NodeOptions nodeOptions, PDConfig.Raft config) { | ||
| nodeOptions.setRpcConnectTimeoutMs(config.getRpcConnectTimeout()); | ||
| nodeOptions.setRpcDefaultTimeout(config.getRpcTimeout()); | ||
| nodeOptions.setRpcInstallSnapshotTimeout(config.getRpcInstallSnapshotTimeout()); | ||
| } | ||
|
|
||
| /** | ||
| * Create a Raft RPC Server for communication between PDs | ||
| */ | ||
|
|
@@ -602,6 +624,69 @@ public Node getRaftNode() { | |
| return raftNode; | ||
| } | ||
|
|
||
| /** | ||
| * Wait until this node has applied every entry committed before the call, so a local | ||
| * read that follows sees every write acknowledged before it. jraft's ReadIndex confirms | ||
| * the commit index with a quorum in the current term and runs the closure once the | ||
| * applied index reaches it; a leader elected a moment ago first waits for an entry of | ||
| * its own term to commit. Bounded by the raft rpc timeout. | ||
| * <p> | ||
| * Never call it on the state machine thread: the closure waits for that thread to apply. | ||
| */ | ||
| public void waitReadIndex() throws PDException { | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Minor: This ReadIndex wait, and the other four stage A commits ( |
||
| waitReadIndex(this.config.getRpcTimeout()); | ||
| } | ||
|
|
||
| void waitReadIndex(long timeoutMs) throws PDException { | ||
| CompletableFuture<Void> future = new CompletableFuture<>(); | ||
| readIndex(future, READ_INDEX_RETRIES); | ||
| try { | ||
| future.get(timeoutMs, TimeUnit.MILLISECONDS); | ||
| } catch (InterruptedException e) { | ||
| Thread.currentThread().interrupt(); | ||
| throw new PDException(Pdpb.ErrorType.UNKNOWN_VALUE, | ||
| "Interrupted while waiting for the raft read index", e); | ||
| } catch (TimeoutException e) { | ||
| throw new PDException(Pdpb.ErrorType.UNKNOWN_VALUE, | ||
| String.format("Raft read index timed out after %d ms", | ||
| timeoutMs)); | ||
| } catch (ExecutionException e) { | ||
| if (e.getCause() instanceof PDException) { | ||
| throw (PDException) e.getCause(); | ||
| } | ||
| throw new PDException(Pdpb.ErrorType.UNKNOWN_VALUE, e.getCause()); | ||
| } | ||
| } | ||
|
|
||
| private void readIndex(CompletableFuture<Void> future, int retries) { | ||
| Node node = this.raftNode; | ||
| if (node == null) { | ||
| future.completeExceptionally(new PDException(Pdpb.ErrorType.NOT_LEADER_VALUE, | ||
| "Raft node is not started")); | ||
| return; | ||
| } | ||
| node.readIndex(new byte[0], new ReadIndexClosure() { | ||
| @Override | ||
| public void run(Status status, long index, byte[] reqCtx) { | ||
| if (status.isOk()) { | ||
| future.complete(null); | ||
| return; | ||
| } | ||
| RaftError error = status.getRaftError(); | ||
| // EAGAIN: no entry of the current term committed yet; EBUSY: transferring | ||
| if (retries > 0 && (error == RaftError.EAGAIN || error == RaftError.EBUSY)) { | ||
| READ_INDEX_RETRY.schedule(() -> readIndex(future, retries - 1), | ||
| READ_INDEX_RETRY_DELAY_MS, TimeUnit.MILLISECONDS); | ||
| return; | ||
| } | ||
| int type = error == RaftError.EPERM ? Pdpb.ErrorType.NOT_LEADER_VALUE : | ||
| Pdpb.ErrorType.UNKNOWN_VALUE; | ||
| future.completeExceptionally( | ||
| new PDException(type, "Raft read index failed: " + status)); | ||
| } | ||
| }); | ||
| } | ||
|
|
||
| private boolean peerEquals(PeerId p1, PeerId p2) { | ||
| if (p1 == null && p2 == null) { | ||
| return true; | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
waitReadIndex()waits forrpcTimeout(10 seconds by default), but this code schedules only five retries 20 ms apart; after roughly 100 ms,readIndex()completes withUNKNOWN_VALUE. A newly elected JRaft leader returnsEAGAINuntil an entry in its current term commits, so a quorum that needs longer than 100 ms makesIdMetaStore.getId()fail even while the configured timeout remains. Please retryEAGAIN/EBUSYuntil thetimeoutMsdeadline and add coverage for recovery after the current retry window.