框架整体设计

核心类 职责
OkHttpClient 配置超时、拦截器、SSL 证书
Retrofit 封装 构建 Retrofit 实例,设置 BaseUrl、Gson 转换器、RxJava 适配器
Service 定义 定义各接口的 HTTP 方法,使用 Flowable<ApiResponse> 作为返回类型
请求工具 提供串行/并发请求、结果订阅、取消请求等能力
业务调用层 面向业务方暴露统一方法,内部调用 NetUtils 组合请求
回调与异常 统一回调接口、业务错误码处理、Toast 提示过滤
辅助工具 JSON 转换、RequestBody 构建、本地存储

设计一个高可用,可扩展的网络请求框架最终的设计理念就是单独写成一个网络请求lib库,相当于网络请求的sdk,对外肯定得提供网络请求库得sdk的初始化方法,里面配置一些基础信心,比如OkHttpClient请求超时,拦截器等基础信息。

网络配置与环境切换

public enum BaseUrlType2 {
    COMMON_BUSINESS;

    public String getBaseUrl() {
        return getHttpsBaseURL();
    }


    public static String getHttpsBaseURL() {
        if (isRelease()) {
            return "http://自己后端接口的生产域名";
        }

        return "http://自己后端接口测试域名";
    }

    public static boolean isRelease() {
        return true;
    }
}

单端配置一个枚举来配置生产和测试的域名,这个只是demo简单的例子,可以随着自身应用的业务复杂度扩展这个统一配置各种环境下的请求域名。开发是需要一些发散性思维,举一反三嘛

还可以将一些统一的后端接口业务错误码统一配置

/**
 * 关于网络状态码进行统一管理
 *
 * @author shan.peng
 */
public enum HttpBusinessCode2 {
    //登录密码错误
    STATE_A33013("A33013");

    private final String code;

    HttpBusinessCode2(String code) {
        this.code = code;
    }

    public String getCode() {
        return code;
    }

    public boolean isMatch(String code) {
        return code != null && Objects.equals(this.code, code);
    }

    /**
     * 通过code获取对应的枚举值
     */
    public static HttpBusinessCode2 fromCode(String code) {
        for (HttpBusinessCode2 businessCode : HttpBusinessCode2.values()) {
            if (businessCode.code.equals(code)) {
                return businessCode;
            }
        }
        return null;
    }
}

比如举个例子,后端返回的A33013代码登录密码错误,后期可以把后端返回的所有业务错误码都匹配到这个枚举当中。

OkHttpClient统一配置

/**
 * 自动链接OkHttp客户端配置类
 * 负责创建和配置OkHttpClient.Builder,提供统一的网络客户端配置
 *
 * @author shan.peng
 */
public class AutoLinkNetWork2 {
    private static final String TAG = "AutoLinkNetWork2";
    public static final String NETWORK_DEFAULT_TAG = "NETWORK_DEFAULT_TAG";
    // 连接超时时间设置为15秒(包括连接、读取、写入超时)
    private static final long TIMEOUT_CONN = 15L;
    private OkHttpClient.Builder globalBuilder = null;

    private AutoLinkNetWork2() {
        globalBuilder = defOkHttpClientBuilder();
    }

    private static class Singleton {
        private static final AutoLinkNetWork2 INSTANCE = new AutoLinkNetWork2();
    }

    public static AutoLinkNetWork2 getInstance() {
        return Singleton.INSTANCE;
    }

