在RxKotlin中实现带动态参数的共享Flowable以处理背压
Got it, I totally get what you're trying to do here—you have three buttons triggering the same launchActivity method with different parameters, and that method handles IO operations before returning an Intent. You want to use a shared Flowable with RxKotlin to manage backpressure using the DROP strategy: when one click's IO task is still running, any subsequent clicks get discarded until the current task completes.
Here's a step-by-step implementation that fits your needs:
1. Set up a shared event source for clicks
First, we'll use a PublishSubject to act as a centralized event emitter for all button clicks. This lets us funnel all three buttons' click events into a single stream:
import io.reactivex.rxjava3.core.PublishSubject import io.reactivex.rxjava3.core.Flowable import io.reactivex.rxjava3.schedulers.Schedulers import androidx.appcompat.app.AppCompatActivity import android.os.Bundle import android.widget.Button import android.content.Intent class MainActivity : AppCompatActivity() { // Define a subject to emit click parameters (adjust the type to match your actual params) private val clickSubject = PublishSubject.create<String>() override fun onCreate(savedInstanceState: Bundle?) { super.onCreate(savedInstanceState) setContentView(R.layout.activity_main) // Link each button to emit its specific parameter when clicked findViewById<Button>(R.id.button1).setOnClickListener { clickSubject.onNext("param_for_button1") } findViewById<Button>(R.id.button2).setOnClickListener { clickSubject.onNext("param_for_button2") } findViewById<Button>(R.id.button3).setOnClickListener { clickSubject.onNext("param_for_button3") } // Set up the shared Flowable with backpressure handling setupSharedClickFlowable() } }
2. Configure the shared Flowable with DROP backpressure
Next, we'll convert the subject to a Flowable, apply the DROP backpressure strategy, share the stream, and ensure sequential execution of IO tasks:
private fun setupSharedClickFlowable() { val sharedClickFlowable = clickSubject // Convert to Flowable with DROP backpressure: discard events when downstream is busy .toFlowable(BackpressureStrategy.DROP) // Share the stream so all subscribers (if any) receive the same event sequence .share() sharedClickFlowable // Use concatMap to execute tasks sequentially: wait for one IO task to finish before the next .concatMap { param -> // Wrap your IO-heavy launchActivity logic in a Flowable Flowable.fromCallable { // Replace this with your actual IO operations + Intent creation performLaunchActivityIO(param) } // Run IO tasks on a background thread .subscribeOn(Schedulers.io()) } // Switch back to main thread to handle Intent (since starting activities needs main thread) .observeOn(AndroidSchedulers.mainThread()) .subscribe( { intent -> // Launch the activity once the Intent is ready startActivity(intent) }, { error -> // Handle any errors from IO operations error.printStackTrace() } ) } // Replace this with your actual launchActivity logic that does IO work private fun performLaunchActivityIO(param: String): Intent { // Simulate IO work (e.g., reading a file, making a network call) Thread.sleep(2000) // Return the constructed Intent return Intent(this, TargetActivity::class.java).apply { putExtra("button_param", param) } }
Key Details Explained
BackpressureStrategy.DROP: When the downstream (our IO task) can't keep up with incoming click events, new events are immediately discarded. This is exactly what you want—no queuing, just drop clicks that happen while a task is running.share(): Ensures the Flowable is shared across any potential subscribers (if you add more later), so we don't create duplicate streams or duplicate work.concatMap: Forces sequential execution of tasks. It waits for the current IO operation to complete before subscribing to the next one. Combined withDROP, this means clicks during an active task get thrown away.- Threading:
subscribeOn(Schedulers.io())moves IO work off the main thread to avoid UI freezes, andobserveOn(AndroidSchedulers.mainThread())brings the result back to the main thread for launching the activity.
How the Backpressure Works in Practice
If you click Button 1, the IO task starts. If you quickly click Button 2 and 3 while that task is running, those two click events get dropped by the
DROPstrategy. Only after Button 1's IO finishes and the activity launches will the next valid click (if any) be processed.
内容的提问来源于stack exchange,提问作者Dan 0

