RxJava - synchronized块 - 内部源

eln*_*ino 5 synchronized synchronized-block thread-synchronization rx-java2

我想使用synchronized块作为flatMap的源码.但我需要使用此构造进行处理(方法processItem),而不仅仅是在创建内部源时.

每隔5分钟调用一次Observable(例如Observable.interval)

Observable.fromIterable(getSourceData())
.flatMapCompletable(item -> {
  synchronized(LockManager.getInstance().getLockObject(item.id)){
    return processItem(item)
          .subscribeOn(Schedulers.io());
  }
})
Run Code Online (Sandbox Code Playgroud)

我的processsItem方法如下所示:

public Completable processItem(Item item){
  return mApiSource.getItemById(item.id)
    .flatMap(item -> 
      mItemRepository.replace(item)
    )
    .toCompletable();
}
Run Code Online (Sandbox Code Playgroud)

两个内部方法都返回Single.

它是从服务器定期更新的方法的一部分.我需要序列化processItem方法调用(从服务器定期同步项目)和修改项目(更新,删除)的方法,这些方法是从其他项目类调用的(项目同步).

主要问题是定期更新可以重写新更新的项目.

其实我用这个解决方案:

用于更新新项目的链:

public Completable updateItem(Item item){
  return Completable.fromAction(() -> {
        synchronized(LockManager.getInstance().getLockObject(item.id)){
          mApiSource.update(item)
          .flatMap(item -> 
            mItemRepository.replace(item)
          )
          .toCompletable()
          .blockingAwait();
        }
      })
      .subscribeOn(Schedulers.io())
}
Run Code Online (Sandbox Code Playgroud)

定期更新链:

Observable.fromIterable(getSourceData())
    .flatMapCompletable(item -> {
      Completable.fromAction(() ->{
        synchronized(LockManager.getInstance().getLockObject(item.id)){
          processItem(item).blockingAwait();
        }
      })
      .subscribeOn(Schedulers.io())
    });
Run Code Online (Sandbox Code Playgroud)

我知道RxJava解决方案并不清楚.

你知道更好的解决方案吗?