Android中Retrofit+RxJava的Flowable提前完成问题排查
问题:Spring RestController返回Flowable,Android Retrofit+RxJava仅收到一个值就触发onComplete
你遇到的这个问题我之前也踩过坑,核心原因是默认的HTTP请求是一次性响应,而你想用Spring的Flowable实现持续的数据流推送,这需要用到Server-Sent Events(SSE)——也就是文本流的传输方式,而你当前的代码并没有配置成这种模式。
为什么浏览器能正常工作?
浏览器原生支持解析SSE格式的响应,即使你没严格按照SSE格式返回,浏览器也会把持续输出的内容当成流来处理;但Retrofit默认会把HTTP响应当成一次性的完整数据,一旦读取到第一个返回片段,就会认为请求完成,触发onComplete。
解决方案分两步走:
1. 调整Spring控制器,输出标准SSE格式的流
首先要告诉Spring,这个接口返回的是文本事件流,并且每条消息要符合SSE的格式(以data: 开头,结尾用两个换行符分隔):
import org.springframework.http.MediaType; import java.time.LocalTime; import io.reactivex.Flowable; import java.util.concurrent.TimeUnit; @GetMapping(value = "/api/reactive", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flowable<String> reactive() { return Flowable.interval(1, TimeUnit.SECONDS) // 按照SSE格式包装每条数据 .map(sequence -> "data: \"Flowable-" + LocalTime.now().toString() + "\"\n\n"); }
produces = MediaType.TEXT_EVENT_STREAM_VALUE会让Spring设置响应头Content-Type: text/event-stream,告诉客户端这是一个持续的流。
2. 修改Retrofit配置,支持流式响应
Retrofit默认会一次性读取整个响应体,所以需要给接口添加@Streaming注解,告诉它这是一个流式响应,然后手动解析SSE格式的内容:
第一步:修改Retrofit接口
import io.reactivex.Flowable; import okhttp3.ResponseBody; import retrofit2.http.GET; import retrofit2.http.Streaming; public interface UserRepository { @GET("reactive") @Streaming // 关键:开启流式处理 Flowable<ResponseBody> testReactive(); }
第二步:在订阅时解析SSE流
把ResponseBody的字节流转换成SSE消息,提取出每条数据:
import okhttp3.ResponseBody; import io.reactivex.Flowable; import io.reactivex.schedulers.Schedulers; import androidx.annotation.NonNull; import io.reactivex.subscribers.ResourceSubscriber; import android.util.Log; import android.widget.Toast; import java.io.BufferedReader; import java.io.IOException; import java.io.InputStreamReader; public void useReactive() { Retrofit retrofit = new Retrofit.Builder() .baseUrl(Values.BASE_URL) // 这里不需要Jackson转换器,因为我们要手动解析文本流 .addCallAdapterFactory(RxJava2CallAdapterFactory.create()) .build(); UserRepository userRepository = retrofit.create(UserRepository.class); Flowable<String> reactiveFlow = userRepository.testReactive() .subscribeOn(Schedulers.io()) .flatMap(responseBody -> { // 读取响应体的字节流 BufferedReader reader = new BufferedReader(new InputStreamReader(responseBody.byteStream())); return Flowable.create(emitter -> { String line; try { while ((line = reader.readLine()) != null) { // 过滤SSE的data行,提取有效内容 if (line.startsWith("data: ")) { String data = line.substring(6).trim(); emitter.onNext(data); } } emitter.onComplete(); } catch (IOException e) { emitter.onError(e); } finally { // 关闭流,避免资源泄漏 try { reader.close(); responseBody.close(); } catch (IOException e) { Log.e("ReactiveTest", "流关闭失败", e); } } }, io.reactivex.BackpressureStrategy.BUFFER); }); Disposable disp = reactiveFlow .observeOn(AndroidSchedulers.mainThread()) .subscribeWith(new ResourceSubscriber<String>() { @Override public void onNext(String s) { Log.i("ReactiveTest", "收到数据:" + s); Toast.makeText(authActivity, s, Toast.LENGTH_SHORT).show(); } @Override public void onError(Throwable t) { Log.e("ReactiveTest", "请求出错", t); Toast.makeText(authActivity, "请求出错", Toast.LENGTH_SHORT).show(); } @Override public void onComplete() { Log.i("ReactiveTest", "流已完成"); Toast.makeText(authActivity, "Completed", Toast.LENGTH_SHORT).show(); } }); }
额外优化方案
如果不想手动解析SSE,可以使用Retrofit的官方SSE适配器retrofit2-sse,它会自动帮你处理SSE格式的消息:
- 添加Gradle依赖:
implementation 'com.squareup.retrofit2:sse:2.9.0'
- 修改Retrofit构建器:
Retrofit retrofit = new Retrofit.Builder() .baseUrl(Values.BASE_URL) .addCallAdapterFactory(RxJava2CallAdapterFactory.create()) .addCallAdapterFactory(SseCallAdapterFactory.create()) .build();
- 接口可以直接返回
Flowable<Event>,不用手动解析流。
内容的提问来源于stack exchange,提问作者Vlad
相关产品推荐
相关产品推荐

