|
19 | 19 | package com.michelin.kafka.test.unit; |
20 | 20 |
|
21 | 21 | import static org.junit.jupiter.api.Assertions.assertEquals; |
| 22 | +import static org.junit.jupiter.api.Assertions.assertNull; |
22 | 23 | import static org.mockito.ArgumentMatchers.any; |
23 | 24 | import static org.mockito.Mockito.*; |
24 | 25 |
|
|
31 | 32 | import java.io.ObjectInputStream; |
32 | 33 | import java.nio.ByteBuffer; |
33 | 34 | import org.apache.kafka.clients.consumer.ConsumerRecord; |
| 35 | +import org.apache.kafka.common.TopicPartition; |
| 36 | +import org.apache.kafka.common.errors.RecordDeserializationException; |
| 37 | +import org.apache.kafka.common.header.Headers; |
| 38 | +import org.apache.kafka.common.record.TimestampType; |
34 | 39 | import org.junit.jupiter.api.BeforeEach; |
35 | 40 | import org.junit.jupiter.api.Test; |
36 | 41 | import org.mockito.ArgumentCaptor; |
@@ -142,4 +147,31 @@ void shouldHandleErrorWithNullThrowable() { |
142 | 147 | assertEquals(record.key(), capturedErrorModel.getKey()); |
143 | 148 | assertEquals(record.value(), capturedErrorModel.getValue()); |
144 | 149 | } |
| 150 | + |
| 151 | + @Test |
| 152 | + void shouldHandleErrorWithRecordDeserializationException() throws IOException { |
| 153 | + // Given |
| 154 | + ConsumerRecord<String, String> record = new ConsumerRecord<>("testTopic", 12, 13458L, "key", "value"); |
| 155 | + RecordDeserializationException rde = |
| 156 | + new RecordDeserializationException( |
| 157 | + RecordDeserializationException.DeserializationExceptionOrigin.KEY, |
| 158 | + new TopicPartition(record.topic(),record.partition()), |
| 159 | + record.offset(), |
| 160 | + 1764603801, |
| 161 | + TimestampType.CREATE_TIME, |
| 162 | + DefaultErrorProcessor.toByteBuffer("MyKey"), |
| 163 | + DefaultErrorProcessor.toByteBuffer("MyValue"), |
| 164 | + null, |
| 165 | + "Record deserialization error", |
| 166 | + null); |
| 167 | + |
| 168 | + // When |
| 169 | + errorProcessor.processError(rde, record, 32L); |
| 170 | + verify(mockDeadLetterProducer, times(1)).send(keyCaptor.capture(), valueCaptor.capture()); |
| 171 | + |
| 172 | + GenericErrorModel capturedErrorModel = valueCaptor.getValue(); |
| 173 | + |
| 174 | + assertEquals("Record deserialization error", capturedErrorModel.getCause()); |
| 175 | + assertEquals(record.topic(), capturedErrorModel.getTopic()); |
| 176 | + } |
145 | 177 | } |
0 commit comments