Exception Handling trong Spring for Apache Kafka Streams – Phần 1

Để handle các exception có thể xảy ra trong các ứng dụng có sử dụng Spring for Apache Kafka Streams, các bạn có thể khai báo các cấu hình như sau:

Configuration Handles Ví dụ
DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG Handle các exception liên quan đến việc deserialization các message trong consumer Chuỗi JSON không hợp lệ, sai đối tượng SerDer, …
PROCESSING_EXCEPTION_HANDLER_CLASS_CONFIG Lỗi xảy ra khi process các message Lỗi NullPointerException, ArithmeticException, lỗi xảy ra khi handle business logic trong các phương thức map(), join(), filter(), …
PRODUCTION_EXCEPTION_HANDLER_CLASS_CONFIG Lỗi xảy ra khi produce các output messages Serialization bị lỗi, broker reject một message, …

Cụ thể như thế nào? Chúng ta hãy cùng nhau tìm hiểu trong bài viết này, các bạn nhé!

Đầu tiên, mình sẽ tạo mới một Maven project, tương tự như trong ví dụ của bài viết trước, như sau:

Các dependency liên quan đến Spring for Apache Kafka Streams, mình sẽ khai báo như sau:

Chúng ta cần khai báo dependency management cho các thư viện Jackson để các thư viện này tương thích với nhau, như sau:

Mình cũng khai báo dependency cho thư viện Slf4J và Logback để chúng ta có thể log tất cả những lỗi xảy ra như sau các bạn nhé:

Tập tin cấu hình của thư viện Logback trong thư mục /src/main/resources sẽ có nội dung như sau:

Các bạn cũng cấu hình thông tin Apache Kafka Server và enable kafka streams như sau các bạn nhé:

Điều đầu tiên mình cần nói với các bạn là giá trị của các cấu hình DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, PROCESSING_EXCEPTION_HANDLER_CLASS_CONFIG, PRODUCTION_EXCEPTION_HANDLER_CLASS_CONFIG lần lượt là các implementation cho các interface là DeserializationExceptionHandler, ProcessingExceptionHandlerProductionExceptionHandler, các bạn nhé!

DeserializationExceptionHandler

Như mình đã nói ở trên, DeserializationExceptionHandler sẽ handle các exception liên quan đến việc deserialization các message từ producer trong ứng dụng của chúng ta.

Cho phần implementation của interface DeserializationExceptionHandler, mặc định thì thư viện Apache Kafka Streams hỗ trợ 2 implementation là:

  • LogAndContinueExceptionHandler
  • LogAndFailExceptionHandler

còn thư viện Spring for Apache Kafka hỗ trợ thêm implementation:

  • RecoveringDeserializationExceptionHandler

các bạn nhé!

LogAndContinueExceptionHandler sẽ log lỗi deserialization, bỏ qua message bị lỗi và tiếp tục process các message khác, còn LogAndFailExceptionHandler cũng sẽ log lỗi deserialization nhưng sẽ stop application. Implementation RecoveringDeserializationExceptionHandler thì thay vì chỉ log lỗi và bỏ qua message bị lỗi, nó sẽ gửi message đó vào một Dead Letter Queue (DLQ) topic để các bạn có thể recover message này sau.

Ví dụ, mình định nghĩa một stream topology như sau:

Stream topology này sẽ consume topic users với phần serialization/ deserialization cho key là String, value là class User. Nội dung của class User như sau:

Class để chạy ứng dụng:

Lúc này, nếu chạy ứng dụng và publish vào topic users một message với key là “001”, value là ‘{“name”:”. Các bạn sẽ thấy ứng dụng của chúng ta sẽ log lỗi deserialization và stop, sử dụng class deserialization LogAndFailExceptionHandler, như sau:

Nếu các bạn cấu hình exception handler cho phần deserialization này sử dụng class LogAndContinueExceptionHandler:

và chạy lại ứng dụng, các bạn sẽ thấy class LogAndContinueExceptionHandler sẽ handle error:

Nó chỉ đơn giản là log lại error và ứng dụng của chúng ta sẽ tiếp tục chạy, các bạn nhé!

ProcessingExceptionHandler

Mặc định thì thư viện Apache Kafka Streams hỗ trợ 2 implementation cho interface ProcessingExceptionHandler là:

  • LogAndContinueProcessingExceptionHandler
  • LogAndFailProcessingExceptionHandler

còn thư viện Spring for Apache Kafka hỗ trợ thêm implementation:

  • RecoveringProcessingExceptionHandler

các bạn nhé!

Tương tự như DeserializationExceptionHandler, LogAndContinueProcessingExceptionHandler sẽ , LogAndFailProcessingExceptionHandlerRecoveringProcessingExceptionHandler, còn RecoveringProcessingExceptionHandler

Ví dụ như mình có một stream topology như sau:

Nếu bây giờ, các bạn chạy ứng dụng và publish một message với key là “001”, value là “ERROR”, các bạn sẽ thấy ứng dụng log lỗi sử dụng class LogAndFailProcessingExceptionHandler và stop như sau:

Các bạn có thể thay đổi class exception handler trong trường hợp này sử dụng cấu hình PROCESSING_EXCEPTION_HANDLER_CLASS_CONFIG như sau:

Chạy lại ví dụ và publish lại message với nội dung như trên, các bạn sẽ thấy ứng dụng của chúng ta cũng sẽ log lỗi nhưng sẽ không stop nữa.

Các bạn có thể đọc tiếp phần 2 tại đây.

Các bạn có thể xem đầy đủ các exception handler trong video ở đây:

Add Comment