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

使用Retrofit2与RxAndroid对接Spring WebFlux:Android端反序列化配置疑问

搞定Retrofit2 + RxAndroid对接Spring WebFlux Flux接口的反序列化问题

嘿,我来帮你解决这个对接的痛点!Spring WebFlux的Flux<Quote>返回的是流式响应,不是普通的一次性JSON数组,所以常规的Retrofit反序列化配置肯定不好使,得针对性调整。下面一步步给你讲清楚:

1. 先搞懂服务端的响应格式

你的Spring WebFlux接口返回Flux<Quote>,默认会输出换行分隔的JSON对象(Content-Type是application/stream+json),大概长这样:

{"id":1,"content":"Stay hungry, stay foolish"}
{"id":2,"content":"Innovation distinguishes between a leader and a follower"}
...

不是把所有Quote打包成一个JSON数组返回,而是逐个输出,所以Retrofit默认的“一次性解析整个响应体”逻辑完全不适用,得处理逐行的流式数据。

2. 添加必要的依赖

首先确保你的build.gradle里有这些依赖(版本可以根据你的项目调整):

// Retrofit核心库
implementation 'com.squareup.retrofit2:retrofit:2.9.0'
// RxJava3适配器(用RxJava2的话换adapter-rxjava2)
implementation 'com.squareup.retrofit2:adapter-rxjava3:2.9.0'
// Gson转换器(用Moshi的话换converter-moshi)
implementation 'com.squareup.retrofit2:converter-gson:2.9.0'
// RxAndroid
implementation 'io.reactivex.rxjava3:rxandroid:3.0.2'
// OkHttp核心(Retrofit底层依赖)
implementation 'com.squareup.okhttp3:okhttp:4.9.3'

3. 配置Retrofit和接口定义

第一步:配置OkHttpClient和Retrofit

流式响应不能设置读取超时(否则会中途断开),而且要确保用IO线程处理:

// 构建OkHttpClient,关闭读取超时,添加日志拦截器方便调试
OkHttpClient okHttpClient = new OkHttpClient.Builder()
        .connectTimeout(10, TimeUnit.SECONDS)
        .readTimeout(0, TimeUnit.SECONDS) // 流式响应必须设置为0,避免超时断开
        .addInterceptor(new HttpLoggingInterceptor().setLevel(HttpLoggingInterceptor.Level.BODY)) // 可选,用来查看响应内容
        .build();

// 构建Retrofit实例
Retrofit retrofit = new Retrofit.Builder()
        .baseUrl("http://你的服务端地址/") // 替换成你的实际地址
        .client(okHttpClient)
        .addConverterFactory(GsonConverterFactory.create())
        .addCallAdapterFactory(RxJava3CallAdapterFactory.create())
        .build();

第二步:定义API接口

必须添加@Streaming注解告诉Retrofit不要一次性加载整个响应体,然后先返回Flowable<ResponseBody>,后续再逐行解析:

public interface QuoteApiService {
    @GET("quotes-reactive-paged")
    @Streaming // 关键:流式处理响应体
    @Headers("Accept: application/stream+json") // 指定接收流式JSON格式
    Flowable<ResponseBody> getQuotesReactivePagedRaw(
            @Query("page") int page,
            @Query("size") int size
    );
}

4. 实现流式反序列化逻辑

现在在调用接口的时候,把ResponseBody转换成字符流,逐行读取并反序列化为Quote对象:

// 创建API服务实例
QuoteApiService apiService = retrofit.create(QuoteApiService.class);

// 发起请求并处理流式数据
apiService.getQuotesReactivePagedRaw(0, 10) // 传入page和size参数
        .subscribeOn(Schedulers.io()) // 必须在IO线程处理流读取
        .flatMap(responseBody -> {
            // 将ResponseBody转换成BufferedReader,逐行读取
            BufferedReader reader = new BufferedReader(
                    new InputStreamReader(responseBody.byteStream(), StandardCharsets.UTF_8)
            );
            // 将行迭代器转换成Flowable
            return Flowable.fromIterable(() -> reader.lines().iterator())
                    // 取消订阅时关闭流,避免资源泄漏
                    .doOnCancel(() -> {
                        try {
                            reader.close();
                        } catch (IOException e) {
                            e.printStackTrace();
                        }
                    });
        })
        .map(line -> {
            // 用Gson将每行JSON反序列化为Quote对象
            return new Gson().fromJson(line, Quote.class);
        })
        .observeOn(AndroidSchedulers.mainThread()) // 切换到主线程更新UI
        .subscribe(
                quote -> {
                    // 处理每个返回的Quote对象,比如更新UI
                    Log.d("QuoteReceived", "ID: " + quote.getId() + ", Content: " + quote.getContent());
                },
                throwable -> {
                    // 处理请求错误
                    Log.e("QuoteError", "请求失败", throwable);
                }
        );

5. 关键注意事项

  • @Streaming注解不能少:如果不加,Retrofit会把整个流式响应加载到内存,不仅失去流式的意义,数据量大时还会OOM。
  • 关闭读取超时:流式响应是持续发送数据的,设置读取超时会导致连接中途断开,所以必须设为0。
  • Quote类要匹配JSON字段:确保你的Quote实体类字段和服务端返回的JSON字段一致,必要时用@SerializedName注解映射,比如:
public class Quote {
    @SerializedName("id")
    private int id;
    @SerializedName("content")
    private String content;

    // Getter和Setter方法
    public int getId() { return id; }
    public void setId(int id) { this.id = id; }
    public String getContent() { return content; }
    public void setContent(String content) { this.content = content; }
}
  • 资源泄漏防护:取消订阅时一定要关闭流,避免IO资源浪费。

6. 如果服务端返回的是SSE格式

如果你的Spring WebFlux接口配置了返回SSE(Content-Type是text/event-stream),响应会是data: {"id":1,...}这种格式,那你需要用OkHttp的SSE扩展来处理,添加依赖:

implementation 'com.squareup.okhttp3:okhttp-sse:4.9.3'

然后自定义转换器解析data:字段里的JSON内容,逻辑类似,只是多一步提取data:后面的内容再反序列化。


内容的提问来源于stack exchange,提问作者James Black

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 11:14:38