Skip to content
Open
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
3 changes: 3 additions & 0 deletions hugegraph-pd/docs/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,9 @@ raft:
|-----------|------|---------|-------------|
| `raft.address` | String | `127.0.0.1:8610` | Raft service address for this PD node. Format: `<ip>:<port>`. Must be unique across all PD nodes. |
| `raft.peers-list` | String | `127.0.0.1:8610` | Comma-separated list of all PD nodes' Raft addresses. Used for cluster formation and leader election. |
| `raft.rpc-timeout` | Integer | `10000` | Timeout of a Raft RPC between PD nodes, in milliseconds. |
| `raft.rpc-connect-timeout` | Integer | `1000` | Timeout for opening a connection to a peer, in milliseconds. A candidate opens its connections one peer at a time while it holds the Raft node lock, so an election waits this long for each peer that accepts the connection but does not answer (a stopped process, a lost host). |
| `raft.rpc-install-snapshot-timeout` | Integer | `0` | Timeout for sending a snapshot to a follower that fell behind the log, in milliseconds. It bounds the transfer of the whole metadata store. `0` uses `raft.rpc-timeout`, the value this timeout had before it became a separate option; raise it if large snapshots time out. |

**Critical Rules**:
1. `raft.address` must be unique for each PD node
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -186,11 +186,24 @@ public class Raft {
private int snapshotInterval;
@Value("${raft.rpc-timeout:10000}")
private int rpcTimeout;
// Bounds the ping that opens a connection to a peer; kept apart from rpc-timeout
// because a candidate opens its connections while it holds the raft node lock
@Value("${raft.rpc-connect-timeout:1000}")
private int rpcConnectTimeout = 1000;
// A follower that fell behind the log receives the whole store within this time;
// 0 keeps the timeout it always had, raft.rpc-timeout
@Value("${raft.rpc-install-snapshot-timeout:0}")
private int rpcInstallSnapshotTimeout;
@Value("${grpc.host}")
private String host;
@Value("${server.port}")
private int port;

public int getRpcInstallSnapshotTimeout() {
return this.rpcInstallSnapshotTimeout > 0 ? this.rpcInstallSnapshotTimeout :
this.rpcTimeout;
}

@Value("${pd.cluster_id:1}")
private long clusterId;
@Value("${grpc.port}")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,10 @@ public long getId(String key, int delta) throws PDException {
Object probableLock = getLock(key);
byte[] keyBs = (ID_PREFIX + key).getBytes(Charset.defaultCharset());
synchronized (probableLock) {
// The read and the put are two steps, not one raft entry: a leader elected a
// moment ago may not have applied its predecessor's last put yet, and would
// hand out the same range again without this wait
getStore().waitReadIndex();
byte[] bs = getOne(keyBs);
long current = bs != null ? bytesToLong(bs) : 0L;
long next = current + delta;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -74,6 +75,15 @@ public class RaftEngine {
*/
private static final long ALIVE_PEERS_REFRESH_MS = 1000L;

private static final int READ_INDEX_RETRIES = 5;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

⚠️ [P2] Keep retrying transient ReadIndex failures until the configured deadline. waitReadIndex() waits for rpcTimeout (10 seconds by default), but this code schedules only five retries 20 ms apart; after roughly 100 ms, readIndex() completes with UNKNOWN_VALUE. A newly elected JRaft leader returns EAGAIN until an entry in its current term commits, so a quorum that needs longer than 100 ms makes IdMetaStore.getId() fail even while the configured timeout remains. Please retry EAGAIN/EBUSY until the timeoutMs deadline and add coverage for recovery after the current retry window.

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";
Expand Down Expand Up @@ -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());
Expand All @@ -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
*/
Expand Down Expand Up @@ -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 {

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.

Minor: This ReadIndex wait, and the other four stage A commits (1d4335fdd snapshot checkpoint on the FSM thread, 8d6e18684 raft.rpc-connect-timeout, be867b28c hstore truncate keeping the schema id counters, c26d96130 the AUTH/ prefix), fix master bugs that do not depend on the schema version work. In this PR they are tied to a design that has an open CHANGES_REQUESTED review and a planned rewrite (stages B to E), so fixes for duplicate PD id ranges, double-applied snapshot entries and wiped AUTHx graphs wait for that design to settle. apache/hugegraph squash-merges, so they would also land as one commit titled fix(server): converge schema caches across HStore servers, which makes any one of them hard to find or revert, and the PR description does not mention them or the new PD options. Requested change: move the stage A commits into their own PR (or one PR per module, PD and server), each with its own description and tests, and keep this PR to the schema cache protocol.

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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -195,45 +195,57 @@ public void onConfigurationCommitted(final Configuration conf) {
log.info("Raft onConfigurationCommitted {}", conf);
}

/**
* jraft calls this on the state machine thread, with the snapshot index set to the last
* applied index. The RocksDB checkpoint is taken here, before the thread applies the next
* entry: a checkpoint taken later holds entries past the snapshot index, and a node that
* installs the snapshot applies them a second time. Only the compression runs on the job
* thread. The checkpoint is hard links under the store's write lock, so it is short.
*/
@Override
public void onSnapshotSave(final SnapshotWriter writer, final Closure done) {
MetadataService.getUninterruptibleJobs().submit(() -> {
lock.lock();
String snapshotDir = writer.getPath() + File.separator + SNAPSHOT_DIR_NAME;
lock.lock();
try {
log.info("start snapshot save");
try {
log.info("start snapshot save");
String snapshotDir = writer.getPath() + File.separator + SNAPSHOT_DIR_NAME;
FileUtils.deleteDirectory(new File(snapshotDir));
FileUtils.forceMkdir(new File(snapshotDir));
} catch (IOException e) {
log.error("Failed to create snapshot directory {}", snapshotDir);
done.run(new Status(RaftError.EIO, e.toString()));
return;
}
for (RaftTaskHandler taskHandler : taskHandlers) {
try {
FileUtils.deleteDirectory(new File(snapshotDir));
FileUtils.forceMkdir(new File(snapshotDir));
} catch (IOException e) {
log.error("Failed to create snapshot directory {}", snapshotDir);
KVOperation op = KVOperation.createSaveSnapshot(snapshotDir);
taskHandler.invoke(op, null);
log.info("Raft onSnapshotSave success");
} catch (PDException e) {
log.error("Raft onSnapshotSave failed. {}", e.toString());
done.run(new Status(RaftError.EIO, e.toString()));
return;
}
for (RaftTaskHandler taskHandler : taskHandlers) {
try {
KVOperation op = KVOperation.createSaveSnapshot(snapshotDir);
taskHandler.invoke(op, null);
log.info("Raft onSnapshotSave success");
} catch (PDException e) {
log.error("Raft onSnapshotSave failed. {}", e.toString());
done.run(new Status(RaftError.EIO, e.toString()));
}
}
}
} catch (Exception e) {
log.error("failed to save snapshot", e);
done.run(new Status(RaftError.EIO, e.toString()));
return;
} finally {
lock.unlock();
}

MetadataService.getUninterruptibleJobs().submit(() -> {
lock.lock();
try {
// compress
try {
compressSnapshot(writer);
FileUtils.deleteDirectory(new File(snapshotDir));
} catch (Exception e) {
log.error("Failed to delete snapshot directory {}, {}", snapshotDir,
e.toString());
done.run(new Status(RaftError.EIO, e.toString()));
return;
}
compressSnapshot(writer);
FileUtils.deleteDirectory(new File(snapshotDir));
done.run(Status.OK());
log.info("snapshot save done");
} catch (Exception e) {
log.error("failed to save snapshot", e);
log.error("Failed to compress snapshot directory {}, {}", snapshotDir,
e.toString());
done.run(new Status(RaftError.EIO, e.toString()));
} finally {
lock.unlock();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -56,4 +56,11 @@ public interface HgKVStore {
List<KV> scanRange(byte[] start, byte[] end);

void close();

/**
* Wait until a local read sees every write committed before the call. A store that
* is not replicated has nothing to wait for.
*/
default void waitReadIndex() throws PDException {
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,15 @@ public void close() {
store.close();
}

/**
* Reads above are served from the local state, which on a new leader can trail the
* entries its predecessor committed; this waits for them to be applied
*/
@Override
public void waitReadIndex() throws PDException {
this.engine.waitReadIndex();
}

/**
* Need to walk the real operation of Raft
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,12 @@ raft:
address: $RAFT_ADDRESS$
# raft cluster
peers-list: $RAFT_PEERS_LIST$
# Timeout for opening a connection to a peer, in milliseconds; an election waits this
# long for each peer that accepts the connection but does not answer
rpc-connect-timeout: 1000
# Timeout for sending a snapshot to a follower that fell behind the log, in milliseconds;
# 0 uses rpc-timeout
rpc-install-snapshot-timeout: 0
# The interval between snapshot generation, in seconds
snapshotInterval: 300
metrics: true
Expand Down
6 changes: 6 additions & 0 deletions hugegraph-pd/hg-pd-service/src/main/resources/application.yml
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,12 @@ raft:
peers-list: 127.0.0.1:8610
# The read and write timeout period of the raft rpc, in milliseconds
rpc-timeout: 10000
# Timeout for opening a connection to a peer, in milliseconds; an election waits this
# long for each peer that accepts the connection but does not answer
rpc-connect-timeout: 1000
# Timeout for sending a snapshot to a follower that fell behind the log, in milliseconds;
# 0 uses rpc-timeout
rpc-install-snapshot-timeout: 0
# The interval between snapshot generation, in seconds
snapshotInterval: 300
metrics: true
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,12 +17,16 @@

package org.apache.hugegraph.pd.core;

import org.apache.hugegraph.pd.core.meta.IdMetaStoreReadIndexTest;
import org.apache.hugegraph.pd.core.meta.MetadataKeyHelperTest;
import org.apache.hugegraph.pd.core.store.HgKVStoreImplTest;
import org.apache.hugegraph.pd.raft.IpAuthHandlerTest;
import org.apache.hugegraph.pd.raft.RaftEngineIpAuthIntegrationTest;
import org.apache.hugegraph.pd.raft.RaftEngineLeaderAddressTest;
import org.apache.hugegraph.pd.raft.RaftEngineReadIndexTest;
import org.apache.hugegraph.pd.raft.RaftEngineReadinessTest;
import org.apache.hugegraph.pd.raft.RaftEngineRpcTimeoutTest;
import org.apache.hugegraph.pd.raft.RaftStateMachineSnapshotTest;
import org.junit.runner.RunWith;
import org.junit.runners.Suite;

Expand All @@ -31,6 +35,7 @@
@RunWith(Suite.class)
@Suite.SuiteClasses({
MetadataKeyHelperTest.class,
IdMetaStoreReadIndexTest.class,
HgKVStoreImplTest.class,
PDConfigTest.class,
ConfigServiceTest.class,
Expand All @@ -45,6 +50,9 @@
RaftEngineIpAuthIntegrationTest.class,
RaftEngineLeaderAddressTest.class,
RaftEngineReadinessTest.class,
RaftEngineReadIndexTest.class,
RaftStateMachineSnapshotTest.class,
RaftEngineRpcTimeoutTest.class,
// StoreNodeServiceTest.class,
})
@Slf4j
Expand Down
Loading
Loading