【发布时间】:2015-03-25 02:42:43
【问题描述】:
我们正在重构从 Couchbase Client 2 迁移到新的 CouchBase SDK 2 的基准工具。
以前的版本具有以下“批量获取”逻辑来批量检索密钥,如果从主服务器读取失败,则会从“副本”读取故障转移
遗留代码:
List<Map.Entry<String, OperationFuture<CASValue<JsonNode>>>> futures = new java.util.ArrayList<>(keys.size());
for (String key : keys) {
futures.add(new AbstractMap.SimpleImmutableEntry<>(key, client.asyncGets(key, transcoder)));
}
Map<String, Long> casValues = new java.util.HashMap<>(keys.size(), 1f);
for (Map.Entry<String, OperationFuture<CASValue<JsonNode>>> e : futures) {
String key = e.getKey();
OperationFuture<CASValue<JsonNode>> future = e.getValue();
try {
CASValue<JsonNode> casVal = future.get();
if (checkStatus(future.getStatus(), errIfNotFound) == OK) {
result.put(key, JsonByteIterator.asMap(casVal.getValue()));
casValues.put(key, casVal.getCas());
} else {
return ERROR;
}
} catch (RuntimeException te) {
if (te.getCause() instanceof CheckedOperationTimeoutException) { ///READ FROM REPLICA
log.warn("Reading from Replica as reading from master has timed out.");
// This is a timeout operation on a read, let's try to read from slave
ReplicaGetFuture<JsonNode> futureReplica = client.asyncGetFromReplica(key, transcoder);
result.put(key, JsonByteIterator.asMap(futureReplica.get()));
} else {
throw te;
}
}
}
使用新的 Couchbase SDK2
根据新的 Couchbase 2 SDK 文档, http://docs.couchbase.com/developer/java-2.0/documents-bulk.html
我有以下逻辑来批量检索。但我不太确定在哪里添加故障转移机制以使用“副本”读取
bucket.async().getFromReplica(key, ReplicaMode.ALL);
List<RawJsonDocument> rawDocs = idObs.flatMap((keys)->{
Observable<RawJsonDocument> rawJsonObs = bucket.async().get(key, RawJsonDocument.class);
return rawJsonObs;
}).toList()
.toBlocking()
.single();
如何使用基于 RxJava 的新 CouchBase SDK 实现这种“从副本读取”故障转移机制?
【问题讨论】:
标签: java couchbase reactive-programming rx-java