RxJava3:基于ID共享Flowable时,结合share操作符与线程安全doFinally操作如何避免死锁及竞态条件?
Core Requirements
- Share a single Flowable instance for all subscribers of the same ID
- Allow individual subscribers to cancel their subscription at any time
- When all subscribers for an ID have canceled, a new subscription request for that ID should create a fresh Flowable instance
The Deadlock & Race Condition Problem
Your initial implementation used a synchronized method alongside RxJava's share() operator, which led to a deadlock. Here's why:
- RxJava's
FlowableRefCount(returned byshare()) uses internal synchronization in itssubscribeActualandcancelmethods - Your
doFinallycallback was executingremoveSubscription(which held theSharedFlowProviderlock) inside theFlowableRefCountlock scope - This created a lock-order inversion: one thread held the
FlowableRefCountlock while waiting for theSharedFlowProviderlock, and another held theSharedFlowProviderlock while waiting for theFlowableRefCountlock
Your attempt to switch to ConcurrentHashMap fixed the deadlock but introduced a race condition: if a subscription was canceled immediately after being created, the Flowable could be removed from the map before your final computeIfAbsent call, leaving an unused Flowable in the map (or worse, a new Flowable being created unnecessarily).
Solution: Atomic Map Operations & Safe Cleanup
The key is to avoid cross-lock dependencies and use atomic operations to manage the Flowable instances in the map. Here's a robust implementation:
class SharedFlowProvider(private val flowProvider: FlowProvider) { // Use ConcurrentHashMap for thread-safe, lock-free map operations private val eventFlowById = ConcurrentHashMap<String, Flowable<ComputedProperties>>() fun subscribeToProperties(subscriber: CancellableSubscriber<ComputedProperties>, id: String) { var flow = eventFlowById[id] if (flow == null) { // Double-checked pattern with putIfAbsent to avoid duplicate Flowable creation val newFlow = buildFlowable(id) val existingFlow = eventFlowById.putIfAbsent(id, newFlow) flow = existingFlow ?: newFlow } // Subscribe the caller to the shared Flowable flow.subscribe(subscriber) } private fun buildFlowable(id: String): Flowable<ComputedProperties> { return flowProvider.getFlowable(id) .doFinally { // Only remove the Flowable from the map if it's still the same instance we created // This prevents accidentally removing a new Flowable created by a concurrent subscription eventFlowById.remove(id) { currentMapFlow -> currentMapFlow === this@buildFlowable } } // share() = publish().refCount(1): maintains a single Flowable until all subscribers cancel .share() } }
How This Works
- No External Synchronization: We rely on
ConcurrentHashMap's thread-safe operations instead ofsynchronizedmethods, eliminating lock-order inversion risks. - Atomic Flowable Creation: The
putIfAbsentmethod ensures only one Flowable is created per ID, even with concurrent subscription requests. - Safe Cleanup: The
doFinallycallback usesConcurrentHashMap'sremove(key, value)overload, which only removes the entry if the stored value matches the Flowable instance we're cleaning up. This avoids removing a new Flowable that might have been created by a concurrent subscription request right after the last subscriber canceled. - Meets All Requirements:
- Subscribers of the same ID share the same Flowable instance
- Individual subscribers can cancel freely
- When all subscribers cancel, the
doFinallycallback removes the Flowable from the map, so the next subscription request creates a fresh instance
Key Notes
- The
share()operator already handles ref-counting: it keeps the upstream Flowable alive as long as there's at least one subscriber, and cleans it up when the last subscriber cancels. - Using
ConcurrentHashMap's atomic operations ensures we never have conflicting updates to the map, even under high concurrency.
内容的提问来源于stack exchange,提问作者Vlad Piscu
相关产品推荐
相关产品推荐

