如何在订阅后指定时间优雅关闭冷Observable并处理分页API?
问题背景
现有一段基于RxJava Flowable处理外部分页API(paginatedAPI)的代码:
Flowable<Data> process = Flowable.generate(() -> new State(), new BiConsumer<State, Emitter<Data>>() { @Override void accept(State state, Emitter<Data> dataEmitter) throws Exception { // 从上游服务获取数据,内部会调用dataEmitter.onNext() String pageId = paginatedAPI.get(pageId, dataEmitter); // 处理数据 ... // 更新状态 state.updatePageId(pageId); } }).subscribeOn(Schedulers.from(executor));
该Flowable由generate创建,accept方法仅会在订阅者准备好接收下一条数据时被调用。可完全自定义State的内容,但无法修改paginatedAPI的实现。
需求
- 从订阅开始计时,时间T到达后:
- 遍历所有剩余页面但不将数据发送给订阅者,调用
paginatedAPI.close(); - 将
paginatedAPI.close()返回的数据提供给订阅者。
- 遍历所有剩余页面但不将数据发送给订阅者,调用
- 若订阅者在时间T之前断开连接:
- 遍历所有剩余页面但不将数据发送给订阅者,调用
paginatedAPI.close()。
- 遍历所有剩余页面但不将数据发送给订阅者,调用
困惑
- 如何在Flowable的逻辑中加入“订阅后计时”的控制?
accept方法最多只能调用一次onNext,如何通过多次调用完成paginatedAPI的收尾处理?
解决思路与实现
核心思路
- 计时控制:用
Flowable.defer为每个订阅创建独立状态,结合takeUntil+timer触发计时收尾逻辑; - 收尾逻辑适配:通过状态标记区分正常数据发送和收尾阶段,在
accept方法内用同步循环遍历剩余页面(不违反generate单次onNext的约束),最后仅调用一次onNext发送close返回的数据; - 订阅断开处理:利用
generate的清理回调,在订阅者断开时自动执行剩余页面遍历和close调用。
完整代码实现
// 扩展自定义State,增加阶段标记 class State { private String currentPageId; private boolean isTimeUp; // 标记是否到达计时时间T private boolean isDisposed; // 标记订阅是否已断开 private boolean hasClosed; // 标记是否已执行close操作 // 原有页面ID更新方法 public void updatePageId(String pageId) { this.currentPageId = pageId; } // Getter/Setter方法 public String getCurrentPageId() { return currentPageId; } public boolean isTimeUp() { return isTimeUp; } public void setTimeUp(boolean timeUp) { isTimeUp = timeUp; } public boolean isDisposed() { return isDisposed; } public void setDisposed(boolean disposed) { isDisposed = disposed; } public boolean hasClosed() { return hasClosed; } public void setHasClosed(boolean hasClosed) { this.hasClosed = hasClosed; } } // 最终处理流 Flowable<Data> finalProcess = Flowable.defer(() -> { State state = new State(); return Flowable.generate( // 初始化状态 () -> state, // 核心处理逻辑 (s, emitter) -> { if (s.hasClosed()) { emitter.onComplete(); return; } // 正常数据发送阶段 if (!s.isTimeUp() && !s.isDisposed()) { String newPageId = paginatedAPI.get(s.getCurrentPageId(), emitter); s.updatePageId(newPageId); // 若已到最后一页,直接触发收尾 if (newPageId == null) { s.setTimeUp(true); } } else { // 收尾阶段:遍历剩余页面但不发送数据 Emitter<Data> noopEmitter = new Emitter<Data>() { @Override public void onNext(Data data) {} // 空实现,不推送给订阅者 @Override public void onError(Throwable error) { emitter.onError(error); } @Override public void onComplete() {} }; // 同步遍历所有剩余页面 String pageId = s.getCurrentPageId(); while (pageId != null) { pageId = paginatedAPI.get(pageId, noopEmitter); } // 调用close并推送返回数据 Data closeResult = paginatedAPI.close(); emitter.onNext(closeResult); s.setHasClosed(true); emitter.onComplete(); } }, // 订阅断开时的清理逻辑 s -> { if (!s.hasClosed()) { s.setDisposed(true); Emitter<Data> noopEmitter = new Emitter<Data>() { @Override public void onNext(Data data) {} @Override public void onError(Throwable error) {} @Override public void onComplete() {} }; // 遍历剩余页面并关闭API String pageId = s.getCurrentPageId(); while (pageId != null) { pageId = paginatedAPI.get(pageId, noopEmitter); } paginatedAPI.close(); } } ).takeUntil( // 计时T到达后标记状态,触发收尾 Flowable.timer(T, TimeUnit.SECONDS) .doOnNext(tick -> state.setTimeUp(true)) ); }).subscribeOn(Schedulers.from(executor));
关键细节说明
- 计时触发:
defer保证每个订阅拥有独立的State,timer触发后修改isTimeUp标记,generate检测到标记后进入收尾流程; generate约束适配:收尾阶段的页面遍历是同步循环完成的,最后仅调用一次onNext发送close结果,完全符合accept方法单次onNext的要求;- 资源清理:
generate的第三个参数是订阅断开时的回调,在这里执行剩余页面遍历和close,避免API资源泄漏; - 空Emitter:自定义不推送数据的
noopEmitter,确保收尾阶段的页面数据不会传递给订阅者。
内容的提问来源于stack exchange,提问作者gGwP
相关产品推荐
相关产品推荐

