|
Ghi chú
|
Đây là bài 8 trong series microservices e-commerce, phần 2 (dữ liệu và sự kiện). Code của bài ở tag |
Ở bài 6, composite tạo một sản phẩm bằng cách gọi REST song song tới product, recommendation và review. Nếu review-service đang tắt đúng lúc đó, sản phẩm được tạo còn đánh giá thì mất, và client nhận lỗi 500 dù một nửa công việc đã xong.
Bài này đổi cách ghi dữ liệu: composite không gọi REST nữa mà gửi event vào Kafka, rồi trả lời ngay. Mỗi service lõi tự đọc event của mình và xử lý khi sẵn sàng. Còn đọc thì vẫn như cũ, qua REST.
Cuối bài bạn sẽ có:
-
Kafka chạy trong Docker Compose, chế độ KRaft, không cần ZooKeeper.
-
POSTvàDELETEtrên composite trả về202 Acceptedvà gửi event vào bốn topic. -
product, recommendation, review và inventory tiêu thụ event qua Spring Cloud Stream.
-
Một thí nghiệm: tắt review-service, tạo sản phẩm, bật lại, và đánh giá vẫn xuất hiện.
Vì sao ghi bằng event
Với REST, composite phải chờ cả bốn service trả lời, và chỉ một service lỗi là cả request lỗi. Với event:
-
Composite chỉ cần Kafka nhận event là xong. Nó không cần biết service nào đang chạy.
-
Kafka lưu event lại. Service nào đang tắt thì khi bật lên sẽ đọc tiếp từ chỗ nó dừng.
-
Thêm một service mới cũng muốn biết khi có sản phẩm mới, chỉ cần cho nó đọc cùng topic, composite không phải sửa gì.
Cái giá phải trả là nhất quán sau cùng (eventual consistency): ngay sau khi nhận 202, dữ liệu có thể chưa có ở service lõi. Client phải chấp nhận điều đó, và test của ta cũng vậy.
Đọc thì khác. Trang sản phẩm cần dữ liệu ngay, nên GET vẫn gọi REST. Sách gọi cách chia này là ghi bất đồng bộ, đọc đồng bộ.
Spring Cloud Stream
Ta không dùng thẳng thư viện Kafka client. Spring Cloud Stream cho phép viết code gửi, nhận message mà không phụ thuộc broker cụ thể. Ba khái niệm cần nhớ:
-
Binder: phần nối với broker thật. Ở đây là Kafka binder.
-
Binding: một đầu vào hoặc đầu ra có tên, ví dụ
products-out-0haymessageProcessor-in-0. -
Destination: tên topic mà binding trỏ tới, khai trong cấu hình.
Code chỉ biết tên binding, còn binding trỏ tới topic nào, trên broker nào, là việc của file cấu hình.
Thêm dependency vào composite và bốn service lõi:
implementation 'org.springframework.cloud:spring-cloud-starter-stream-kafka'
testImplementation 'org.springframework.cloud:spring-cloud-stream-test-binder'
Phiên bản lấy từ Spring Cloud BOM 2025.0.3 đã khai trong build.gradle gốc từ bài 1. Test binder là một binder giả trong bộ nhớ, dùng cho test để không phải chạy Kafka.
Event
Mọi event có cùng một dạng, đặt trong module api để hai bên gửi và nhận dùng chung:
/**
* Event dùng chung để tạo và xoá dữ liệu ở các core service.
* key là productId, data là đối tượng cần tạo (null với DELETE).
*/
public record Event<K, T>(
Type eventType,
K key,
T data,
@JsonFormat(shape = JsonFormat.Shape.STRING) ZonedDateTime eventCreatedAt) {
public enum Type {
CREATE,
DELETE
}
public Event(Type eventType, K key, T data) {
this(eventType, key, data, ZonedDateTime.now());
}
}
@JsonFormat để thời điểm tạo event được ghi thành chuỗi ISO-8601 dễ đọc, thay vì một con số.
Composite gửi event
ProductCompositeIntegration giữ nguyên các method GET dùng WebClient. Các method tạo, xoá thì giờ gửi event qua StreamBridge:
@Override
public Mono<Product> createProduct(Product body) {
return Mono.fromCallable(() -> {
sendMessage("products-out-0", new Event<>(Event.Type.CREATE, body.productId(), body));
return body;
}).subscribeOn(publishEventScheduler);
}
@Override
public Mono<Void> deleteProduct(int productId) {
return Mono.fromRunnable(() -> sendMessage("products-out-0", new Event<>(Event.Type.DELETE, productId, null)))
.subscribeOn(publishEventScheduler).then();
}
/** Gửi event qua StreamBridge; binding name quyết định topic (xem application.yml). */
private void sendMessage(String bindingName, Event<Integer, ?> event) {
LOG.debug("Sending a {} message to {}", event.eventType(), bindingName);
streamBridge.send(bindingName, event);
}
Recommendation, review và inventory viết y hệt, chỉ khác tên binding: recommendations-out-0, reviews-out-0, inventory-out-0.
streamBridge.send() là lời gọi blocking: lần đầu gửi tới một topic, Kafka producer phải chờ lấy metadata từ broker. Vì vậy ta lặp lại đúng mẹo của bài 6, chạy nó trên một thread pool riêng:
/**
* StreamBridge.send() có thể bị chặn (chờ metadata từ Kafka), nên chạy nó trên thread pool riêng
* thay vì event loop của WebFlux, giống cách review-service chạy JPA ở bài 6.
*/
@Bean
public Scheduler publishEventScheduler(
@Value("${app.threadPoolSize:10}") Integer threadPoolSize,
@Value("${app.taskQueueSize:100}") Integer taskQueueSize) {
return Schedulers.newBoundedElastic(threadPoolSize, taskQueueSize, "publish-pool");
}
Tên binding nối với topic trong application.yml của composite:
spring.cloud.stream:
defaultBinder: kafka
default.contentType: application/json
bindings:
products-out-0.destination: products
recommendations-out-0.destination: recommendations
reviews-out-0.destination: reviews
inventory-out-0.destination: inventory
kafka.binder.brokers: localhost:9092
Profile docker đổi broker thành kafka:9092.
ProductCompositeServiceImpl gần như không đổi: nó vẫn gọi integration.createProduct(…) và các method khác, rồi chờ bằng Mono.when. Chỉ khác là giờ "xong" nghĩa là event đã được gửi, không phải dữ liệu đã được lưu. Interface trong module api báo điều đó cho client:
@ResponseStatus(HttpStatus.ACCEPTED)
@PostMapping(value = "/product-composite", consumes = "application/json")
Mono<Void> createProduct(@RequestBody ProductAggregate body);
DELETE /product-composite/{productId} cũng trả 202 Accepted. Mã 202 nghĩa là "đã nhận yêu cầu, sẽ xử lý sau", khác với 200 "đã xong".
Một hệ quả: lỗi xảy ra sau khi gửi event, ví dụ trùng productId, không còn trả về client được nữa. Những lỗi kiểm tra được ngay thì nên kiểm tra ở composite, trước khi gửi:
if (body.productId() < 1) {
throw new InvalidInputException("Invalid productId: " + body.productId());
}
Service lõi nhận event
Mỗi service lõi khai một bean kiểu java.util.function.Consumer. Spring Cloud Stream thấy bean này và tự gắn nó vào một binding đầu vào:
@Configuration
public class MessageProcessorConfig {
private static final Logger LOG = LoggerFactory.getLogger(MessageProcessorConfig.class);
private final ProductService productService;
public MessageProcessorConfig(ProductService productService) {
this.productService = productService;
}
/** Spring Cloud Stream gắn bean này vào binding messageProcessor-in-0, tức topic products. */
@Bean
public Consumer<Event<Integer, Product>> messageProcessor() {
return event -> {
LOG.info("Process message created at {}...", event.eventCreatedAt());
switch (event.eventType()) {
case CREATE -> {
LOG.info("Create product with ID: {}", event.key());
productService.createProduct(event.data()).block();
}
case DELETE -> {
LOG.info("Delete products with ProductID: {}", event.key());
productService.deleteProduct(event.key()).block();
}
default -> throw new EventProcessingException(
"Incorrect event type: " + event.eventType() + ", expected a CREATE or DELETE event");
}
LOG.info("Message processing done!");
};
}
}
Vài điểm:
-
Kiểu
Consumer<Event<Integer, Product>>đủ để Spring Cloud Stream chuyển JSON thành đúngEventchứaProduct, không cần viết code đọc JSON. -
Logic vẫn nằm trong
ProductServiceImplcủa bài 6. Message processor chỉ là một lối vào mới. -
.block()ở đây không sao: consumer chạy trên thread riêng của Kafka listener, không phải event loop của WebFlux. Chờ xử lý xong mới báo Kafka là đã đọc, nên lỗi khi lưu sẽ được Kafka biết.
Cấu hình đầu vào:
spring.cloud.function.definition: messageProcessor
spring.cloud.stream:
defaultBinder: kafka
default.contentType: application/json
bindings.messageProcessor-in-0:
destination: products
group: productsGroup
kafka.binder.brokers: localhost:9092
Tên binding messageProcessor-in-0 là quy ước: tên bean, in cho đầu vào, 0 là vị trí tham số. group quan trọng hơn vẻ ngoài của nó:
-
Mọi instance cùng một group chia nhau đọc topic, mỗi event chỉ được một instance xử lý. Khi scale product-service lên hai instance, sản phẩm không bị tạo hai lần.
-
Kafka nhớ vị trí đã đọc (offset) của từng group. Service tắt rồi bật lại sẽ đọc tiếp từ đó.
-
Group mới tạo bắt đầu đọc từ đầu topic. Không có
group, Spring Cloud Stream tạo một group tạm ngẫu nhiên, bắt đầu từ cuối topic, và event gửi lúc service đang tắt sẽ bị bỏ qua.
Recommendation, review và inventory giống hệt, với topic và group của riêng mình. Inventory gọi setStock cho event CREATE.
Bỏ REST tạo và xoá ở service lõi
Nếu service lõi vẫn mở POST /product, sẽ có hai đường ghi dữ liệu, và sớm muộn ai đó sẽ dùng nhầm đường. Nên mình gỡ annotation REST khỏi các method tạo, xoá trong module api, chỉ giữ method để message processor gọi:
public interface ProductService {
/** Không còn là REST endpoint: chỉ được gọi khi nhận event CREATE từ topic products. */
Mono<Product> createProduct(Product body);
@GetMapping(value = "/product/{productId}", produces = "application/json")
Mono<Product> getProduct(@PathVariable int productId);
/** Không còn là REST endpoint: chỉ được gọi khi nhận event DELETE từ topic products. */
Mono<Void> deleteProduct(int productId);
}
Sách cũng làm vậy ở chương 7.
Kafka trong Docker Compose
kafka:
image: apache/kafka:3.9.1
mem_limit: 1024m
environment:
# Chế độ KRaft: một node vừa là broker vừa là controller, không cần ZooKeeper
- KAFKA_NODE_ID=1
- KAFKA_PROCESS_ROLES=broker,controller
- KAFKA_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093
- KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092
- KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER
- KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
- KAFKA_CONTROLLER_QUORUM_VOTERS=1@kafka:9093
- KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1
- KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR=1
- KAFKA_TRANSACTION_STATE_LOG_MIN_ISR=1
- KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS=0
healthcheck:
test: ["CMD-SHELL", "/opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list"]
interval: 10s
timeout: 10s
retries: 30
-
KAFKA_ADVERTISED_LISTENERSlà địa chỉ broker báo cho client. Các service chạy trong cùng mạng compose nên dùng tênkafka. -
Ba biến
…_REPLICATION_FACTORvàMIN_ISRbằng 1 vì ta chỉ có một broker. Mặc định là 3, và một broker thì không đủ để tạo các topic nội bộ. -
GROUP_INITIAL_REBALANCE_DELAY_MS=0để consumer group nhận việc ngay, không chờ 3 giây như mặc định.
Composite và bốn service lõi thêm depends_on để chỉ khởi động khi Kafka đã sẵn sàng:
depends_on:
kafka:
condition: service_healthy
|
Ghi chú
|
Chỗ khác sách: sách dùng RabbitMQ làm mặc định, Kafka chỉ là lựa chọn thứ hai, chạy bằng image Chương 7 của sách còn bật partition, retry và dead-letter queue ngay từ đầu. Mình để ba thứ này sang bài 9, vì mỗi thứ cần một ví dụ lỗi riêng mới thấy được giá trị. |
Test
Composite: event có được gửi đúng không
Test binder ghi lại mọi message gửi đi, và OutputDestination cho phép đọc chúng theo tên topic:
@SpringBootTest(webEnvironment = RANDOM_PORT)
@Import(TestChannelBinderConfiguration.class)
class MessagingTests {
@Autowired
private WebTestClient client;
@Autowired
private OutputDestination target;
@Test
void createCompositeProduct() throws Exception {
ProductAggregate composite = new ProductAggregate(1, "name", new BigDecimal("199000"), 10,
List.of(new RecommendationSummary(1, "a", 1, "c")),
List.of(new ReviewSummary(1, "a", "s", "c")), null);
client.post().uri("/product-composite").bodyValue(composite).exchange().expectStatus().isAccepted();
List<Map<String, Object>> products = getMessages("products");
assertEquals(1, products.size());
assertEquals("CREATE", products.get(0).get("eventType"));
assertEquals(1, products.get(0).get("key"));
// ...
List<Map<String, Object>> inventory = getMessages("inventory");
assertEquals(1, inventory.size());
assertEquals(10, ((Map<?, ?>) inventory.get(0).get("data")).get("quantity"));
}
// deleteCompositeProduct, createCompositeProductWithoutStock, createWithInvalidProductIdSendsNothing ...
}
getMessages(topic) gọi target.receive(0, topic) tới khi hết message, rồi đọc JSON thành Map. Sách viết một Hamcrest matcher IsSameEvent để so sánh event mà bỏ qua eventCreatedAt. Mình so từng field cần kiểm tra, ngắn hơn và đủ dùng.
Service lõi: gọi thẳng message processor
Test của service lõi không cần Kafka, cũng không cần test binder gửi message. Nó lấy chính bean messageProcessor và gọi accept():
@Testcontainers
@SpringBootTest(webEnvironment = RANDOM_PORT)
@Import(TestChannelBinderConfiguration.class)
class ProductServiceApplicationTests {
@Container
@ServiceConnection
static MongoDBContainer database = new MongoDBContainer("mongo:7.0");
// Bean xử lý event thật; test gọi thẳng vào nó thay vì gửi qua Kafka
@Autowired
@Qualifier("messageProcessor")
private Consumer<Event<Integer, Product>> messageProcessor;
@Test
void duplicateError() {
int productId = 1;
sendCreateProductEvent(productId);
InvalidInputException thrown = assertThrows(InvalidInputException.class,
() -> sendCreateProductEvent(productId));
assertEquals("Duplicate key, Product Id: " + productId, thrown.getMessage());
}
private void sendCreateProductEvent(int productId) {
Product product = new Product(productId, "Name " + productId, new BigDecimal("199000"), null);
messageProcessor.accept(new Event<>(Event.Type.CREATE, productId, product));
}
}
@Import(TestChannelBinderConfiguration.class) thay Kafka binder bằng test binder, nên context khởi động được mà không cần broker. Các test GET qua REST giữ nguyên.
./gradlew build
47 test, tất cả đều qua.
Test end-to-end với dữ liệu bất đồng bộ
test-em-all.bash phải đổi ba chỗ.
Tạo và xoá giờ trả 202:
assertCurl 202 "curl -X DELETE http://$HOST:$PORT/product-composite/${productId} -s"
assertEqual 202 $(curl -X POST -s http://$HOST:$PORT/product-composite -H "Content-Type: application/json" \
--data "$composite" -w "%{http_code}")
Sau khi tạo dữ liệu test, phải chờ các service lõi xử lý xong event rồi mới kiểm tra:
# Ghi giờ là bất đồng bộ: chờ tới khi các service lõi đã xử lý xong event của dữ liệu test
function waitForMessageProcessing() {
echo "Wait for messages to be processed... "
local n=0
until [ "$(curl -s http://$HOST:$PORT/product-composite/$PROD_ID_REVS_RECS \
| jq -c '[(.recommendations | length), (.reviews | length), .stock]' 2>/dev/null)" = "[3,3,10]" ] \
&& [ "$(curl -s -o /dev/null -w "%{http_code}" http://$HOST:$PORT/product-composite/$PROD_ID_NO_RECS)" = "200" ] \
&& [ "$(curl -s -o /dev/null -w "%{http_code}" http://$HOST:$PORT/product-composite/$PROD_ID_NO_REVS)" = "200" ]; do
n=$((n + 1))
if [[ $n == 40 ]]; then
echo " Give up"
exit 1
fi
sleep 1
echo -n ", retry #$n "
done
echo "All messages are now processed!"
}
Lần chạy đầu của hàm này thất bại sau 40 giây dù dữ liệu đã có đủ. Lỗi nằm ở biểu thức jq: viết [.recommendations | length, (.reviews | length), .stock] thì dấu | đầu tiên áp lên cả phần còn lại, và jq cố đọc .reviews bên trong mảng recommendations. Thêm ngoặc quanh (.recommendations | length) là hết. Chờ bất đồng bộ dễ che lỗi kiểu này, vì triệu chứng chỉ là "chờ mãi không xong".
Cuối cùng, kiểm tra xoá cũng phải chờ:
# Xoá là idempotent: gửi hai lần đều được nhận (202), sau khi event được xử lý thì sản phẩm không còn
assertCurl 202 "curl -X DELETE http://$HOST:$PORT/product-composite/$PROD_ID_NO_REVS -s"
assertCurl 202 "curl -X DELETE http://$HOST:$PORT/product-composite/$PROD_ID_NO_REVS -s"
waitForHttpCode 404 http://$HOST:$PORT/product-composite/$PROD_ID_NO_REVS
assertCurl 404 "curl http://$HOST:$PORT/product-composite/$PROD_ID_NO_REVS -s"
Ban đầu script chờ hệ thống khởi động bằng một lệnh DELETE. Giờ DELETE trả 202 ngay khi composite gửi được event, nên nó không còn cho biết service lõi đã chạy chưa. Thay vào đó, script chờ GET /product-composite/13 trả 404: nghĩa là composite đã gọi được product-service và product-service trả lời đúng.
$ ./gradlew build && docker compose build
$ ./test-em-all.bash start
...
Wait for: http://localhost:8080/product-composite/13... , retry #1 , ... , retry #7 DONE, continues...
...
Wait for messages to be processed...
, retry #1 All messages are now processed!
...
End, all tests OK: Thu Oct 8 03:36:35 UTC 2026
Nhìn vào Kafka
Tạo một sản phẩm mới:
$ curl -s -i -X POST localhost:8080/product-composite -H "Content-Type: application/json" \
--data '{"productId":2,"name":"Balo laptop","price":650000,"stock":5,
"reviews":[{"reviewId":1,"author":"an","subject":"Ben","content":"Dung tot"}]}' | head -1
HTTP/1.1 202 Accepted
Kafka đã tự tạo bốn topic khi composite gửi tới lần đầu:
$ docker compose exec -T kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
__consumer_offsets
inventory
products
recommendations
reviews
Đọc topic products từ đầu, ba event cuối là:
$ docker compose exec -T kafka /opt/kafka/bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic products --from-beginning --timeout-ms 5000
...
{"eventType":"DELETE","key":213,"data":null,"eventCreatedAt":"2026-10-08T03:36:34.139536497Z"}
{"eventType":"DELETE","key":213,"data":null,"eventCreatedAt":"2026-10-08T03:36:34.17359037Z"}
{"eventType":"CREATE","key":2,"data":{"productId":2,"name":"Balo laptop","price":650000,"serviceAddress":null},"eventCreatedAt":"2026-10-08T03:36:46.254659284Z"}
Hai event DELETE sản phẩm 213 là hai lần xoá trong kiểm tra idempotent ở trên. Event không mất đi sau khi được đọc: topic là một log chỉ ghi thêm, và mỗi consumer group tự nhớ mình đã đọc tới đâu.
$ docker compose exec -T kafka /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group productsGroup
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG ...
productsGroup products 0 9 9 0 ...
LAG bằng 0: product-service đã đọc hết mọi event trong topic.
Log của product-service khi nhận event tạo sản phẩm 2:
INFO 1 --- [product] [container-0-C-1] c.e.c.p.services.MessageProcessorConfig : Process message created at 2026-10-08T03:36:46.254659284Z...
INFO 1 --- [product] [container-0-C-1] c.e.c.p.services.MessageProcessorConfig : Create product with ID: 2
INFO 1 --- [product] [container-0-C-1] c.e.c.p.services.MessageProcessorConfig : Message processing done!
container-0-C-1 là thread của Kafka listener, không phải event loop, nên .block() ở đây an toàn như đã nói.
Thí nghiệm: tắt review-service
Đây là lý do chính để dùng event. Tắt review-service, rồi tạo một sản phẩm có đánh giá:
$ docker compose stop review
$ curl -s -o /dev/null -w "POST %{http_code}\n" -X POST localhost:8080/product-composite \
-H "Content-Type: application/json" \
--data '{"productId":3,"name":"Mu len","price":150000,
"reviews":[{"reviewId":1,"author":"binh","subject":"Am","content":"Doi mua dong rat am"}]}'
POST 202
$ curl -s localhost:8080/product-composite/3 | jq -c '{productId,name,reviews: (.reviews|length)}'
{"productId":3,"name":"Mu len","reviews":0}
Ở bài 6, request này sẽ lỗi. Giờ nó được nhận, sản phẩm đã có, chỉ thiếu đánh giá. Event đánh giá đang nằm trong Kafka chờ:
$ docker compose exec -T kafka /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group reviewsGroup
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID
reviewsGroup reviews 0 12 13 1 -
LAG là 1, và không có consumer nào (-). Bật review-service lại:
$ docker compose start review
$ curl -s localhost:8080/product-composite/3 | jq -c '{productId,name,reviews}'
{"productId":3,"name":"Mu len","reviews":[{"reviewId":1,"author":"binh","subject":"Am","content":"Doi mua dong rat am"}]}
Log của review-service cho thấy event được tạo lúc 03:37:12 và được xử lý lúc 03:37:40, sau khi service khởi động xong:
INFO 1 --- [review] [container-0-C-1] c.e.c.r.services.MessageProcessorConfig : Process message created at 2026-10-08T03:37:12.421690509Z...
INFO 1 --- [review] [container-0-C-1] c.e.c.r.services.MessageProcessorConfig : Create review with ID: 3
INFO 1 --- [review] [container-0-C-1] c.e.c.r.services.MessageProcessorConfig : Message processing done!
Không ai phải gọi lại, không dữ liệu nào mất.
Còn một lỗ hổng: event lỗi
Event nào mà service lõi không xử lý được thì sao? Thử tạo lại sản phẩm 1, vốn đã tồn tại:
$ curl -s -o /dev/null -w "POST %{http_code}\n" -X POST localhost:8080/product-composite \
-H "Content-Type: application/json" --data '{"productId":1,"name":"Giay sneaker","price":899000}'
POST 202
Composite trả 202 vì nó không biết sản phẩm đã có. Log của product-service (rút gọn):
ERROR 1 --- [product] [container-0-C-1] o.s.integration.handler.LoggingHandler : org.springframework.messaging.MessageHandlingException: error occurred in message handler [...], failedMessage=GenericMessage [payload=byte[162], headers={kafka_offset=10, ..., deliveryAttempt=3, ..., kafka_receivedTopic=products, ..., kafka_groupId=productsGroup}]
Caused by: com.ecommerce.api.exceptions.InvalidInputException: Duplicate key, Product Id: 1
ERROR 1 --- [product] [container-0-C-1] o.s.kafka.listener.DefaultErrorHandler : Backoff none exhausted for products-0@10
Spring Cloud Stream đã thử 3 lần (deliveryAttempt=3, mặc định), rồi ghi log và bỏ qua event. LAG của productsGroup lại về 0. Với một event trùng lặp thì bỏ qua cũng được, nhưng nếu lỗi là do database tạm thời không truy cập được, event đó mất luôn, và chỉ còn dấu vết trong log.
Bài 9 sẽ xử lý chuyện này: cấu hình retry có backoff, chuyển event lỗi vào dead-letter queue để xem lại sau, và chia topic thành nhiều partition để chạy nhiều instance mà vẫn giữ đúng thứ tự event của từng sản phẩm.
Commit:
git add .
git commit -m "Bài 8: event-driven với Kafka"
git tag blog-08
Tóm lại
-
Composite ghi bằng cách gửi event qua
StreamBridge, trả202 Accepted. Đọc vẫn qua REST. -
Service lõi nhận event bằng một bean
Consumer<Event<…>>. REST tạo, xoá ở service lõi được gỡ bỏ để chỉ còn một đường ghi. -
Consumer group giúp Kafka nhớ vị trí đã đọc; service tắt rồi bật lại vẫn xử lý đủ event.
-
Test dùng test binder thay cho Kafka; test end-to-end phải chờ dữ liệu xuất hiện thay vì kiểm tra ngay.
-
Event lỗi hiện vẫn bị bỏ qua sau 3 lần thử. Đó là việc của bài 9.