基于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
相关产品推荐
相关产品推荐

