thắc mắc Retry trong spring boot kafka

  • Người tạo chủ đề Người tạo chủ đề hungdo1208
  • Ngày bắt đầu Ngày bắt đầu

hungdo1208

Senior Member
Xin chào mọi người, em đang tìm hiểu về kafka và có 1 số thắc mắc mong mọi người giải thích ạ.

Em đang sử dụng RetryTopicConfiguration để cấu hình retry + DLT cho consumer khi các retryable exception như SocketTimeoutException, TimeoutException xảy ra. Đầu tiên khi cố tình throw ra exception bên trong hàm consume thì retry hoạt động tốt. Sau đó e test đến lỗi disconnect bằng cách debug, trước khi commit offset thủ công bằng ack.acknowledge() thì tắt mạng để tạo lỗi connect tới kafka thì không thấy retry hoạt động mà chỉ thấy log disconnect ở console.

Em muốn hỏi mọi người 1 số thắc mắc như sau ạ:
  1. Retry có thể bắt và sử dụng cho các lỗi ở tầng kafka như connect broker, commit offset hay nó chỉ áp dụng cho các lỗi được throw ra bên trong hàm consumer (call api thứ 3 bị timeout, disconnect db khi save,...)

  2. Sau khi save db mà commit offset bị lỗi do disconnect khiến cho db đã lưu data nhưng message chưa được commit làm xảy ra nguy cơ message bị xử lý trùng lặp. trường hợp này hướng xử lý là như thế nào ạ (em có dùng @transactional nhưng không rollback được vì k có exception được throw ra, chỉ có log disconnect)
Mong mọi người giải đáp giúp em ạ.
 
1. cái này chắc phụ thuộc nhiều vào implementation của framework, mình cũng ko rõ lắm, nhưng miễn là chưa commit offset thì kiểu gì message đó cũng sẽ đc xử lý lại
2. nói rộng ra thì đây là bài toán dual write, có nhiều cách để xử lý tùy theo tình huống. Còn trong trường hợp này với kafka và database thì mình nghĩ có thể áp dụng vài cách sau:
  • xử lý message idempotent: nôm na là khi message đến thì bạn check xem đã xử lý cái message đó chưa và data đã đc lưu chưa, đã được thay đổi chưa. Cách này phụ thuộc vào nhiều business rule và data bạn lưu trong database.
  • mỗi khi nhận đc message thì sẽ lưu 1 job trong database, commit offset luôn rồi sau đó service lấy job đó ra xử lý: cách này cũng na ná như cái cách trên, xử lý ko đc lưu duplicate job, tuy nhiên thì ko phải phụ thuộc vào business rule, còn cái phần xử lý job khi có nhiều instances thì cần setup cũng có library hỗ trợ rồi
  • commit offset trước rồi mới lưu database sau: cách này thì lưu database lỗi thì ko rollback đc offset ở kafka. Chắc ít trường hợp phù hợp với cách này lắm
 
Sửa lần cuối:
Xin chào mọi người, em đang tìm hiểu về kafka và có 1 số thắc mắc mong mọi người giải thích ạ.

Em đang sử dụng RetryTopicConfiguration để cấu hình retry + DLT cho consumer khi các retryable exception như SocketTimeoutException, TimeoutException xảy ra. Đầu tiên khi cố tình throw ra exception bên trong hàm consume thì retry hoạt động tốt. Sau đó e test đến lỗi disconnect bằng cách debug, trước khi commit offset thủ công bằng ack.acknowledge() thì tắt mạng để tạo lỗi connect tới kafka thì không thấy retry hoạt động mà chỉ thấy log disconnect ở console.

Em muốn hỏi mọi người 1 số thắc mắc như sau ạ:
  1. Retry có thể bắt và sử dụng cho các lỗi ở tầng kafka như connect broker, commit offset hay nó chỉ áp dụng cho các lỗi được throw ra bên trong hàm consumer (call api thứ 3 bị timeout, disconnect db khi save,...)

  2. Sau khi save db mà commit offset bị lỗi do disconnect khiến cho db đã lưu data nhưng message chưa được commit làm xảy ra nguy cơ message bị xử lý trùng lặp. trường hợp này hướng xử lý là như thế nào ạ (em có dùng @transactional nhưng không rollback được vì k có exception được throw ra, chỉ có log disconnect)
Mong mọi người giải đáp giúp em ạ.
nghe thì có vẻ như liên quan đến exactly once processing (EOP), nhưng EOP hầu như là bất khả thể trong hệ thống phân tán

mình nghĩ là chỉ còn nghĩ cách làm sao để không gây ra vấn đề khi re-process thôi, như code logic hệ thống sao cho nó idempotent + đảm bảo message luôn được deliver ít nhất một lần (at least once processing, ALOP) chẳng hạn
 

Thống kê chủ đề

Ngày tạo
hungdo1208,
Người trả lời cuối
small-lambda,
Trả lời
2
Lượt xem
785
Quay lại
Lên đầu trang