使用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
相关产品推荐
相关产品推荐

