如何用RxJava获取各线程最后一个Observable及实现K8s状态轮询?
Alright, let's break down your questions and walk through practical, working solutions for each one.
1. Getting the Last Emission from Each Thread's Observable
If you're dealing with a scenario where multiple threads are emitting events through Observables, and you need to capture the last emission from each individual thread, here's a straightforward approach using groupBy and last():
Example Implementation
import io.reactivex.rxjava3.core.Observable; import io.reactivex.rxjava3.schedulers.Schedulers; import java.util.concurrent.Executors; import javafx.util.Pair; // Or use your own pair class if preferred public class ThreadLastEmissionExample { public static void main(String[] args) { Observable.just("A", "B", "C", "D", "E") .flatMap(item -> // Simulate processing across multiple threads using a fixed thread pool Observable.just(item) .subscribeOn(Schedulers.from(Executors.newFixedThreadPool(3))) .map(data -> new Pair<>(Thread.currentThread().getName(), data)) ) // Group emissions by the thread name that produced them .groupBy(Pair::getKey) // For each thread group, take only the final emission .flatMap(threadGroup -> threadGroup.last()) .subscribe(result -> System.out.printf("Thread %s's last emission: %s%n", result.getKey(), result.getValue()) ); } }
How It Works
groupBy(Pair::getKey)groups all emissions by the thread name they originated from, creating a separate Observable for each thread.threadGroup.last()extracts the final emission from each thread's Observable sequence.- If you have independent Observables running on different threads, you can simply call
last()on each individual Observable and combine results usingzipormergeas needed.
2. Processing a List of Objects with K8s Polling
Your existing polling code is a solid starting point, but let's refine it to meet the requirement of polling until the status returns true (plus add safeguards like timeouts). We'll also cover scaling this to process an entire list of objects.
Step 1: Refine the Polling Method
First, let's fix and enhance your startPolling method:
import io.reactivex.rxjava3.core.Observable; import io.reactivex.rxjava3.schedulers.Schedulers; import java.util.concurrent.TimeUnit; import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class K8sPollingService { private static final Logger log = LoggerFactory.getLogger(K8sPollingService.class); private final CheckSvcStatus checkSvcStatus = new CheckSvcStatus(); // Reuse instance if possible private Observable<Boolean> startPolling(String content) { log.info("Starting polling for content: {}", content); return Observable.interval(2, TimeUnit.SECONDS) // Defer ensures we create a fresh check operation for each interval tick .flatMap(intervalCount -> Observable.defer(() -> // Run the status check on an IO thread to avoid blocking the scheduler Observable.just(checkSvcStatus.check(content)) .subscribeOn(Schedulers.io()) .retry(2) // Retry transient check failures up to 2 times )) // Stop polling immediately once we get a true result .takeUntil(checkResult -> checkResult) // Return the last result (will be true if stopped via takeUntil) .lastOrDefault(false) // Add a timeout to prevent infinite polling (adjust duration to fit your needs) .timeout(5, TimeUnit.MINUTES, Observable.just(false)); } }
Key Improvements
- Removed
take(3)since we need to poll until we gettrue(not just 3 iterations). - Used
subscribeOn(Schedulers.io())to run the status check on an IO thread, preventing scheduler thread blocking. - Added
retry(2)to handle temporary failures in the status check. - Added
timeoutto avoid infinite polling if the status never transitions to true.
Step 2: Process the Entire List of Objects
Now, let's process your List<Object>—you can choose between serial processing (one object at a time) or parallel processing (multiple objects concurrently):
Serial Processing (One After Another)
Use this if you need sequential processing:
import io.reactivex.rxjava3.core.Observable; import java.util.List; public class ObjectProcessor { private final K8sPollingService pollingService = new K8sPollingService(); private static final Logger log = LoggerFactory.getLogger(ObjectProcessor.class); public void processObjects(List<Object> objList) { Observable.fromIterable(objList) // Process one object at a time (concurrency = 1) .flatMap(obj -> sendToK8s(obj) .flatMap(content -> pollingService.startPolling(content)) .doOnNext(result -> { if (result) { log.info("Successfully processed object: {}", obj); } else { log.error("Processing failed for object: {}", obj); } }), 1) .subscribe( ignored -> {}, // No need to handle individual results here error -> log.error("Global processing error occurred", error), () -> log.info("All objects have been processed!") ); } // Replace this with your actual K8s send logic private Observable<String> sendToK8s(Object obj) { log.info("Sending object to K8s: {}", obj); // Simulate sending to K8s (return the content needed for polling) return Observable.just(obj.toString()); } }
Parallel Processing (Multiple Objects at Once)
Adjust the concurrency parameter to process multiple objects in parallel:
// Inside processObjects() .flatMap(obj -> sendToK8s(obj) .flatMap(content -> pollingService.startPolling(content)) .doOnNext(result -> { if (result) { log.info("Successfully processed object: {}", obj); } else { log.error("Processing failed for object: {}", obj); } }), 5) // Process up to 5 objects concurrently
Notes
- Ensure your
sendToK8smethod is wrapped as an Observable to properly handle errors and asynchronous logic. - If
CheckSvcStatus.check()is already asynchronous, return its Observable directly instead of wrapping it withObservable.just().
内容的提问来源于stack exchange,提问作者FakeAlcohol

