如何用RxSwift将async/await方法包装为Observable<AuthSession>
使用RxSwift包装AWS Amplify的async/await方法返回Observable
我正在使用AWS Amplify库与Cognito服务通信,该库的多数函数已采用新的async/await方式重写。现有如下方法:
func fetchAuthSession() async throws -> AuthSession { return try await Amplify.Auth.fetchAuthSession() }
请问如何通过RxSwift包装上述await调用,使其返回Observable
解决方案
你可以通过Observable.create包装async/await操作,同时处理任务生命周期,确保订阅销毁时能取消异步任务:
import RxSwift import Amplify func fetchAuthSessionObservable() -> Observable<AuthSession> { return Observable.create { observer in // 启动Task执行异步操作 let task = Task { do { let session = try await Amplify.Auth.fetchAuthSession() observer.onNext(session) observer.onCompleted() } catch { observer.onError(error) } } // 返回Disposable,在订阅销毁时取消Task return Disposables.create { task.cancel() } } }
如果需要更通用的包装方式,可以给Observable添加扩展,方便复用:
extension Observable { static func fromAsync<T>(_ operation: @escaping () async throws -> T) -> Observable<T> { return Observable.create { observer in let task = Task { do { let result = try await operation() observer.onNext(result) observer.onCompleted() } catch { observer.onError(error) } } return Disposables.create { task.cancel() } } } }
使用这个扩展时,代码会更简洁:
func fetchAuthSessionObservable() -> Observable<AuthSession> { return Observable.fromAsync(Amplify.Auth.fetchAuthSession) }
说明
- 异步操作通过
Task执行,正确处理async/await的抛出错误 - 成功时发送
onNext传递结果,调用onCompleted结束序列 - 失败时通过
onError传递错误信息 - 返回的
Disposable会在订阅取消或销毁时取消Task,避免不必要的资源消耗和内存泄漏
内容的提问来源于stack exchange,提问作者Florian Ldt
相关产品推荐
相关产品推荐