    /**
     * 构造一个通用的OkHttpClient.Builder
     * 提供默认的网络客户端配置,包括:
     * 1. 超时时间设置
     * 2. 请求日志记录
     * 3. 公共请求头添加
     *
     * @return 配置好的OkHttpClient.Builder实例
     */
    private OkHttpClient.Builder defOkHttpClientBuilder() {
        /*
         * 全局日志控制配置
         * 创建一个HttpLoggingInterceptor,用于记录HTTP请求和响应的详细信息
         * 日志级别设置为BODY,会记录请求头、请求体、响应头、响应体等信息
         */
        HttpLoggingInterceptor defLoggingInterceptor = new HttpLoggingInterceptor(message ->
                // 使用统一的日志工具记录网络日志,方便统一管理和过滤
                Log.d(NETWORK_DEFAULT_TAG, message));
        // 设置日志级别为BODY,记录完整的请求和响应信息
        defLoggingInterceptor.setLevel(HttpLoggingInterceptor.Level.BODY);
        // 创建OkHttpClient.Builder并添加基本配置
        OkHttpClient.Builder okBuilder = new OkHttpClient.Builder();
        // 连接超时:建立与服务器的连接所需的最长时间
        okBuilder.connectTimeout(TIMEOUT_CONN, TimeUnit.SECONDS);
        // 读取超时:从服务器读取数据的最长等待时间
        okBuilder.readTimeout(TIMEOUT_CONN, TimeUnit.SECONDS);
        // 写入超时:向服务器发送数据的最长等待时间
        okBuilder.writeTimeout(TIMEOUT_CONN, TimeUnit.SECONDS);
        okBuilder.pingInterval(TIMEOUT_CONN, TimeUnit.SECONDS);
        // 添加HeaderInterceptor,用于在请求中添加统一的请求头
        okBuilder.addInterceptor(new HeaderInterceptor(false));
        //添加日志拦截器
        okBuilder.addInterceptor(defLoggingInterceptor);
        return okBuilder;
    }

    public OkHttpClient.Builder getOkhttpClientBuilder() {
        return globalBuilder;
    }

    public static void initSdk(Context context) {
        Log.i(TAG, "jitOkHttp initSDK");
        if (context instanceof Application) {
            Utils.init((Application) context);
        }
        GsonUtilsWrapper.initGsonUtilsWrapper();
        ThrowableManager.initThrowableManager();
        Singleton.INSTANCE.defOkHttpClientBuilder();
    }
}

里面有配置日志拦截器:HttpLoggingInterceptor
有请求头拦截器:HeaderInterceptor
还有一些读写超时基础配置,当然对于这个配置超时时间也可以通过接口配置动态配置请求时间

Retrofit接口定义与构建

定义retrofit的时候可以考虑其中的可扩展性,让上层使用着根据自己的业务走不同配置的网络请求

/**
 * Retrofit构建器抽象基类
 */
public abstract class AbstractRetrofitBuilder {

    private static volatile Retrofit mRetrofit;

    public Retrofit getRetrofit() {
        if (mRetrofit == null) {
            synchronized (this) {
                if (mRetrofit == null) {
                    mRetrofit = createRetrofit();
                }
            }
        }
        return mRetrofit;
    }

    protected Retrofit createRetrofit() {
        return createBuilder().build();
    }

    protected Retrofit.Builder createBuilder() {
        OkHttpClient.Builder builder = okHttpClientBuilder();
        return new Retrofit.Builder()
                .baseUrl(apiBaseUrl())
                .addConverterFactory(GsonConverterFactory.create(GsonUtils.getGson()))
                .addCallAdapterFactory(RxJava3CallAdapterFactory.createWithScheduler(Schedulers.io()))
                .client(builder.build());
    }

    protected OkHttpClient.Builder okHttpClientBuilder() {
        AutoLinkNetWork autoLinkNetWork = AutoLinkNetWork.getInstance();
        autoLinkNetWork.setIsPrivacyPolicy(isPrivacyType());
        return autoLinkNetWork.getOkhttpClientBuilder();
    }

    protected String apiBaseUrl() {
        return baseUrlType().getBaseUrl();
    }

    protected abstract BaseUrlType baseUrlType();

    protected abstract boolean isPrivacyType();

    public static void releaseBase() {
        mRetrofit = null;
    }
}

然后定义一个用于相关业务的实现类

public class AppRetrofitBuilder extends AbstractRetrofitBuilder {
    private static volatile Retrofit mRetrofit;

    public static Retrofit getMRetrofit() {
        if (mRetrofit == null) {
            synchronized (AppRetrofitBuilder.class) {
                if (mRetrofit == null) {
                    mRetrofit = new AppRetrofitBuilder().getRetrofit();
                }
            }
        }
        return mRetrofit;
    }

    @Override
    protected BaseUrlType baseUrlType() {
        return BaseUrlType.COMMON_BUSINESS;
    }

