How to enforce only 1 subscriber per multiple instances of a same Single/Observable?

Viewed 150

I have this sync modeled as a Single, and only 1 sync can be running at a time.

I'm trying to subscribe the "job" on a Schedulers.single() which mostly works, but inside the chain there are schedulers hops (to db writes scheduler), which unblocks the natural queue created by single()

Then I looked at flatMap(maxConcurrency=1) but this won't work, as that requires always the same instance. I.e. from what I understand, some sort of a Subject of sync requests, which however is uncomposable as my usecase mostly looks like this

fun someAction1AndSync(): Single<Unit> {
   return someAction1()
     .flatMap { sync() }
}

fun someAction2AndSync(): Single<Unit> {
   return someAction2()
     .flatMap { sync() }
}

...

as you can see, its separate sync Single instances :/

Also note someActionXAndSync should not emit until the sync is also done

Basically I'm looking for coroutines Semaphore

2 Answers

I can think of three ways:

  • use a single thread for whole sync operation (decoupling through queue)
  • use semaphore to protect sync method from entering multiple times (not recommended, because will block callee)
  • fast return, when sync is in progress (AtomicBoolean)

There might by other solutions, which I am not aware of.

fast return, when sync is in progress

Also note someActionXAndSync should not emit until the sync is also done

This solution will not queue up sync requests, but will fail fast. The callee must handle the error appropriately

SyncService

class SyncService {
    val isSync: AtomicBoolean = AtomicBoolean(false)

    fun sync(): Completable {
        return if (isSync.compareAndSet(false, true)) {
            Completable.fromCallable { "" }.doOnEvent { isSync.set(false) }
        } else {
            Completable.error(IllegalStateException("whatever"))
        }
    }
}

Handling

When sync process is already happening, you will receive an onError. This issue must be handled somehow, because the onError will be emitted to the subscriber. Either you are fine with it, or you could just ignore it with onErrorComplete

    fun someAction1AndSync(): Completable {
        return Single.just("")
            .flatMapCompletable {
                sync().onErrorComplete()
            }
    }

use a single thread for whole sync operation

You have to make sure, that the whole sync-process is processed in a single job. When the sync-process is composed of multiple reactive steps on other threads, it could happen, that another sync process is started, while one sync process is already in progress.

How?

You have to have a scheduler with one thread. Each sync invocation must be invoked from given scheduler. The sync operation must complete sync in one running job.

I would use this:

fun Observable<Unit>.forceOnlyOneSubscriber(): Observable<Unit> {
    val subscriberCount = AtomicInteger(0)
    return doOnSubscribe { subscriberCount.incrementAndGet() }
        .doFinally { subscriberCount.decrementAndGet() }
        .doOnSubscribe { check(subscriberCount.get() <= 1) }
}

You can always generify Unit using generics if you need.

Related