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

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格式的消息:

  1. 添加Gradle依赖:
implementation 'com.squareup.retrofit2:sse:2.9.0'
  1. 修改Retrofit构建器:
Retrofit retrofit = new Retrofit.Builder()
        .baseUrl(Values.BASE_URL)
        .addCallAdapterFactory(RxJava2CallAdapterFactory.create())
        .addCallAdapterFactory(SseCallAdapterFactory.create())
        .build();
  1. 接口可以直接返回Flowable<Event>,不用手动解析流。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 08:27:36