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

RxSwift:如何重复调用Observable直到JSON解码成功

嘿,我太懂你卡了两天的感受了——这种要动态串行调用、还得靠解码结果判断终止的场景,确实容易在RxSwift的操作符里绕晕。别慌,我给你一个精准解决问题的方案:

首先,咱们先明确核心需求:必须等前一次readResult()完全完成,再取下一个片段;累积所有片段直到能成功解码完整JSON,不需要提前知道总长度。

实现步骤

先假设你有个要解码的目标模型(比如叫YourTargetModel),先写个解码的辅助函数:

import RxSwift
import Foundation

// 替换成你实际的模型结构
struct YourTargetModel: Codable {
    // 你的模型字段,比如id、content之类的
}

// 尝试解码累积的字符串,失败就抛错
private func tryDecodeFullJSON(_ accumulated: String) throws -> YourTargetModel {
    guard let data = accumulated.data(using: .utf8) else {
        throw NSError(domain: "JSONError", code: -1, userInfo: [NSLocalizedDescriptionKey: "字符串转Data失败"])
    }
    return try JSONDecoder().decode(YourTargetModel.self, from: data)
}

然后用递归的Observable来实现串行调用和终止判断:

func fetchCompleteJSON() -> Observable<YourTargetModel> {
    var accumulatedContent = ""
    
    // 递归函数:每次读片段、累积、尝试解码,失败就继续读
    func fetchNext() -> Observable<YourTargetModel> {
        return readResult()
            .map { fragment in
                // 把新片段拼到累积字符串里
                accumulatedContent += fragment
                return accumulatedContent
            }
            .flatMap { fullContent in
                do {
                    // 解码成功!直接返回模型,Observable就结束了
                    let model = try tryDecodeFullJSON(fullContent)
                    return Observable.just(model)
                } catch {
                    // 解码失败,继续递归调用,取下一个片段
                    return fetchNext()
                }
            }
    }
    
    return fetchNext()
}

为啥这个方案能解决你的问题?

  • 串行保证:flatMap会老老实实等前一个readResult()的Observable完成(也就是API请求+回调完成),才会处理结果并决定要不要订阅下一次fetchNext(),完全满足你"前一次读完再读下一次"的要求,根本不需要定时器那种不靠谱的东西。
  • 动态终止:不需要提前知道要调用多少次readResult()——只要解码成功,就立刻发射结果并结束;解码失败就自动继续调用,直到成功为止。
  • 错误处理灵活:如果readResult()本身抛出API错误,整个Observable会终止。要是你想给API请求加重试,直接在readResult()后面加.retry(3)就行(数字改成你要的重试次数):
    return readResult()
        .retry(3) // API失败重试3次
        .map { fragment in
            accumulatedContent += fragment
            return accumulatedContent
        }
        ...
    

调用示例

let disposeBag = DisposeBag()

fetchCompleteJSON()
    .subscribe(onNext: { model in
        print("终于拿到完整解码后的模型啦:\(model)")
    }, onError: { error in
        print("出问题了:\(error.localizedDescription)")
    })
    .disposed(by: disposeBag)

这个方案完全避开了你之前踩的坑——既不用timer搞异步竞态,也不用concat硬写调用次数,纯靠RxSwift的响应式逻辑解决问题,应该能直接解决你的困扰。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:57:23