    @Override
    protected boolean isPrivacyType() {
        return true;
    }

    public static void release() {
        mRetrofit = null;
    }
}

如果比如说有个其他模块的业务,可以单独起一个类实现一下这个AbstractRetrofitBuilder 类,根据自己情况走不同的配置。

统一响应封装

既然有请求,那肯定就有响应,那网络请求势必可以封装成一个统一的响应回调ResponseCallback了!下面通过代码的形式来揭开它的面纱

public class ResponseCallback<T> {
    private final NetConfig netConfig;
    private final String TAG = getClass().getSimpleName();

    /**
     * 请求开始状态Listener
     */
    private OnStartListener mStartListenerAction;

    /**
     * 请求成功状态Listener(非空数据)
     */
    private OnSuccessListener<T> mSuccessListenerAction;

    /**
     * 请求成功Listener(允许空数据)
     */
    private OnSuccessNullableListener<T> mSuccessNullableListenerAction;

    /**
     * 请求错误Listener
     */
    private OnErrorListener mErrorListenerAction;

    /**
     * 请求完成Listener(无论成功失败都会执行)
     */
    private OnCompleteListener mCompleteListenerAction;

    public ResponseCallback() {
        this(null);
    }

    public ResponseCallback(NetConfig netConfig) {
        this.netConfig = netConfig;
    }

    /**
     * 请求开始状态Listener
     */
    public interface OnStartListener {
        void invoke();
    }

    /**
     * 请求成功状态Listener(非空数据)
     */
    public interface OnSuccessListener<T> {
        void invoke(T t);
    }

    /**
     * 请求成功Listener(允许空数据)
     */
    public interface OnSuccessNullableListener<T> {
        void invoke(T t);
    }

    /**
     * 请求错误Listener
     */
    public interface OnErrorListener {
        boolean invoke(String code, String msg, Throwable throwable);
    }

    /**
     * 请求完成Listener(无论成功失败都会执行)
     */
    public interface OnCompleteListener {
        void invoke();
    }

    /**
     * 开始请求调用方法
     */
    public void _onStart() {
        if (mStartListenerAction != null) {
            mStartListenerAction.invoke();
        }
    }

    /**
     * 请求成功状态调用方法
     */
    public void _onSuccess(T t) {
        if (mSuccessListenerAction != null) {
            mSuccessListenerAction.invoke(t);
        }
    }

    /**
     * 提供了处理空数据的状态方法调用
     */
    public void _onSuccessNullable(T t) throws ApiNullException {
        if (mSuccessNullableListenerAction != null) {
            mSuccessNullableListenerAction.invoke(t);// 用户处理空值
        } else {
            if (t == null) {
                throw new ApiNullException();// 默认情况下禁止空值
            }
            _onSuccess(t);// 非空时传递给普通成功回调
        }
    }

    /**
     * 请求error回调方法的调用
     */
    public void _onError(String code, String msg, Throwable throwable) {
        if (throwable instanceof CancelException) {
            return;
        }
        if (mErrorListenerAction == null) {
            NetUtils.defThrowableHandle(code, msg, throwable, netConfig);
        } else {
            if (mErrorListenerAction.invoke(code, msg, throwable)) {
                NetUtils.defThrowableHandle(code, msg, throwable, netConfig);
            }
        }
    }

    /**
     * 请求完成(无论成功失败都会执行)调用方法
     */
    public void _onComplete() {
        if (mCompleteListenerAction != null) {
            mCompleteListenerAction.invoke();
        }
    }

    /*********************** 最终使用的 **********************/
    public ResponseCallback<T> onStart(OnStartListener action) {
        this.mStartListenerAction = action;
        return this;
    }

    public ResponseCallback<T> onSuccess(OnSuccessListener<T> action) {
        this.mSuccessListenerAction = action;
        return this;
    }

    public ResponseCallback<T> onSuccessNullable(OnSuccessNullableListener<T> action) {
        this.mSuccessNullableListenerAction = action;
        return this;
    }

    public ResponseCallback<T> onError(OnErrorListener action) {
        this.mErrorListenerAction = action;
        return this;
    }

