RxJava使用retryWhen处理CASMismatch异常时未触发完整查询更新流程
问题根因与修复方案
1. 直接原因:重试配置未生效
你当前的updateDetectionWithRetry方法没有返回应用了retryWhen的Observable,RxJava的冷Observable只有被订阅时才会执行逻辑,你没有返回该Observable也没有主动订阅,相当于重试配置完全没有生效,只会执行一次updateDetection的逻辑,自然不会触发完整的查询+更新重试流程。
2. 修复代码
第一步:修正updateDetectionWithRetry方法
private Observable<String> updateDetectionWithRetry(DetectionFeed detectionFeed, String userId, String detectionPath) { // 新增return,把加了重试逻辑的Observable返回给调用方订阅 return updateDetection(detectionFeed, userId, detectionPath) .retryWhen(retryHandlers.getRetryOnCASMismatchExceptionHandler("Failed to update persisted UserAccount with detection data [" + detectionFeed.toString() + "]")); }
第二步:优化重试逻辑,处理超最大重试次数场景
你当前的重试逻辑在重试次数超过maxAttempts时会直接结束流,不会抛出异常,不符合预期,优化后代码:
public Func1<Observable<? extends Throwable>, Observable<?>> getRetryOnCASMismatchExceptionHandler(String unrecoverableErrorMessage) { return observable -> observable .zipWith(Observable.range(1, maxAttempts + 1), ImmutablePair::of) // 次数范围+1,区分首次执行和重试 .flatMap(pair -> { var throwable = pair.left; var attemptsCounter = pair.right; if (throwable instanceof CASMismatchException && attemptsCounter <= maxAttempts) { // 未超过最大重试次数,走延迟重试 return Observable.timer((long) attemptsCounter * backoffMs, TimeUnit.MILLISECONDS); } // 非CAS异常/超过最大重试次数,直接抛出业务异常 return Observable.error(new RuntimeException(unrecoverableErrorMessage, throwable)); }); }
3. 额外确认项
只要你的userRepo.lookupDetection、userRepo.replaceDetection返回的是Couchbase SDK原生的冷Observable(没有用share、publish等操作符转为热Observable),每次重试触发订阅时都会完整执行「查询带CAS的文档 -> 执行更新」的全流程,完全符合需求。
内容的提问来源于stack exchange,提问作者Etay Ceder
相关产品推荐
相关产品推荐

