You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何仅重试耗时<1s的5xx HTTP响应请求?(RxJava1/Retrofit1/Spring4)

问题描述

技术栈

  • rx-java 1.1.x
  • retrofit 1.9.x
  • spring 4.3.x

背景

刚接触rx-java,我的服务A有一个API端点会频繁调用服务B的搜索接口,但该接口失败率偏高。部分错误是深层服务超时(耗时>30s),但也有大量快速失败的请求(耗时<1s)。

需求

仅重试满足以下条件的失败请求:

  • 请求耗时低于阈值(比如<1s)
  • 返回5xx HTTP状态码

已尝试的无效思路

  • 常规Observable.timeout()不适用:不想中断长耗时请求,仅需重试已返回5xx失败的请求。
  • retry()不适用:不想重试所有失败请求。
  • retryWhen()可能可行,但不清楚如何从Throwable中提取HTTP状态码,以及如何测量Observable调用的耗时。

代码示例

Controller代码

@RestController
@RequestMapping(...)
public class MyController {

    @RequestMapping(method = GET)
    public DeferredResult<MyJsonWrapper> fetchSomething(
            MySearchRequest searchRequest, 
            BindingResult bindingResult, 
            HttpServletRequest request) {

        return new MyDeferredResult(
            serviceB.searchSomething(...)
            .doOnNext( result -> /* log size of search */ ));
    }
}

serviceB.searchSomething(...)返回Observable<MyJsonWrapper>

MyDeferredResult定义

class MyDeferredResult<T> extends DeferredResult<T> {

    public MyDeferredResult(Observable<T> observable) {

        onTimeout(this::handleTimeout);
        ConnectableObservable<T> publication = observable.publish();
        publication.subscribe(this::onNext, this::onError, this::onCompleted);
        publication.connect(subscription -> this.subscription = subscription);
    }

    (...)
    private void handleTimeout() {
        setErrorResult(new MyTimeoutException( /* some info about request */ ));
        subscription.unsubscribe();
    }
}

解决方案

要实现这个需求,需要结合请求耗时统计和错误信息解析,通过retryWhen()实现条件重试,具体步骤如下:

1. 统计请求耗时

在调用服务B的Observable前后记录时间戳,将耗时信息附加到错误Throwable中,这里用自定义异常携带耗时:

// 自定义异常,用于传递请求耗时
public class RequestDurationException extends RuntimeException {
    private final long durationMs;

    public RequestDurationException(long durationMs) {
        this.durationMs = durationMs;
    }

    public long getDurationMs() {
        return durationMs;
    }
}

// 包装Observable,添加耗时统计逻辑
private Observable<MyJsonWrapper> wrapWithTiming(Observable<MyJsonWrapper> source) {
    return Observable.defer(() -> {
        long startTime = System.currentTimeMillis();
        return source
                .doOnError(throwable -> {
                    long duration = System.currentTimeMillis() - startTime;
                    // 将耗时异常附加到原错误中
                    if (throwable instanceof RetrofitError) {
                        throwable.addSuppressed(new RequestDurationException(duration));
                    }
                });
    });
}

2. 解析Retrofit的HTTP错误状态码

Retrofit 1.9.x中失败请求会抛出RetrofitError,可以从中提取响应状态码:

private boolean is5xxError(Throwable throwable) {
    if (!(throwable instanceof RetrofitError)) {
        return false;
    }
    RetrofitError retrofitError = (RetrofitError) throwable;
    Response response = retrofitError.getResponse();
    return response != null && response.getStatus() >= 500 && response.getStatus() < 600;
}

3. 实现retryWhen的条件判断

通过retryWhen()接收错误流,判断每个错误是否满足重试条件,满足则发送重试信号,否则终止:

private Observable<?> retryIfFast5xx(Observable<? extends Throwable> errors) {
    return errors.flatMap(throwable -> {
        // 先判断是否是5xx错误
        if (!is5xxError(throwable)) {
            return Observable.error(throwable);
        }
        // 提取耗时信息
        long duration = -1;
        for (Throwable suppressed : throwable.getSuppressed()) {
            if (suppressed instanceof RequestDurationException) {
                duration = ((RequestDurationException) suppressed).getDurationMs();
                break;
            }
        }
        // 判断耗时是否小于1000ms
        if (duration != -1 && duration < 1000) {
            // 满足条件,延迟100ms后重试(可根据需求调整延迟时间)
            return Observable.timer(100, TimeUnit.MILLISECONDS);
        } else {
            // 不满足条件,抛出原错误终止流
            return Observable.error(throwable);
        }
    });
}

4. 整合到原有代码中

修改Controller中的调用逻辑,将上述步骤整合:

@RequestMapping(method = GET)
public DeferredResult<MyJsonWrapper> fetchSomething(
        MySearchRequest searchRequest, 
        BindingResult bindingResult, 
        HttpServletRequest request) {

    Observable<MyJsonWrapper> searchObservable = serviceB.searchSomething(...)
            .compose(this::wrapWithTiming) // 添加耗时统计
            .retryWhen(this::retryIfFast5xx) // 条件重试
            .doOnNext(result -> /* log size of search */);

    return new MyDeferredResult(searchObservable);
}

额外说明

  • 如果需要限制重试次数,可以在retryWhen()中添加计数器,比如用zipWith(Observable.range(1, maxRetries), (err, count) -> ...)控制重试次数。
  • 延迟重试的时间可根据实际场景调整,避免短时间内大量重试给服务B造成额外压力。

内容的提问来源于stack exchange,提问作者Eel Lee

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.26 03:48:23