    public ResponseCallback<T> onComplete(OnCompleteListener action) {
        this.mCompleteListenerAction = action;
        return this;
    }

    // 为了方便使用,提供函数式接口版本(Java 8+)
    @FunctionalInterface
    public interface OnStart {
        void onStart();
    }

    @FunctionalInterface
    public interface OnSuccess<T> {
        void onSuccess(T result);
    }

    @FunctionalInterface
    public interface OnSuccessNullable<T> {
        void onSuccess(T result);
    }

    @FunctionalInterface
    public interface OnError {
        boolean onError(String code, String msg, Throwable throwable);
    }

    @FunctionalInterface
    public interface OnComplete {
        void onComplete();
    }

    // 重载方法以支持Java 8 lambda表达式
    public ResponseCallback<T> onStart(OnStart action) {
        this.mStartListenerAction = action::onStart;
        return this;
    }

    public ResponseCallback<T> onSuccess(OnSuccess<T> action) {
        this.mSuccessListenerAction = action::onSuccess;
        return this;
    }

    public ResponseCallback<T> onSuccessNullable(OnSuccessNullable<T> action) {
        this.mSuccessNullableListenerAction = action::onSuccess;
        return this;
    }

    public ResponseCallback<T> onError(OnError action) {
        this.mErrorListenerAction = action::onError;
        return this;
    }

    public ResponseCallback<T> onComplete(OnComplete action) {
        this.mCompleteListenerAction = action::onComplete;
        return this;
    }
}

有请求开始,请求成功,请求失败,请求完成四种状态。

全局异常处理

包括请求业务的时候我们可以封装一个专门用来处理接口业务异常的场景,比如弹出toast提示用户,给予客户必要的说明。

import android.util.Log;

import com.autolink.net.config.NetConfig;

import kotlin.Pair;

/**
 * 网络请求框架的全局异常处理
 * 对API返回错误码的统一处理
 *
 * @author shan.peng
 */
public final class HttpErrorHandle {
    private static final String TAG = "HttpErrorHandle";
    private static long lastShowLoginPage = 0L;

    // 私有构造器,防止实例化
    private HttpErrorHandle() {
        throw new IllegalStateException("Utility class");
    }

    /**
     * 业务异常处理 - 完整参数版本(核心方法)
     *
     * @param e         异常对象,可为null
     * @param code      错误码字符串,例如"401"、"500"
     * @param message   错误信息,用于显示或记录
     * @param netConfig 网络配置对象,控制是否显示Toast等行为
     */
    public static void handleApiBusinessExceptions(
            Throwable e,
            String code,
            String message,
            NetConfig netConfig) {

        // 封装错误码和消息
        Pair<String, String> pair = null;
        if (code != null && message != null) {
            pair = new Pair<>(code, message);
        }

        //判断是否显示Toast
        boolean showErrorToast;
        if (isNoShowToastCode(code) || ThrowableManager.isLocalNetException(e != null ? e.getClass() : null)) {
            // 如果是指定的错误码则除了配置过的请求外,其他的都不显示错误Toast
            Log.e(TAG, "出现指定错误码或者类型的异常------>Throwable:" + e + "    Code:" + code);
            if (netConfig != null && netConfig.getShowErrorToast2User()) {
                showErrorToast = true;
            } else {
                showErrorToast = false;
            }
        } else {
            // 非指定错误码默认都显示
            showErrorToast = netConfig == null || netConfig.getShowErrorToast2User();
        }

        ThrowableManager.showToast2User(
                new Object(),
                e,    // 异常对象
                pair, // 错误码+信息
                null, // 预留参数
                showErrorToast // 是否显示标志
        );

        // 处理特定业务错误码(如401跳转登录)
        if (code != null) {
            try {
                handleSpecificErrorCodes(code);
            } catch (Exception ex) {
                ex.printStackTrace();
            }
        }
    }

    /**
     * 业务Error处理 - 重载版本(无 e 参数)
     */
    public static void handleApiBusinessExceptions(
            String code,
            String message,
            NetConfig netConfig) {
        handleApiBusinessExceptions(null, code, message, netConfig);
    }

