Trong các bài viết trước Exception Handling trong Spring for Apache Kafka Streams – Phần 1 và Exception 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, ProcessingExceptionHandler và ProductionExceptionHandler 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:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 |
@Bean(name = KafkaStreamsDefaultConfiguration.DEFAULT_STREAMS_CONFIG_BEAN_NAME) public KafkaStreamsConfiguration kafkaStreamsConfiguration() { Map<String, Object> props = new HashMap(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "spring-kafka-streams-example"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); props.put(StreamsConfig.ERRORS_DEAD_LETTER_QUEUE_TOPIC_NAME_CONFIG, "default-dlq-topic"); props.put( StreamsConfig.DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, LogAndContinueExceptionHandler.class.getName()); props.put( StreamsConfig.PROCESSING_EXCEPTION_HANDLER_CLASS_CONFIG, LogAndContinueProcessingExceptionHandler.class.getName()); props.put( StreamsConfig.PRODUCTION_EXCEPTION_HANDLER_CLASS_CONFIG, MyProductionExceptionHandler.class.getName()); return new KafkaStreamsConfiguration(props); } |
Ở đâ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à:
|
1 2 3 4 5 6 7 8 |
@Bean public KStream<String, User> users(StreamsBuilder builder) { KStream<String, User> stream = builder.stream("users"); stream.to("output", Produced.with(Serdes.String(), new JacksonJsonSerde<>(User.class))); return stream; } |
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:
- Sử dụng implementation của functional interface
KafkaStreamsDeadLetterDestinationResolver - Sử dụng cấu hình
StreamsConfig.ERRORS_DEAD_LETTER_QUEUE_TOPIC_NAME_CONFIGcủa Apache Kafka Streams như mình đã đề cập ở trên. - Sử dụng implementation của interface
ConsumerRecordRecoverervớ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, ConsumerRecord và Exception trong tham số của phương thức resolve(), với định nghĩa bean, ví dụ như sau:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 |
@Bean public KafkaStreamsDeadLetterDestinationResolver resolver() { return (context, message, exception) -> { if (context.processorNodeId().equals("huongdanjava")) { return new TopicPartition("huongdanjava-dlq-topic", -1); } if (message.value() instanceof String m && m.equals("ERROR")) { return new TopicPartition("error-message-dlq-topic", -1); } if (exception instanceof NumberFormatException) { return new TopicPartition("invalid-dlq-topic", -1); } return new TopicPartition("default1-dlq-topic", 0); }; } |
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:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 |
@Bean(name = KafkaStreamsDefaultConfiguration.DEFAULT_STREAMS_CONFIG_BEAN_NAME) public KafkaStreamsConfiguration kafkaStreamsConfiguration( KafkaStreamsDeadLetterDestinationResolver resolver) { Map<String, Object> props = new HashMap(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "spring-kafka-streams-example"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); props.put(StreamsConfig.ERRORS_DEAD_LETTER_QUEUE_TOPIC_NAME_CONFIG, "default-dlq-topic"); props.put( StreamsConfig.DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, RecoveringDeserializationExceptionHandler.class.getName()); props.put( StreamsConfig.PROCESSING_EXCEPTION_HANDLER_CLASS_CONFIG, RecoveringProcessingExceptionHandler.class.getName()); props.put( StreamsConfig.PRODUCTION_EXCEPTION_HANDLER_CLASS_CONFIG, RecoveringProductionExceptionHandler.class.getName()); props.put(RecoveringDeserializationExceptionHandler.DLQ_DESTINATION_RESOLVER, resolver); props.put(RecoveringProcessingExceptionHandler.DLQ_DESTINATION_RESOLVER, resolver); props.put(RecoveringProductionExceptionHandler.DLQ_DESTINATION_RESOLVER, resolver); return new KafkaStreamsConfiguration(props); } |
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:
|
1 2 3 4 |
@Bean public StreamsBuilderFactoryBeanConfigurer streamsBuilderFactoryBeanConfigurer() { return sfb -> sfb.setDeadLetterTopicName("default2-dlq-topic"); } |
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:
|
1 2 3 4 5 |
@Bean public DeadLetterPublishingRecoverer recoverer(KafkaTemplate kafkaTemplate) { return new DeadLetterPublishingRecoverer( kafkaTemplate, (record, ex) -> new TopicPartition("default3-dlq-topic", 0)); } |
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é:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 |
@Bean KafkaTemplate<byte[], byte[]> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); } @Bean public ProducerFactory<byte[], byte[]> producerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class); return new DefaultKafkaProducerFactory<>(props); } |
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:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 |
@Bean(name = KafkaStreamsDefaultConfiguration.DEFAULT_STREAMS_CONFIG_BEAN_NAME) public KafkaStreamsConfiguration kafkaStreamsConfiguration( KafkaStreamsDeadLetterDestinationResolver resolver, DeadLetterPublishingRecoverer recoverer) { Map<String, Object> props = new HashMap(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "spring-kafka-streams-example"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); props.put(StreamsConfig.ERRORS_DEAD_LETTER_QUEUE_TOPIC_NAME_CONFIG, "default-dlq-topic"); props.put( StreamsConfig.DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, RecoveringDeserializationExceptionHandler.class.getName()); props.put( StreamsConfig.PROCESSING_EXCEPTION_HANDLER_CLASS_CONFIG, RecoveringProcessingExceptionHandler.class.getName()); props.put( StreamsConfig.PRODUCTION_EXCEPTION_HANDLER_CLASS_CONFIG, RecoveringProductionExceptionHandler.class.getName()); props.put(RecoveringDeserializationExceptionHandler.DLQ_DESTINATION_RESOLVER, resolver); props.put(RecoveringProcessingExceptionHandler.DLQ_DESTINATION_RESOLVER, resolver); props.put(RecoveringProductionExceptionHandler.DLQ_DESTINATION_RESOLVER, resolver); props.put(RecoveringDeserializationExceptionHandler.RECOVERER, recoverer); props.put(RecoveringProcessingExceptionHandler.RECOVERER, recoverer); props.put(RecoveringProductionExceptionHandler.RECOVERER, recoverer); return new KafkaStreamsConfiguration(props); } |
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é:

