基于RxJava实现API串行调用:先获取餐厅ID再逐个查询详情
Awesome question — this is such a common pattern in API-driven apps, and RxJava is perfect for handling this kind of sequential-then-parallel workflow. Let’s break down a solid RxJava implementation, plus some optimizations and alternatives you should consider.
First, let’s assume you have an API interface that looks like this (adjust based on your actual setup):
interface RestaurantApi { // Returns an Observable emitting a list of restaurant IDs Observable<List<String>> getRestaurantIds(); // Returns details for a single restaurant ID Observable<RestaurantDetails> getRestaurantDetails(String restaurantId); }
Core Implementation
Here’s a clean, robust implementation that handles the sequential fetch of IDs, then parallel fetch of details with proper error handling and concurrency control:
// Inject or initialize your RestaurantApi instance RestaurantApi restaurantApi = ...; restaurantApi.getRestaurantIds() // Convert the list of IDs into a stream of individual ID emissions .flatMapIterable(restaurantIds -> restaurantIds) // Fetch details for each ID, with concurrency control and error handling .flatMap(restaurantId -> restaurantApi.getRestaurantDetails(restaurantId) // Handle failures for individual detail requests without breaking the whole stream .onErrorResumeNext(throwable -> { Log.e("RestaurantFetcher", "Failed to load details for ID: " + restaurantId, throwable); return Observable.empty(); // Skip failed entries }), 5 // Limit concurrent requests to 5 (adjust based on API rate limits) ) // Collect all successful results into a single list (optional, remove if you want to process items one-by-one) .toList() // Subscribe to the stream to trigger execution .subscribe( successfulDetails -> { Log.d("RestaurantFetcher", "Successfully loaded " + successfulDetails.size() + " restaurant details"); // Process the final list here }, globalError -> { Log.e("RestaurantFetcher", "Failed to fetch initial restaurant ID list", globalError); // Handle critical failure (e.g., no internet, API down) } );
Key Optimizations & Notes
- Concurrency Control: The second parameter in
flatMap(5in this example) limits how many detail requests run in parallel. This prevents overwhelming the API or hitting rate limits — adjust this number based on your API’s constraints. - Graceful Error Handling:
onErrorResumeNextensures a single failed detail request doesn’t crash the entire workflow. You can customize this to return fallback data instead of skipping, if needed. - Stream Processing: If you don’t need to wait for all details to load before processing, remove
toList()and handle eachRestaurantDetailsitem directly insubscribe’s onNext callback. - Progress Tracking: Add a
doOnNext(details -> updateProgress())betweenflatMapandtoList()to track individual item completion.
If you’re working in a Kotlin codebase, Coroutines and Flow offer a lighter, more readable alternative to RxJava:
// Assume your API uses suspend functions (retrofit supports this natively) interface RestaurantApi { suspend fun getRestaurantIds(): List<String> suspend fun getRestaurantDetails(restaurantId: String): RestaurantDetails } suspend fun fetchAllRestaurantDetails(): List<RestaurantDetails> = coroutineScope { val ids = restaurantApi.getRestaurantIds() // Fetch details concurrently, with error handling ids.map { id -> async(Dispatchers.IO) { try { restaurantApi.getRestaurantDetails(id) } catch (e: Exception) { Log.e("RestaurantFetcher", "Failed to load details for ID: $id", e) null } } }.awaitAll() // Wait for all requests to finish .filterNotNull() // Remove failed entries }
If your backend offers a batch endpoint (e.g., getRestaurantDetailsBatch(List<String> ids)), use that instead of individual requests. This reduces network overhead, minimizes latency, and is far more efficient than any client-side parallelization. It’s the optimal solution if available.
内容的提问来源于stack exchange,提问作者Anuj Jindal