    /**
     * 业务Error处理 - 重载版本(无 message 参数)
     */
    public static void handleApiBusinessExceptions(
            Throwable e,
            String code,
            NetConfig netConfig) {
        handleApiBusinessExceptions(e, code, null, netConfig);
    }

    /**
     * 处理特定错误码
     */
    private static void handleSpecificErrorCodes(String code) {
        //先预留,后期好拓展
    }

    /**
     * 判断是否为不显示Toast的错误码
     */
    private static boolean isNoShowToastCode(String code) {
        //先预留
        return false;
    }

    /**
     * 业务Error处理 - 简化版本(仅错误码和消息)
     */
    public static void handleApiBusinessExceptions(String code, String message) {
        handleApiBusinessExceptions(null, code, message, null);
    }

    /**
     * 业务Error处理 - 简化版本(仅错误码)
     */
    public static void handleApiBusinessExceptions(String code) {
        handleApiBusinessExceptions(null, code, null, null);
    }
}

统一请求工具

上面的一切准备好了,就是最最重要的请求工具类,结合rxjava的特性进行封装。下面直接看具体工具类

/**
 * 网络请求工具类 - 基于RxJava封装的网络请求工具
 * 提供并发请求、串行请求、结果处理、错误处理等功能
 * 采用响应式编程模式,支持异步网络请求管理
 */
public final class NetUtils {
    private static final String TAG = "NetUtils";

    private NetUtils() {
        // 工具类,防止实例化
    }

    /**
     * 并发执行多个网络请求
     * 所有请求并行执行,每个请求独立回调,全部完成后触发onComplete
     *
     * @param requestArray 请求Callable列表,每个请求返回ApiResponse
     * @param callback     响应回调接口
     * @param <T>          响应数据类型
     * @return CompositeDisposable 用于管理所有请求的Disposable集合
     */
    @NotNull
    public static <T> CompositeDisposable doMultipleRequest(
            @NotNull final Consumer<ResponseCallback<T>> callback,
            @NotNull final Callable<ApiResponse<? extends T>>... requestArray) {
        return doMultipleRequest(null, callback, requestArray);
    }

    /**
     * 并发执行多个网络请求(完整版)
     * 可配置网络请求参数,并行执行多个请求
     *
     * @param config       网络配置消费者(可空)
     * @param requestArray 请求Callable列表
     * @param callback     响应回调接口
     * @param <T>          响应数据类型
     * @return CompositeDisposable 用于管理所有请求的Disposable集合
     */
    @NotNull
    public static <T> CompositeDisposable doMultipleRequest(
            @Nullable final Consumer<NetConfig> config,
            @NotNull final Consumer<ResponseCallback<T>> callback,
            @NotNull final Callable<ApiResponse<? extends T>>... requestArray) {

        final CompositeDisposable disposables = new CompositeDisposable();
        final NetConfig netConfig = getNetConfigByConsumer(config);
        final ResponseCallback<T> listener = getRealCallback(callback, netConfig);

        listener._onStart();// 请求开始回调

        // 检查是否为空数组
        if (requestArray.length == 0) {
            listener._onComplete();
            return disposables;
        }

        for (int i = 0; i < requestArray.length; i++) {
            final Callable<ApiResponse<? extends T>> request = requestArray[i];
            final boolean isLast = (i == requestArray.length - 1);// 判断是否为最后一个请求

            // 创建并执行单个请求
            final Disposable disposable = Single.fromCallable(request)
                    .subscribeOn(Schedulers.io())  // IO线程执行网络请求
                    .observeOn(AndroidSchedulers.mainThread()) // 主线程回调
                    .subscribe(new io.reactivex.rxjava3.functions.Consumer<ApiResponse<? extends T>>() {
                        @Override
                        public void accept(ApiResponse<? extends T> response) throws Throwable {
                            handleEach(listener, response); // 处理单个响应
                            if (isLast) {
                                listener._onComplete(); // 最后一个请求完成后回调
                            }
                        }
                    }, new io.reactivex.rxjava3.functions.Consumer<Throwable>() {
                        @Override
                        public void accept(Throwable throwable) throws Throwable {
                            listener._onError(null, null, throwable); // 错误处理
                        }
                    });

            disposables.add(disposable); // 添加到Disposable集合
        }

        return disposables;
    }

