Exception handling với dead letter queue trong Spring for Apache Kafka


Trong các bài viết trước Exception Handling trong Spring for Apache Kafka Streams – Phần 1Exception Handling trong Spring for Apache Kafka Streams – Phần 2, chúng ta đã tìm hiểu về các loại exception handlers trong Spring for Apache Kafka và Apache Kafka Streams. Đối với các implementation LogAndContinue… và LogAndFail… thì mặc định khi có lỗi xảy ra, ứng dụng của chúng ta sẽ log error và tiếp tục chạy hoặc dừng lại. Các bạn cũng có thể publish message bị lỗi này vào một dead letter queue (DLQ) topic để có thể review lại message bị lỗi và re-process message đó lại sau. 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é!

MÌnh sẽ sử dụng lại ví dụ trong các bài viết trước.

Đầu tiên, mình sẽ nói với các bạn một xíu về DLQ topic. DLQ topic là một topic sẽ lưu lại tất cả các message được publish vào Apache Kafka server nhưng ứng dụng của chúng ta không thể process chúng được, sau vài lần retry. Chúng ta sẽ cần review manually để tìm phương án thích hợp cho những message bị lỗi này.

Đối với các implementation của thư viện Apache Kafka Streams cho các interface DeserializationExceptionHandler, ProcessingExceptionHandlerProductionExceptionHandler thì để publish các message bị lỗi vào DLQ topic thì các bạn chỉ cần cấu hình thêm cho property StreamsConfig.ERRORS_DEAD_LETTER_QUEUE_TOPIC_NAME_CONFIG tên của DLQ topic là được, ví dụ như sau:

Ở đây mình đang cấu hình tên của DLQ topic là “default-dlq-topic” đó các bạn!

Bây giờ thì chạy lại ví dụ với stream tolology là:

Sau đó thì publish vào topic users với key là 001, value là Khanh, các bạn sẽ thấy ứng dụng của chúng ta sẽ log lỗi và message lỗi sẽ được publish vào DLQ topic như sau:

Đối với các implementation của thư viện Spring for Apache Kafka như RecoveringDeserializationExceptionHandler, RecoveringProcessingExceptionHandler, RecoveringProductionExceptionHandler thì mặc định, các implementation này sẽ luôn sử dụng một DLQ topic để publish message bị lỗi và recover stream topology các bạn nhé! Nếu thông tin DLQ topic không được khai báo thì ứng dụng của chúng ta sẽ log lỗi và stop lại.

Có 3 cách chúng ta có thể định nghĩa thông tin DLQ topic với các implementation này theo thứ tự ưu tiên như sau:

  1. Sử dụng implementation của functional interface KafkaStreamsDeadLetterDestinationResolver
  2. Sử dụng cấu hình StreamsConfig.ERRORS_DEAD_LETTER_QUEUE_TOPIC_NAME_CONFIG của Apache Kafka Streams như mình đã đề cập ở trên.
  3. Sử dụng implementation của interface ConsumerRecordRecoverer với implementation mặc định là DeadLetterPublishingRecoverer

Với functional interface KafkaStreamsDeadLetterDestinationResolver, các bạn có thể resolve DLQ topic name bằng cách sử dụng thông tin của các đối tượng ErrorHandlerContext, ConsumerRecordException trong tham số của phương thức resolve(), với định nghĩa bean, ví dụ như sau:

Sau khi đã định nghĩa resolver bean thì các bạn cần cấu hình nó cho các Recovering…ExceptionHandler của Spring for Apache Kakfa như sau:

Với cách thứ 2 sử dụng cấu hình StreamsConfig.ERRORS_DEAD_LETTER_QUEUE_TOPIC_NAME_CONFIG thì các bạn có thể định nghĩa như trên hoặc định nghĩa một bean của interface StreamsBuilderFactoryBeanConfigurer để set giá trị của DLQ topic name như sau:

Với cách thứ 3 thì các bạn có thể định nghĩa bean cho interface ConsumerRecordRecoverer sử dụng implementation DeadLetterPublishingRecoverer như sau:

Với cách này thì các bạn cần định nghĩa bean cho class KafkaTemplate các bạn nhé:

Kiểu dữ liệu của KafkaTemplate phải là <byte[], byte[]> và class serializer phải ByteArraySerializer các bạn nhé! Bởi vì raw source record sẽ được gửi tới DLQ topic.

Tương tự như resolver, cho bean recoverer này, các bạn cũng cần cấu hình nó cho các Recovering…ExceptionHandler của Spring for Apache Kafka như sau:

Chạy lại ví dụ với 3 cách cấu hình trên, các bạn sẽ thấy ứng dụng của chúng ta sẽ consume lại message vừa bị lỗi ở trên. Vì không process được nên nó sẽ publish vào DLQ topic theo thứ tự sẽ là “default1-dlq-topic”, các bạn nhé:

Add Comment