Current state: Under Discussion
Discussion thread: here
JIRA: here
Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).
KIP-334 introduced into the Consumer the RecordDeserializationException with offsets information. That is useful to skip a poison pill but as you do not have access to the Record, it still prevents easy implementation of dead letter queue or simply logging the faulty data.
The change in the RecordDeserializationException would be modification of the constructor and a new method to get the ConsumerRecord.
As it’s only addition, it would still be compatible with existing consumer code.
public class RecordDeserializationException extends SerializationException {
@Deprecated
p
ublic RecordDeserializationException(TopicPartition partition,
long offset,
String message,
Throwable cause);
// New constructor
public RecordDeserializationException(TopicPartition partition,
Record record,
String message,
Throwable cause);
// New method
public Record record();
} |
We propose to include the Record object to the RecordDeserializationException and provide a getConsumerRecord method that will lazily create the needed wrapper object if needed. This would still allow us to access offsets and then skip the record but also to fetch the concerned data for specific processing by sending it to a dead letter queue or whatever action that makes sense.
Here is an example of basic usage to implement a DLQ feature
{
// …
try (KafkaConsumer<String, SimpleValue> consumer = new KafkaConsumer<>(settings())) {
// Subscribe to our topic
LOGGER.info("Subscribing to topic " + KAFKA_TOPIC);
consumer.subscribe(List.of(KAFKA_TOPIC));
LOGGER.info("Subscribed !");
try (KafkaProducer<byte[], byte[]> dlqProducer = new KafkaProducer<>(producerSettings())) {
//noinspection InfiniteLoopStatement
while (true) {
try {
final var records = consumer.poll(POLL_TIMEOUT);
LOGGER.info("poll() returned {} records", records.count());
for (var record : records) {
LOGGER.info("Fetch record key={} value={}", record.key(), record.value());
// Any processing
// ...
}
} catch (RecordDeserializationException re) {
long offset = re.offset();
Throwable t = re.getCause();
LOGGER.error("Failed to consume at partition={} offset={}", re.topicPartition().partition(), offset, t);
sendDlqRecord(dlqProducer, re.record());
LOGGER.info("Skipping offset={}", offset);
consumer.seek(re.topicPartition(), offset + 1);
} catch (Exception e) {
LOGGER.error("Failed to consume", e);
}
}
}
} finally {
LOGGER.info("Closing consumer");
}
}
void sendDlqRecord(KafkaProducer<byte[], byte[]> dlqProducer, Record record) {
var dlqRecord = new ProducerRecord<>(DLQ_TOPIC, Utils.toNullableArray(record.key()), Utils.toNullableArray(record.value()));
try {
dlqProducer.send(dlqRecord).get();
LOGGER.info("Record sent to DLQ");
} catch (Exception e) {
LOGGER.error("Failed to send corrupted record to DLQ", e);
}
} |
This change is backward compatible and will have no impact on existing application.