    /**
     * 创建串行请求的Flowable流
     * 适用于多个请求之间有依赖关系,需要顺序执行的场景
     *
     * @param requestArray 可变参数,多个请求Callable
     * @param <T>          响应数据类型
     * @return Flowable<ApiResponse> 串行请求的数据流
     */
    @NotNull
    public static <T> Flowable<ApiResponse<? extends T>> doRequestByFlow(
            @NotNull final Callable<ApiResponse<? extends T>>... requestArray) {

        return Flowable.fromArray(requestArray)
                .concatMap(new Function<Callable<ApiResponse<? extends T>>, Flowable<ApiResponse<? extends T>>>() {
                    @Override
                    public Flowable<ApiResponse<? extends T>> apply(Callable<ApiResponse<? extends T>> request) throws Throwable {
                        return Single.fromCallable(new Callable<ApiResponse<? extends T>>() {
                                    @Override
                                    public ApiResponse<? extends T> call() throws Exception {
                                        try {
                                            return request.call();// 执行请求
                                        } catch (Exception e) {
                                            e.printStackTrace();
                                            if (e instanceof HttpException) {
                                                HttpException httpException = (HttpException) e;
                                                return new ApiResponse<>(null, String.valueOf(httpException.code()), e, e.getMessage());
                                            } else {
                                                return new ApiResponse<>(null, null, e, e.getMessage());
                                            }
                                        }
                                    }
                                })
                                .subscribeOn(Schedulers.io()) // 添加:确保网络请求在 IO 线程
                                .toFlowable(); // 转换为Flowable
                    }
                });
    }

    /**
     * 订阅Flowable流并处理响应结果(完整版)
     * 只提取ApiResponse中的data字段,自动切换到主线程回调
     *
     * @param flowable 数据流
     * @param config   网络配置(可空)
     * @param callback 响应回调
     * @param <T>      响应数据类型
     * @return Disposable 用于取消订阅
     */
    @NotNull
    public static <T> Disposable result(
            @NotNull final Flowable<ApiResponse<? extends T>> flowable,
            @Nullable final Consumer<NetConfig> config,
            @NotNull final Consumer<ResponseCallback<T>> callback) {

        final NetConfig netConfig = getNetConfigByConsumer(config);
        final ResponseCallback<T> listener = getRealCallback(callback, netConfig);

        return flowable
                .subscribeOn(Schedulers.io()) // 确保在 IO 线程执行网络请求
                .observeOn(AndroidSchedulers.mainThread()) //切换到主线程回调
                .doOnSubscribe(new io.reactivex.rxjava3.functions.Consumer<Subscription>() {
                    @Override
                    public void accept(Subscription subscription) throws Throwable {
                        listener._onStart();// 订阅开始时回调
                    }
                })
                .doFinally(new io.reactivex.rxjava3.functions.Action() {
                    @Override
                    public void run() throws Throwable {
                        listener._onComplete();// 流结束时回调
                    }
                })
                .subscribe(new io.reactivex.rxjava3.functions.Consumer<ApiResponse<? extends T>>() {
                    @Override
                    public void accept(ApiResponse<? extends T> response) throws Throwable {
                        handleEach(listener, response);// 处理每个响应
                    }
                }, new io.reactivex.rxjava3.functions.Consumer<Throwable>() {
                    @Override
                    public void accept(Throwable throwable) throws Throwable {
                        listener._onError(null, null, throwable);// 错误处理
                    }
                });
    }

    /**
     * 订阅Flowable流并处理响应结果(简化版)
     * 使用默认配置
     *
     * @param flowable 数据流
     * @param callback 响应回调
     * @param <T>      响应数据类型
     * @return Disposable 用于取消订阅
     */
    @NotNull
    public static <T> Disposable result(
            @NotNull final Flowable<ApiResponse<? extends T>> flowable,
            @NotNull final Consumer<ResponseCallback<T>> callback) {
        return result(flowable, null, callback);
    }

    /**
     * 取消网络请求
     *
     * @param disposable 要取消的Disposable对象
     */
    public static void cancelReq(@NotNull Disposable disposable) {
        if (!disposable.isDisposed()) {
            disposable.dispose();// 取消订阅
        }
    }

    /**
     * 从Consumer创建NetConfig对象
     *
     * @param configConsumer 配置消费者
     * @return NetConfig 网络配置对象(可能为空)
     */
    @SuppressLint("NewApi")
    @Nullable
    private static NetConfig getNetConfigByConsumer(@Nullable final Consumer<NetConfig> configConsumer) {
        if (configConsumer != null) {
            final NetConfig netConfig = new NetConfig();
            try {
                configConsumer.accept(netConfig);
            } catch (Throwable throwable) {
                throwable.printStackTrace();
            }
            return netConfig;
        }
        return null;
    }


    /**
     * 订阅Flowable流并获取原始ApiResponse(完整版)
     * 返回整个ApiResponse对象,包含code、message等完整信息
     *
     * @param flowable 数据流
     * @param callback 响应回调
     * @param <T>      响应数据类型
     * @return Disposable 用于取消订阅
     */
    @NotNull
    public static <T> Disposable resultByRawModel(
            @NotNull final Flowable<ApiResponse<? extends T>> flowable,
            @NotNull final Consumer<ResponseCallback<ApiResponse<? extends T>>> callback) {

        final ResponseCallback<ApiResponse<? extends T>> listener = getRealCallback(callback, null);

        return flowable
                .subscribeOn(Schedulers.io()) // 确保在 IO 线程执行网络请求
                .observeOn(AndroidSchedulers.mainThread()) // 修复:切换到主线程回调
                .doOnSubscribe(new io.reactivex.rxjava3.functions.Consumer<Subscription>() {
                    @Override
                    public void accept(Subscription subscription) throws Throwable {
                        listener._onStart();
                    }
                })
                .doFinally(new io.reactivex.rxjava3.functions.Action() {
                    @Override
                    public void run() throws Throwable {
                        listener._onComplete();
                    }
                })
                .subscribe(new io.reactivex.rxjava3.functions.Consumer<ApiResponse<? extends T>>() {
                    @Override
                    public void accept(ApiResponse<? extends T> response) throws Throwable {
                        handleEachByRaw(listener, response);
                    }
                }, new io.reactivex.rxjava3.functions.Consumer<Throwable>() {
                    @Override
                    public void accept(Throwable throwable) throws Throwable {
                        listener._onError(null, null, throwable);
                    }
                });
    }

    /**
     * 处理原始ApiResponse响应
     *
     * @param listener 响应回调
     * @param data     ApiResponse数据
     * @param <T>      响应数据类型
     */
    private static <T> void handleEachByRaw(
            @NotNull ResponseCallback<ApiResponse<? extends T>> listener,
            @NotNull ApiResponse<? extends T> data) {

        if (data.getError() != null) {
            listener._onError(data.getCode(), data.getMessage(), data.getError());
        } else {
            if (data.isSuccess()) {
                final ApiResponse<? extends T> readData;
                try {
                    readData = AutolinkPreconditions.apiParamsRequireNotNull(data);
                } catch (BaseException e) {
                    throw new RuntimeException(e);
                }
                listener._onSuccess(readData);
            } else {
                listener._onError(
                        data.getCode(),
                        data.getMessage(),
                        new ApiBusinessExceptions(data.getCode(), data.getMessage())
                );
            }
        }
    }

    /**
     * 处理单个响应(提取data字段)
     *
     * @param listener 响应回调
     * @param data     ApiResponse数据
     * @param <T>      响应数据类型
     */
    private static <T> void handleEach(
            @NotNull ResponseCallback<? super T> listener,
            @NotNull ApiResponse<? extends T> data) {

        if (data.getError() != null) {
            listener._onError(data.getCode(), data.getMessage(), data.getError());
        } else {
            Log.d(TAG, "code=" + data.getCode());
            if (data.isSuccess()) {
                try {
                    listener._onSuccessNullable(data.getData());
                } catch (ApiNullException e) {
                    throw new RuntimeException(e);
                }
            } else {
                listener._onError(
                        data.getCode(),
                        data.getMessage(),
                        new ApiBusinessExceptions(data.getCode(), data.getMessage())
                );
            }
        }
    }

