You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RxJava3:基于ID共享Flowable时,结合share操作符与线程安全doFinally操作如何避免死锁及竞态条件?

How to Safely Share Flowables by ID with Subscriber Cleanup Without Deadlocks and Race Conditions

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:

  1. RxJava's FlowableRefCount (returned by share()) uses internal synchronization in its subscribeActual and cancel methods
  2. Your doFinally callback was executing removeSubscription (which held the SharedFlowProvider lock) inside the FlowableRefCount lock scope
  3. This created a lock-order inversion: one thread held the FlowableRefCount lock while waiting for the SharedFlowProvider lock, and another held the SharedFlowProvider lock while waiting for the FlowableRefCount lock

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

  1. No External Synchronization: We rely on ConcurrentHashMap's thread-safe operations instead of synchronized methods, eliminating lock-order inversion risks.
  2. Atomic Flowable Creation: The putIfAbsent method ensures only one Flowable is created per ID, even with concurrent subscription requests.
  3. Safe Cleanup: The doFinally callback uses ConcurrentHashMap's remove(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.
  4. Meets All Requirements:
    • Subscribers of the same ID share the same Flowable instance
    • Individual subscribers can cancel freely
    • When all subscribers cancel, the doFinally callback 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.29 07:07:29