如何利用Spring WebFlux的WebClient聚合两个响应式微服务数据?
Hey there! Let's walk through how to build that third reactive service to aggregate your questions and their corresponding options using Spring WebFlux's WebClient. I'll break this down with practical code examples and best practices to keep things reactive and efficient.
Instead of creating a new WebClient instance every time (which is inefficient), we'll set it up as a singleton bean in a configuration class. This follows Spring's best practices and ensures reuse of resources:
@Configuration public class WebClientConfig { @Bean public WebClient webClient(WebClient.Builder builder) { return builder.build(); } // If you're using service discovery (like Eureka/Nacos), add @LoadBalanced to the builder: // @Bean // @LoadBalanced // public WebClient.Builder loadBalancedWebClientBuilder() { // return WebClient.builder(); // } }
To return a clean, combined response of a question and its options, create a simple DTO class:
public class QuestionWithOptions { private Question question; private List<Option> options; // Constructor, getters, and setters public QuestionWithOptions(Question question, List<Option> options) { this.question = question; this.options = options; } // Add getters for serialization (or use Lombok to simplify) public Question getQuestion() { return question; } public List<Option> getOptions() { return options; } }
Now, let's build the service that calls both your question and option services, then combines their data. We'll use reactive operators like flatMap to handle asynchronous calls for each question, and collectList to bundle options into a list:
@Service public class QuestionAggregationServiceImpl implements QuestionAggregationService { private final WebClient webClient; // Replace these URLs with your actual service endpoints (or use service names if using load balancing) private static final String QUESTION_SERVICE_BASE_URL = "http://question-service/api/questions"; private static final String OPTION_SERVICE_BASE_URL = "http://option-service/api/options"; // Constructor injection (preferred over @Autowired for better testability) public QuestionAggregationServiceImpl(WebClient webClient) { this.webClient = webClient; } @Override public Flux<QuestionWithOptions> getQuestionsWithOptions(String categoryId) { // Step 1: Fetch all questions for the given category return webClient.get() .uri("{baseUrl}?categoryId={categoryId}", QUESTION_SERVICE_BASE_URL, categoryId) .retrieve() .bodyToFlux(Question.class) // Step 2: For each question, fetch its corresponding options .flatMap(question -> webClient.get() .uri("{baseUrl}?questionId={questionId}", OPTION_SERVICE_BASE_URL, question.getId()) .retrieve() .bodyToFlux(Option.class) .collectList() // Convert Flux<Option> to Mono<List<Option>> // Step 3: Combine the question and its options into our DTO .map(options -> new QuestionWithOptions(question, options)) // Add error handling to avoid breaking the entire flux if one option call fails .onErrorResume(error -> Mono.just(new QuestionWithOptions(question, Collections.emptyList()))) ); } }
Key Reactive Concepts Here:
flatMap: This operator lets us take eachQuestionfrom the first flux, make an asynchronous call to the option service, and merge the resultingMono<QuestionWithOptions>back into a singleFlux<QuestionWithOptions>.collectList(): Converts the stream of options (FluxonErrorResume: Gracefully handles failures in the option service call—instead of crashing the entire response, we return the question with an empty list of options.
Finally, create a controller to expose the aggregated data as a REST endpoint:
@RestController @RequestMapping("/api/aggregated") public class QuestionAggregationController { private final QuestionAggregationService aggregationService; public QuestionAggregationController(QuestionAggregationService aggregationService) { this.aggregationService = aggregationService; } @GetMapping("/questions") public Flux<QuestionWithOptions> getQuestionsWithOptions(@RequestParam String categoryId) { return aggregationService.getQuestionsWithOptions(categoryId); } }
To prevent hanging requests if your downstream services are slow, add timeout configurations to your WebClient:
@Bean public WebClient webClient(WebClient.Builder builder) { return builder .clientConnector(new ReactorClientHttpConnector( HttpClient.create() .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 3000) // 3s connect timeout .responseTimeout(Duration.ofSeconds(5)) // 5s response timeout )) .build(); }
内容的提问来源于stack exchange,提问作者Khanh Nguyen