    /**
     * 创建真实的ResponseCallback实例
     *
     * @param callback  回调消费者
     * @param netConfig 网络配置
     * @param <T>       响应数据类型
     * @return ResponseCallback 回调实例
     */
    @SuppressLint("NewApi")
    @NotNull
    private static <T> ResponseCallback<T> getRealCallback(
            @NotNull final Consumer<ResponseCallback<T>> callback,
            @Nullable final NetConfig netConfig) {

        final ResponseCallback<T> listener;
        listener = new ResponseCallback<>(netConfig);
        try {
            callback.accept(listener);
        } catch (Throwable throwable) {
            throwable.printStackTrace();
        }
        return listener;
    }

    /**
     * 默认异常处理
     *
     * @param code      错误码
     * @param msg       错误信息
     * @param throwable 异常对象
     * @param netConfig 网络配置
     */
    public static void defThrowableHandle(
            String code,
            String msg,
            Throwable throwable,
            NetConfig netConfig) {

        if (throwable instanceof ApiBusinessExceptions) {
            handleApiBusinessExceptions(null, code, msg, netConfig);
        } else {
            if (throwable instanceof CancellationException) {
                return;
            }

            if (throwable instanceof SSLHandshakeException) {
                AutoLinkNetWork.release();
                return;
            }

            handleApiBusinessExceptions(throwable, code, netConfig);
        }
    }
}

框架使用示例

举一个例子,有个隐私协议接口,定义一个DemoCommonService

public interface DemoCommonService {
    //隐私协议相关
    @POST("/hus/privacy/auth")
    Flowable<ApiResponse<PrivacyAuthBean>> authPrivacy(@Body RequestBody body);

    /**
     * 服务工厂 - 静态内部类实现单例模式
     */
    class ServiceFactory {
        private ServiceFactory() {
            // 私有构造,防止外部实例化
        }

        /**
         * 静态内部类实现单例模式
         */
        private static class SingletonHolder {
            private static final DemoCommonService INSTANCE =
                    AppRetrofitBuilder.getMRetrofit().create(DemoCommonService.class);
        }

        /**
         * 获取 CommonService 单例
         */
        public static DemoCommonService getInstance() {
            return SingletonHolder.INSTANCE;
        }
    }
}

现在再来看结合我们网络请求工具类完成具体业务接口的请求。

public void authPrivacyNet(PrivacyAuthRequest privacyAuthRequest, ResponseCommonCallback<PrivacyAuthBean> responseCommonCallback) {
        if (privacyAuthRequest == null) return;
        NetUtils.result(NetUtils.doRequestByFlow(() -> {
            CommonService commonService = CommonService.ServiceFactory.getInstance();
            return commonService.authPrivacy(RequestBodyWrapper.data2RequestBody(privacyAuthRequest)).blockingFirst();
        }), callback -> {
            callback.onStart((ResponseCallback.OnStartListener) () -> {
                if (responseCommonCallback != null) {
                    responseCommonCallback.onStart();
                }
            }).onSuccessNullable((ResponseCallback.OnSuccessNullableListener<PrivacyAuthBean>) o -> {
                if (responseCommonCallback != null) {
                    responseCommonCallback.onSuccess(o);
                }
            }).onError((ResponseCallback.OnErrorListener) (code, msg, throwable) -> {
                if (responseCommonCallback != null) {
                    responseCommonCallback.onError(code, msg, throwable);
                }
                return true;
            });
        });
    }

这样调用这个方法完成整个请求的闭环。

框架亮点总结

统一响应模型与标准化
ApiResponse 封装了业务数据、状态码、消息和异常,强制统一了前后端交互的契约,使得上层业务无需关心底层是成功还是失败,只需处理 isSuccess() 即可。

分层清晰,职责分离
1.统一管理 OkHttp 客户端,BaseUrlType 管理多环境 URL。
2.HeaderInterceptor 集中处理请求头,HttpLoggingInterceptor 统一日志。
3.NetUtils 封装 RxJava 并发/串行请求,ThrowableManager 统一异常处理。
希望给在读的你提供一套网络框架的解决思路。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