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

基于RxJava2实现带进度处理的文件下载技术问询

哥们,你的思路完全在线!因为要持续发射多条进度数据(这是流式多事件场景),而Single只能发射单个事件,所以确实得用Flowable来处理进度流。下面我给你一步步拆解完整实现方案:

一、为什么选Flowable?

进度是持续产生的多事件,虽然Observable也能发多事件,但Flowable天然支持背压——如果进度更新太快,下游UI处理不过来的话,背压能避免事件积压导致的OOM,对于下载这种IO密集型场景,用Flowable更稳妥。

二、完整实现步骤

1. 自定义OkHttp ResponseBody监听下载进度

OkHttp本身没有进度回调,我们得包装原始ResponseBody,在读取流的过程中计算并回调进度:

public class ProgressResponseBody extends ResponseBody {
    private final ResponseBody delegate;
    private final ProgressListener listener;
    private BufferedSource bufferedSource;

    public ProgressResponseBody(ResponseBody delegate, ProgressListener listener) {
        this.delegate = delegate;
        this.listener = listener;
    }

    @Override
    public MediaType contentType() {
        return delegate.contentType();
    }

    @Override
    public long contentLength() {
        return delegate.contentLength();
    }

    @Override
    public BufferedSource source() {
        if (bufferedSource == null) {
            bufferedSource = Okio.buffer(new ProgressSource(delegate.source()));
        }
        return bufferedSource;
    }

    private class ProgressSource extends ForwardingSource {
        long totalBytesRead = 0L;
        final long contentLength = contentLength();

        ProgressSource(Source delegate) {
            super(delegate);
        }

        @Override
        public long read(Buffer sink, long byteCount) throws IOException {
            long bytesRead = super.read(sink, byteCount);
            // 累加已读字节,计算进度(bytesRead=-1表示读取完成)
            totalBytesRead += bytesRead != -1 ? bytesRead : 0;
            float progress = contentLength == 0 ? 1.0f : (float) totalBytesRead / contentLength;
            listener.onProgress(progress);
            return bytesRead;
        }
    }

    public interface ProgressListener {
        void onProgress(float progress);
    }
}

2. 用Interceptor注入进度监听

创建OkHttp拦截器,把原始ResponseBody替换成我们自定义的带进度监听的版本:

public class ProgressInterceptor implements Interceptor {
    private final ProgressListener listener;

    public ProgressInterceptor(ProgressListener listener) {
        this.listener = listener;
    }

    @Override
    public Response intercept(Chain chain) throws IOException {
        Response originalResponse = chain.proceed(chain.request());
        return originalResponse.newBuilder()
                .body(new ProgressResponseBody(originalResponse.body(), listener))
                .build();
    }
}

3. 结合RxJava实现从Url到进度流的转换

这里要把获取Response的Single,转换成能持续发射进度的Flowable,同时完成文件写入逻辑:

public Flowable<Float> downloadFile(String url, File saveFile) {
    // 先初始化进度监听器,后续关联到Flowable的事件发射
    final ProgressListener[] progressListener = new ProgressListener[1];

    // 创建带进度拦截器的OkHttpClient
    OkHttpClient client = new OkHttpClient.Builder()
            .addInterceptor(chain -> {
                Response originalResponse = chain.proceed(chain.request());
                return originalResponse.newBuilder()
                        .body(new ProgressResponseBody(originalResponse.body(), progressListener[0]))
                        .build();
            })
            .build();

    // 创建获取Response的Single
    Single<Response> responseSingle = Single.create(emitter -> {
        Request request = new Request.Builder().url(url).build();
        try {
            Response response = client.newCall(request).execute();
            if (!response.isSuccessful()) {
                emitter.onError(new IOException("请求失败: " + response.code()));
                return;
            }
            emitter.onSuccess(response);
        } catch (IOException e) {
            emitter.onError(e);
        }
    });

    // 将Single转换为Flowable,发射进度事件
    return Flowable.create(emitter -> {
        // 关联进度监听器,把进度转发到Flowable
        progressListener[0] = progress -> {
            if (!emitter.isCancelled()) {
                emitter.onNext(progress);
            }
        };

        // 订阅Single处理下载逻辑
        responseSingle.subscribe(
                response -> {
                    try {
                        // 写入文件到指定路径
                        BufferedSource source = response.body().source();
                        BufferedSink sink = Okio.buffer(Okio.sink(saveFile));
                        source.readAll(sink);
                        // 关闭资源
                        sink.close();
                        source.close();
                        // 完成时发射100%进度,再触发完成事件
                        emitter.onNext(1.0f);
                        emitter.onComplete();
                    } catch (IOException e) {
                        emitter.onError(e);
                    }
                },
                emitter::onError
        );
    }, BackpressureStrategy.LATEST); // 背压策略选LATEST,只保留最新进度,避免事件积压
}

4. 订阅Flowable处理进度和结果

在UI层订阅,切换线程处理下载和UI更新:

downloadFile("https://example.com/large-file.apk", new File(getExternalFilesDir(null), "download.apk"))
        .subscribeOn(Schedulers.io()) // 下载逻辑放在IO线程
        .observeOn(AndroidSchedulers.mainThread()) // 进度回调切换到主线程更新UI
        .subscribe(
                progress -> {
                    // 更新进度条,progress范围0~1
                    progressBar.setProgress((int) (progress * 100));
                },
                error -> {
                    // 处理下载错误
                    Toast.makeText(this, "下载失败: " + error.getMessage(), Toast.LENGTH_SHORT).show();
                },
                () -> {
                    // 下载完成提示
                    Toast.makeText(this, "下载完成", Toast.LENGTH_SHORT).show();
                }
        );
三、关键注意事项
  • 背压策略:选BackpressureStrategy.LATEST是因为进度是高频事件,UI只需要最新的进度值,不需要处理每一次更新,能有效避免事件积压。
  • 线程切换:必须用subscribeOn(Schedulers.io())把下载放在IO线程,避免阻塞主线程;用observeOn(AndroidSchedulers.mainThread())把进度回调切到主线程更新UI。
  • 资源释放:文件写入完成后要关闭source和sink;如果用户取消下载(订阅取消),要确保OkHttp请求也能中断,可以在Flowable的create回调里监听emitter.isCancelled(),触发请求取消。
  • 异常处理:要处理文件写入失败、网络请求失败等异常,避免崩溃。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:26:37