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

如何用RxJava获取各线程最后一个Observable及实现K8s状态轮询?

RxJava Solutions for Your Two Scenarios

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 using zip or merge as 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 get true (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 timeout to 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 sendToK8s method 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 with Observable.just().

内容的提问来源于stack exchange,提问作者FakeAlcohol

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:02:42