【问题标题】:Couchbase SDK 2 : bulk read operations , how to failover to replicasCouchbase SDK 2:批量读取操作,如何故障转移到副本
【发布时间】: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


    【解决方案1】:

    我想我找到了答案:

    Observable<RawJsonDocument> rawDocs = idObs.flatMap((key)->{
                System.out.println("key "+key);
                Observable<RawJsonDocument> rawJsonObs = bucket.async().get(key, RawJsonDocument.class);
    
    
                return rawJsonObs.onErrorResumeNext(new Func1<Throwable, Observable<RawJsonDocument>>() {
    
                    @Override
                    public Observable<RawJsonDocument> call(Throwable t1) {
                        if (t1.getCause() instanceof TimeoutException) { //we have a timeout
                            return bucket.async().getFromReplica(key, ReplicaMode.FIRST, RawJsonDocument.class).first();
                        }
                        throw OnErrorThrowable.from(t1);
                    }
                });
    
            });
    

    【讨论】:

    • 这看起来是正确的做法。我认为更惯用的 rx 是用 Observable.error(t1) 替换 throw ,但那是次要的。很酷的提示:有时你会从 rxjava 得到一个OnErrorThrowable.OnNextValue,知道它通常包含导致问题的值,以便你的代码可以检查它:)
    • 您正在使用first()。如果没有发出任何值,这将抛出,从而导致onError,如果键不存在,这是可能的。在这种情况下,您可能想要另一种行为?其他选项是 takeFirst() 只会导致 empty 可观察或 firstOrDefault(null)发出 null 而不是抛出(或您传递给运算符)。
    • oups,意思是 take(1),而不是 takeFirst()
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-10-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多