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

如何在Android中使用RxJava2实现Firebase Storage图片上传?

用RxJava2重构Firebase Storage图片上传逻辑

我来帮你把原来的嵌套回调代码转换成RxJava2的链式风格,这样代码会更简洁易读,也更容易管理异步操作~

第一步:准备工作

首先确保你的项目已经引入了RxJava2和RxAndroid的依赖,也可以选择引入firebase-storage-rxjava扩展(如果不想加额外依赖,我们也可以手动包装Task成Observable)。

第二步:封装Firebase上传操作为Observable

我们先把Firebase的putFile和putBytes这两个异步操作包装成Observable,这样就能用Rx的操作符来组合它们:

// 封装putFile为Observable
fun uploadFileToStorage(uri: Uri, storagePath: String): Observable<UploadTask.TaskSnapshot> {
    val storageRef = FirebaseStorage.getInstance().reference.child(storagePath)
    return Observable.create { emitter ->
        storageRef.putFile(uri)
            .addOnCompleteListener { task ->
                if (task.isSuccessful) {
                    emitter.onNext(task.result!!)
                    emitter.onComplete()
                } else {
                    emitter.onError(task.exception ?: Exception("Upload failed"))
                }
            }
            .addOnFailureListener { exception ->
                emitter.onError(exception)
            }
    }
}

// 封装putBytes为Observable
fun uploadBytesToStorage(bytes: ByteArray, storagePath: String): Observable<UploadTask.TaskSnapshot> {
    val storageRef = FirebaseStorage.getInstance().reference.child(storagePath)
    return Observable.create { emitter ->
        storageRef.putBytes(bytes)
            .addOnCompleteListener { task ->
                if (task.isSuccessful) {
                    emitter.onNext(task.result!!)
                    emitter.onComplete()
                } else {
                    emitter.onError(task.exception ?: Exception("Upload failed"))
                }
            }
            .addOnFailureListener { exception ->
                emitter.onError(exception)
            }
    }
}

第三步:重构onActivityResult中的逻辑

现在我们可以用RxJava的链式调用来替代原来的嵌套回调,同时处理线程调度(压缩和上传都放到IO线程,UI操作回到主线程),还要记得用CompositeDisposable来管理订阅,防止内存泄漏:

首先在你的Activity/Fragment里声明一个CompositeDisposable:

private val compositeDisposable = CompositeDisposable()

然后重构onActivityResult:

public override fun onActivityResult(requestCode: Int, resultCode: Int, data: Intent?) {
    super.onActivityResult(requestCode, resultCode, data)
    if (requestCode == CropImage.CROP_IMAGE_ACTIVITY_REQUEST_CODE) {
        val result = CropImage.getActivityResult(data)
        if (resultCode == Activity.RESULT_OK) {
            val resultUri = result.uri
            val actualImageFile = File(resultUri.path!!)
            val dialogs = SpotsDialog(this, "upload")
            val imageCompressor = Compressor(this)

            // 显示上传对话框
            dialogs.show()

            // 构建Rx链式操作流
            val uploadDisposable = Observable.fromCallable {
                // 在IO线程执行图片压缩,避免阻塞主线程
                imageCompressor.setMaxWidth(200)
                    .setMaxHeight(200)
                    .setQuality(75)
                    .compressToBitmap(actualImageFile)
            }
            .subscribeOn(Schedulers.io())
            .observeOn(AndroidSchedulers.mainThread())
            .doOnNext { bitmap ->
                // 回到主线程更新头像UI
                profile_image?.setImageBitmap(bitmap)
            }
            .observeOn(Schedulers.io())
            // 先上传原图,同时传递bitmap到下一步
            .flatMap { bitmap ->
                val userId = FirebaseAuth.getInstance().currentUser?.uid ?: throw Exception("User not logged in")
                uploadFileToStorage(resultUri, "profile_images/$userId.jpg")
                    .map { bitmap }
            }
            // 再上传缩略图
            .flatMap { bitmap ->
                val userId = FirebaseAuth.getInstance().currentUser?.uid ?: throw Exception("User not logged in")
                val baos = ByteArrayOutputStream()
                bitmap.compress(Bitmap.CompressFormat.JPEG, 100, baos)
                val thumbnailBytes = baos.toByteArray()
                uploadBytesToStorage(thumbnailBytes, "profile_images/thumbs_images/$userId.jpg")
            }
            .observeOn(AndroidSchedulers.mainThread())
            .subscribe(
                {
                    // 全部上传完成的回调
                    dialogs.dismiss()
                    showMessage("thumbnail uploaded")
                },
                { error ->
                    // 统一处理所有错误(压缩/上传失败都会走到这里)
                    dialogs.dismiss()
                    error.printStackTrace()
                    val errorMsg = when {
                        error.message?.contains("User not logged in") == true -> "User not authenticated"
                        else -> error.message ?: "Upload failed"
                    }
                    showMessage(errorMsg)
                }
            )

            // 将订阅加入CompositeDisposable统一管理
            compositeDisposable.add(uploadDisposable)

        } else if (resultCode == CropImage.CROP_IMAGE_ACTIVITY_RESULT_ERROR_CODE) {
            val error = result.error
            error?.printStackTrace()
            showMessage("Image crop error")
        }
    }
}

第四步:清理订阅避免内存泄漏

最后,在Activity的onDestroy方法里清理CompositeDisposable,防止Activity销毁后异步操作还在执行:

override fun onDestroy() {
    super.onDestroy()
    compositeDisposable.dispose()
}

这样重构的优势

  • 原来的嵌套回调变成线性链式调用,逻辑更清晰,维护成本更低
  • 线程调度统一管理,压缩和上传都在IO线程执行,不会卡顿主线程
  • 错误处理集中统一,不用在每个回调里重复写错误提示逻辑
  • 用CompositeDisposable管理订阅,彻底避免内存泄漏风险

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:00:54