【问题标题】:Room Flowable does not emit data on insertionRoom Flowable 不会在插入时发出数据
【发布时间】:2018-06-07 04:29:27
【问题描述】:

我很难理解 Flowable 在 Room 中的工作原理。我有 Dao 这样的方法

@Insert(onConflict = OnConflictStrategy.REPLACE)
void upsert(List<Site> sites);

@Query("SELECT * FROM site ORDER BY distance ASC")
Flowable<List<Site>> getSites();

我希望每当我调用upsert 订阅者到Flowable 时,getSites() 返回的对象总是会被调用。我的假设是真的吗?

这就是我订阅这个 flowable 的方式

private final Flowable<List<Site>> siteFlowable;
ApiService apiService;
FuelDatabase database;

@Override
public void getSites(boolean showOnlyKeySites) {
    // add sites from cache first, then fetch network -> update cache -> update ui
    disposable = siteFlowable.flatMap(Flowable::fromIterable)
        .filter(site -> site.isValid())
        .buffer(100, TimeUnit.MILLISECONDS, 20)
        .takeUntil(sites -> sites.size() == 0)
        .observeOn(AndroidSchedulers.mainThread())
        .doOnNext(mapView::addPins)
        .subscribe(sites -> {
            Timber.d("Flowable emitted %d items", sites.size());
        }, Timber::e);

    apiService.getSites()
        .map(SiteListResponse::getData)
        .flatMap(Observable::fromIterable)
        .filter(Site::isValidSite)
        .toList().toObservable()
        .subscribe(sites -> {
            Timber.i("Success Fetching %d sites", sites.size());
            database.siteDao().clear();
            database.siteDao().upsert(sites);
        }, throwable -> Timber.e(throwable, "Error fetching sites"));
}

在调用 upsert() 之后不会调用这个 flowable。 API 正在返回有效数据,数据正在输入数据库中。

【问题讨论】:

  • AFAIK 房间查询是无限的,因此 toList 不会对它们起作用。见stackoverflow.com/a/47260768/61158
  • @akarnokd 你在哪里看到toList() 查询?
  • 我认为在 toObservable() 之后你需要使用 flatMap 并隐蔽 upsert 到另一个可流动的流中
  • @Rahul 你能详细说明一下吗?
  • getSites 与 Room 数据库对话,对吗?它返回的 Flowable 是无限的,因此 toList 永远不会完成。

标签: android rx-java2 android-room android-database android-architecture-components


【解决方案1】:

试试这个:

apiService.getSites()
    .map(SiteListResponse::getData)
    .flatMap(result -> 
         Observable.fromIterable(result)
        .filter(Site::isValidSite)
        .toList()
        .toFlowable()
    )
    .subscribe(sites -> {
        Timber.i("Success Fetching %d sites", sites.size());
        database.siteDao().clear();
        database.siteDao().upsert(sites);
    }, throwable -> Timber.e(throwable, "Error fetching sites"));

【讨论】:

  • 这将如何使siteFlowable 发出项目?
  • 另外,.buffer(100, TimeUnit.MILLISECONDS, 20) 可能会导致一个空列表,然后您可以使用该列表来通过 takeUntil() 停止整个流程。
  • @akarnod 是的,这可能导致它只发射一次。还添加了相关定义。
  • 您可能缺少数据。我建议您阅读this 并在各个地方申请doOnNexts,以找出数据消失的位置。
  • buffertakeUntill() 导致了这个问题。我之所以应用它们是因为我不希望我的onNext 被频繁调用。暂时删除它们。
猜你喜欢
  • 2018-08-14
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多