|
Note
|
This is part 3 of the e-commerce microservices series. The code for this post is at tag |
A shop’s product page needs three things: the product information, reviews from buyers, and recommended products. Three different services own them. The client should not have to make three calls and stitch the results together, so we put a service in front that does it: product-composite-service.
By the end of this post you will have:
-
review-service(port 7003) andrecommendation-service(port 7002). -
product-composite-service(port 7000), calling the three services in parallel and merging the results. -
A product page that still works when review or recommendation is down.
GET /product-composite/1
|
product-composite :7000
+----------------+----------------+
| | |
GET /product/1 GET /review?... GET /recommendation?...
product :7001 review :7003 recommendation :7002
The two remaining core services
review-service and recommendation-service have exactly the same structure as product-service, so I only cover what differs.
The contract in the api module
package com.ecommerce.api.core.review;
public record Review(
int productId,
int reviewId,
String author,
String subject,
String content,
String serviceAddress) {
}
package com.ecommerce.api.core.review;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import reactor.core.publisher.Flux;
public interface ReviewService {
@GetMapping(value = "/review", produces = "application/json")
Flux<Review> getReviews(@RequestParam(value = "productId", required = true) int productId);
}
This time the return type is Flux<Review> because a product has many reviews. productId is passed as a query parameter (/review?productId=1) rather than in the path, because we are filtering a set of reviews, not fetching one specific review.
Recommendation and RecommendationService are similar, with a rate field instead of subject:
public record Recommendation(
int productId,
int recommendationId,
String author,
int rate,
String content,
String serviceAddress) {
}
Simulated implementation
@RestController
public class ReviewServiceImpl implements ReviewService {
private static final Logger LOG = LoggerFactory.getLogger(ReviewServiceImpl.class);
private final ServiceUtil serviceUtil;
public ReviewServiceImpl(ServiceUtil serviceUtil) {
this.serviceUtil = serviceUtil;
}
@Override
public Flux<Review> getReviews(int productId) {
if (productId < 1) {
throw new InvalidInputException("Invalid productId: " + productId);
}
// Simulated: product 213 has no reviews yet.
if (productId == 213) {
LOG.debug("No reviews found for productId: {}", productId);
return Flux.empty();
}
String address = serviceUtil.getServiceAddress();
List<Review> list = List.of(
new Review(productId, 1, "Author 1", "Subject 1", "Content 1", address),
new Review(productId, 2, "Author 2", "Subject 2", "Content 2", address),
new Review(productId, 3, "Author 3", "Subject 3", "Content 3", address));
LOG.debug("/reviews response size: {}", list.size());
return Flux.fromIterable(list);
}
}
RecommendationServiceImpl is identical, except that product 113 has no recommendations. Add these two rules to the list from the previous post:
| productId | Meaning |
|---|---|
|
Product does not exist (404) |
|
No recommendations |
|
No reviews |
negative |
Invalid input (422) |
Each service has a build.gradle, a main class (still with @ComponentScan("com.ecommerce")) and an application.yml like product-service, with server.port set to 7003 for review and 7002 for recommendation. The tests follow the WebTestClient pattern from part 2, plus a missing-parameter case:
@Test
void getReviewsMissingParameter() {
client.get()
.uri("/review")
.accept(APPLICATION_JSON)
.exchange()
.expectStatus().isBadRequest()
.expectBody()
.jsonPath("$.path").isEqualTo("/review");
}
Remember to add the three new modules to settings.gradle:
include ':microservices:review-service'
include ':microservices:recommendation-service'
include ':microservices:product-composite-service'
The composite contract
The composite returns an aggregate. I use separate summary records because the product page does not need productId repeated inside every review:
package com.ecommerce.api.composite.product;
public record ProductAggregate(
int productId,
String name,
BigDecimal price,
List<RecommendationSummary> recommendations,
List<ReviewSummary> reviews,
ServiceAddresses serviceAddresses) {
}
public record ReviewSummary(int reviewId, String author, String subject, String content) {
}
public record RecommendationSummary(int recommendationId, String author, int rate, String content) {
}
/** Address of the instance that handled each part: composite, product, review, recommendation. */
public record ServiceAddresses(String cmp, String pro, String rev, String rec) {
}
Each record lives in its own file in the com.ecommerce.api.composite.product package. The interface:
public interface ProductCompositeService {
@GetMapping(value = "/product-composite/{productId}", produces = "application/json")
Mono<ProductAggregate> getProduct(@PathVariable int productId);
}
The integration layer: talking to the other services
The book puts the HTTP calls in a separate class, ProductCompositeIntegration. The nice part: this class implements the very interfaces from api. To the rest of the composite, calling integration.getProduct(1) looks exactly like calling a local service; the network is hidden in here.
@Component
public class ProductCompositeIntegration implements ProductService, RecommendationService, ReviewService {
private static final Logger LOG = LoggerFactory.getLogger(ProductCompositeIntegration.class);
private final WebClient webClient;
private final ObjectMapper mapper;
private final String productServiceUrl;
private final String recommendationServiceUrl;
private final String reviewServiceUrl;
public ProductCompositeIntegration(
WebClient.Builder webClientBuilder,
ObjectMapper mapper,
@Value("${app.product-service.host}") String productServiceHost,
@Value("${app.product-service.port}") int productServicePort,
@Value("${app.recommendation-service.host}") String recommendationServiceHost,
@Value("${app.recommendation-service.port}") int recommendationServicePort,
@Value("${app.review-service.host}") String reviewServiceHost,
@Value("${app.review-service.port}") int reviewServicePort) {
this.webClient = webClientBuilder.build();
this.mapper = mapper;
productServiceUrl = "http://" + productServiceHost + ":" + productServicePort + "/product/";
recommendationServiceUrl = "http://" + recommendationServiceHost + ":" + recommendationServicePort
+ "/recommendation?productId=";
reviewServiceUrl = "http://" + reviewServiceHost + ":" + reviewServicePort + "/review?productId=";
}
@Override
public Mono<Product> getProduct(int productId) {
String url = productServiceUrl + productId;
LOG.debug("Will call getProduct API on URL: {}", url);
return webClient.get().uri(url).retrieve()
.bodyToMono(Product.class)
.onErrorMap(WebClientResponseException.class, this::handleException); // (1)
}
@Override
public Flux<Review> getReviews(int productId) {
String url = reviewServiceUrl + productId;
LOG.debug("Will call getReviews API on URL: {}", url);
return webClient.get().uri(url).retrieve()
.bodyToFlux(Review.class)
.onErrorResume(ex -> { // (2)
LOG.warn("Got an exception while requesting reviews, return zero reviews: {}", ex.getMessage());
return Flux.empty();
});
}
// getRecommendations(...) is identical to getReviews(...)
private Throwable handleException(WebClientResponseException ex) {
switch (ex.getStatusCode().value()) {
case 404:
return new NotFoundException(getErrorMessage(ex));
case 422:
return new InvalidInputException(getErrorMessage(ex));
default:
LOG.warn("Got an unexpected HTTP error: {}, will rethrow it", ex.getStatusCode());
LOG.warn("Error body: {}", ex.getResponseBodyAsString());
return ex;
}
}
private String getErrorMessage(WebClientResponseException ex) {
try {
return mapper.readValue(ex.getResponseBodyAsString(), HttpErrorInfo.class).message(); // (3)
} catch (IOException ioex) {
return ex.getMessage();
}
}
}
-
When product-service answers 404, WebClient throws a
WebClientResponseException. We turn it back into aNotFoundException, so the shared error handler in the composite returns 404 to the client. Without this step, every error from a downstream service becomes a 500. -
Reviews and recommendations are the secondary parts of the product page. An error here returns an empty list instead of breaking the whole page. This is a crude form of fallback; the Resilience4j post does it properly.
-
Read the
HttpErrorInfoback from the error body to keep the original message, for exampleNo product found for productId: 13.
|
Note
|
Different from the book: in chapter 3 the book uses |
Merging the results: Mono.zip
@RestController
public class ProductCompositeServiceImpl implements ProductCompositeService {
private final ServiceUtil serviceUtil;
private final ProductCompositeIntegration integration;
public ProductCompositeServiceImpl(ServiceUtil serviceUtil, ProductCompositeIntegration integration) {
this.serviceUtil = serviceUtil;
this.integration = integration;
}
@Override
public Mono<ProductAggregate> getProduct(int productId) {
// Call all 3 services in parallel, wait for every result, then merge.
return Mono.zip(
integration.getProduct(productId),
integration.getRecommendations(productId).collectList(),
integration.getReviews(productId).collectList())
.map(t -> createProductAggregate(t.getT1(), t.getT2(), t.getT3(), serviceUtil.getServiceAddress()));
}
private ProductAggregate createProductAggregate(Product product, List<Recommendation> recommendations,
List<Review> reviews, String serviceAddress) {
List<RecommendationSummary> recommendationSummaries = recommendations.stream()
.map(r -> new RecommendationSummary(r.recommendationId(), r.author(), r.rate(), r.content()))
.toList();
List<ReviewSummary> reviewSummaries = reviews.stream()
.map(r -> new ReviewSummary(r.reviewId(), r.author(), r.subject(), r.content()))
.toList();
String productAddress = product.serviceAddress();
String reviewAddress = reviews.isEmpty() ? "" : reviews.get(0).serviceAddress();
String recommendationAddress = recommendations.isEmpty() ? "" : recommendations.get(0).serviceAddress();
ServiceAddresses serviceAddresses =
new ServiceAddresses(serviceAddress, productAddress, reviewAddress, recommendationAddress);
return new ProductAggregate(product.productId(), product.name(), product.price(),
recommendationSummaries, reviewSummaries, serviceAddresses);
}
}
Mono.zip subscribes to all three sources at once, so the three HTTP requests go out in parallel. The page’s response time is that of the slowest service, not the sum of all three. If getProduct fails, the whole zip fails, which is what we want: no product, no page.
Configuring the service addresses
server.port: 7000
server.error.include-message: always
spring.application.name: product-composite
app:
product-service:
host: localhost
port: 7001
recommendation-service:
host: localhost
port: 7002
review-service:
host: localhost
port: 7003
logging:
level:
root: INFO
com.ecommerce: DEBUG
Hosts and ports are hard-coded for now. Post 4 overrides them for Docker, and the Eureka post removes them entirely in favor of service names.
Running all four services
./gradlew build
java -jar microservices/product-service/build/libs/*.jar &
java -jar microservices/review-service/build/libs/*.jar &
java -jar microservices/recommendation-service/build/libs/*.jar &
java -jar microservices/product-composite-service/build/libs/*.jar &
curl -s localhost:7000/product-composite/1 | jq .
{
"productId": 1,
"name": "name-1",
"price": 199000,
"recommendations": [
{ "recommendationId": 1, "author": "Author 1", "rate": 1, "content": "Content 1" },
{ "recommendationId": 2, "author": "Author 2", "rate": 2, "content": "Content 2" },
{ "recommendationId": 3, "author": "Author 3", "rate": 3, "content": "Content 3" }
],
"reviews": [
{ "reviewId": 1, "author": "Author 1", "subject": "Subject 1", "content": "Content 1" },
{ "reviewId": 2, "author": "Author 2", "subject": "Subject 2", "content": "Content 2" },
{ "reviewId": 3, "author": "Author 3", "subject": "Subject 3", "content": "Content 3" }
],
"serviceAddresses": {
"cmp": "vm/127.0.0.1:7000",
"pro": "vm/127.0.0.1:7001",
"rev": "vm/127.0.0.1:7003",
"rec": "vm/127.0.0.1:7002"
}
}
serviceAddresses shows that four different processes handled the four parts. Errors from product-service pass through the composite with the right status and message:
$ curl -s localhost:7000/product-composite/13
{"timestamp":"2026-10-02T05:57:56.405471652Z","path":"/product-composite/13","httpStatus":"NOT_FOUND","message":"No product found for productId: 13"}
Shutting down review-service
This is my favorite experiment in this post. Stop review-service (kill its process) and call again:
$ curl -s localhost:7000/product-composite/1 | jq -c '{reviews: (.reviews|length), recommendations: (.recommendations|length)}'
{"reviews":0,"recommendations":3}
The product page still returns 200, only without reviews, and the composite logs a WARN … return zero reviews line. This is the "failures do not spread" benefit from part 1, even in its simplest form.
Stop everything when you are done:
kill $(jobs -p)
Testing the composite without the other three services
The composite’s tests should not depend on whether the other three services are running. We replace ProductCompositeIntegration with a mock:
@SpringBootTest(webEnvironment = RANDOM_PORT)
class ProductCompositeServiceApplicationTests {
private static final int PRODUCT_ID_OK = 1;
private static final int PRODUCT_ID_NOT_FOUND = 2;
private static final int PRODUCT_ID_INVALID = 3;
@Autowired
private WebTestClient client;
@MockitoBean
private ProductCompositeIntegration compositeIntegration;
@BeforeEach
void setUp() {
when(compositeIntegration.getProduct(PRODUCT_ID_OK))
.thenReturn(Mono.just(new Product(PRODUCT_ID_OK, "name", new BigDecimal("199000"), "mock-address")));
when(compositeIntegration.getRecommendations(PRODUCT_ID_OK))
.thenReturn(Flux.just(new Recommendation(PRODUCT_ID_OK, 1, "author", 1, "content", "mock address")));
when(compositeIntegration.getReviews(PRODUCT_ID_OK))
.thenReturn(Flux.just(new Review(PRODUCT_ID_OK, 1, "author", "subject", "content", "mock address")));
when(compositeIntegration.getProduct(PRODUCT_ID_NOT_FOUND))
.thenReturn(Mono.error(new NotFoundException("NOT FOUND: " + PRODUCT_ID_NOT_FOUND)));
when(compositeIntegration.getRecommendations(PRODUCT_ID_NOT_FOUND)).thenReturn(Flux.empty());
when(compositeIntegration.getReviews(PRODUCT_ID_NOT_FOUND)).thenReturn(Flux.empty());
// PRODUCT_ID_INVALID returns InvalidInputException the same way
}
@Test
void getProductById() {
client.get()
.uri("/product-composite/" + PRODUCT_ID_OK)
.accept(APPLICATION_JSON)
.exchange()
.expectStatus().isOk()
.expectBody()
.jsonPath("$.productId").isEqualTo(PRODUCT_ID_OK)
.jsonPath("$.recommendations.length()").isEqualTo(1)
.jsonPath("$.reviews.length()").isEqualTo(1);
}
@Test
void getProductNotFound() {
client.get()
.uri("/product-composite/" + PRODUCT_ID_NOT_FOUND)
.accept(APPLICATION_JSON)
.exchange()
.expectStatus().isEqualTo(NOT_FOUND)
.expectBody()
.jsonPath("$.path").isEqualTo("/product-composite/" + PRODUCT_ID_NOT_FOUND)
.jsonPath("$.message").isEqualTo("NOT FOUND: " + PRODUCT_ID_NOT_FOUND);
}
}
|
Note
|
Different from the book: the book uses |
./gradlew build
All tests of the four services must pass. Commit:
git add .
git commit -m "Post 3: review, recommendation and product-composite"
git tag blog-03
Summary
-
The composite calls three services in parallel with
WebClientandMono.zip, so the page is only as slow as the slowest service. -
The integration layer hides the network behind the interfaces from
api, and translates HTTP errors back into our exceptions. -
Errors in secondary parts (reviews, recommendations) do not break the product page.
-
The composite is tested with
@MockitoBean, without the other three services running.
Four Java processes started by hand with four commands does not scale. In part 4, I package each service as a Docker image and start the whole system with one docker compose up.