Handle unavailable cluster nodes explicitly
This commit is contained in:
1 parent
2276fee60c
commit
350611d562
3 files changed
+31
-18
No files matched your search
@@ -400,23 +400,21 @@ final class ClusterStore implements ObjectStorage {
|
||||
for (UUID id : segment.replicas()) {
|
||||
int node = nodes.index(id);
|
||||
if (node < 0) continue;
|
||||
try {
|
||||
byte[] candidate = nodes.get(node, segment.id(), segment.length(), segment.hash());
|
||||
if (copy == null) copy = candidate;
|
||||
healthy.add(id);
|
||||
healthyHosts.add(nodes.faultDomain(node, testNodeDomains));
|
||||
} catch (IOException error) { }
|
||||
byte[] candidate = readableReplica(node, segment);
|
||||
if (candidate == null) continue;
|
||||
if (copy == null) copy = candidate;
|
||||
healthy.add(id);
|
||||
healthyHosts.add(nodes.faultDomain(node, testNodeDomains));
|
||||
}
|
||||
if (copy == null) { unrecoverable++; continue; }
|
||||
for (int node : PlacementPolicy.candidates(segment.id(), nodes, testNodeDomains)) {
|
||||
UUID host = nodes.faultDomain(node, testNodeDomains);
|
||||
if (healthyHosts.contains(host)) continue;
|
||||
try {
|
||||
nodes.repair(node, segment.id(), copy, segment.hash());
|
||||
if (repairReplica(node, segment, copy)) {
|
||||
healthy.add(nodes.node(node).id());
|
||||
healthyHosts.add(host);
|
||||
restored++;
|
||||
} catch (IOException error) { }
|
||||
}
|
||||
if (healthyHosts.size() == 3) break;
|
||||
}
|
||||
if (healthyHosts.size() < 3) underReplicated++;
|
||||
@@ -439,6 +437,19 @@ final class ClusterStore implements ObjectStorage {
|
||||
} catch (SQLException error) { throw databaseError(error); }
|
||||
return new RepairReport(scanned, restored, underReplicated, unrecoverable);
|
||||
}
|
||||
|
||||
private byte[] readableReplica(int node, Segment segment) {
|
||||
try { return nodes.get(node, segment.id(), segment.length(), segment.hash()); }
|
||||
catch (IOException unavailable) { return null; }
|
||||
}
|
||||
|
||||
private boolean repairReplica(int node, Segment segment, byte[] copy) {
|
||||
try {
|
||||
nodes.repair(node, segment.id(), copy, segment.hash());
|
||||
return true;
|
||||
} catch (IOException unavailable) { return false; }
|
||||
}
|
||||
|
||||
@Override public void close() {}
|
||||
|
||||
private final class SegmentStream extends InputStream {
|
||||
|
||||
@@ -74,6 +74,11 @@ final class NodeClient {
|
||||
}
|
||||
}
|
||||
|
||||
static NodeIdentity probeIfAvailable(URI url, String token) {
|
||||
try { return probe(url, token); }
|
||||
catch (IOException offline) { return null; }
|
||||
}
|
||||
|
||||
static void validateUrl(URI url) {
|
||||
if (url == null || !"http".equals(url.getScheme()) || url.getHost() == null ||
|
||||
url.getPort() < 1 || url.getRawUserInfo() != null ||
|
||||
@@ -86,12 +91,11 @@ final class NodeClient {
|
||||
Set<UUID> healthy = new HashSet<>();
|
||||
for (int i = 0; i < nodes.size(); i++) {
|
||||
Node node = nodes.get(i);
|
||||
try {
|
||||
NodeIdentity actual = probe(node.url(), token);
|
||||
if (actual.nodeId().equals(node.id()) && actual.hostId().equals(node.hostId()))
|
||||
healthy.add(faultDomain(i, testNodeDomains));
|
||||
if (healthy.size() >= required) return true;
|
||||
} catch (IOException error) { }
|
||||
NodeIdentity actual = probeIfAvailable(node.url(), token);
|
||||
if (actual == null) continue;
|
||||
if (actual.nodeId().equals(node.id()) && actual.hostId().equals(node.hostId()))
|
||||
healthy.add(faultDomain(i, testNodeDomains));
|
||||
if (healthy.size() >= required) return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
@@ -126,9 +126,7 @@ final class NodeRegistry {
|
||||
for (URI url : urls) {
|
||||
NodeClient.Node node = stored.get(url.toString());
|
||||
if (node == null) throw new IOException("Unregistered storage node URL: " + url);
|
||||
NodeIdentity actual = null;
|
||||
try { actual = NodeClient.probe(url, token); }
|
||||
catch (IOException offline) { }
|
||||
NodeIdentity actual = NodeClient.probeIfAvailable(url, token);
|
||||
if (actual != null && (!actual.nodeId().equals(node.id()) || !actual.hostId().equals(node.hostId())))
|
||||
throw new IOException("Storage node identity changed at " + url);
|
||||
configured.add(node);
|
||||
|
||||
Reference in new issue
Block a user