如何为RxJava中flatMap的processItem处理逻辑添加synchronized同步块
嘿,我一眼就看出问题所在了:你当前的synchronized块只在创建processItem返回的Completable实例时持有锁,而processItem的实际业务逻辑是在subscribeOn(Schedulers.io())指定的IO线程执行的——等IO线程开始处理的时候,锁早就被释放了,相当于同步完全没生效!
要解决这个问题,核心是让synchronized块真正包裹住processItem的整个处理过程,而不是仅仅包裹Completable的创建步骤。这里有两种靠谱的方案:
方案1:把同步逻辑内嵌到processItem中
最简洁的方式是直接修改processItem方法,把锁加在实际工作逻辑的外层,确保处理全程受锁保护:
public Completable processItem(Item item) { return Completable.fromRunnable(() -> { // 在这里加锁,让整个业务处理都在锁的保护下 synchronized(LockManager.getInstance().getLockObject(item.id)) { // 替换成你的实际业务代码,比如数据库操作、外部API调用等 doActualBusinessWork(item); } }).subscribeOn(Schedulers.io()); }
之后调用的时候就不需要在flatMap里额外加锁了,直接链式调用即可:
Observable.fromIterable(getSourceData()) .flatMapCompletable(this::processItem)
这种方式的好处是逻辑内聚,processItem的所有调用方都不需要关心同步细节,维护起来更省心。
方案2:用Completable.defer延迟创建,绑定锁到执行线程
如果不想修改processItem的原有实现,可以用Completable.defer()来延迟Completable的创建过程,让锁的获取和processItem的执行都在指定的IO线程中完成:
Observable.fromIterable(getSourceData()) .flatMapCompletable(item -> { return Completable.defer(() -> { // defer的回调会在subscribeOn指定的IO线程执行,所以锁在这里获取 synchronized(LockManager.getInstance().getLockObject(item.id)) { // 返回processItem的Completable,此时整个处理过程都会持有锁 return processItem(item); } }).subscribeOn(Schedulers.io()); })
⚠️ 注意:如果你的processItem内部已经使用了subscribeOn切换线程,一定要把那部分去掉!线程切换会导致锁被提前释放,同步逻辑失效。统一在外部flatMap中指定线程就好。
另外提一句:你用Observable.interval定时触发的话,它默认在computation线程执行上游逻辑,但这完全不影响我们的同步方案——因为我们已经把锁和处理逻辑绑定到了IO线程,上游线程的切换不会干扰锁的有效性。
内容的提问来源于stack exchange,提问作者elnino

