From 035cc1a0054e7dad58d62ba212944015bcad4e56 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Wed, 8 Nov 2017 12:56:45 -0500 Subject: [PATCH] GH-228: Enhance DLQ messages with failure info Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/228 --- .../src/main/asciidoc/overview.adoc | 4 ++- .../kafka/KafkaMessageChannelBinder.java | 32 ++++++++++++++++++- .../stream/binder/kafka/KafkaBinderTests.java | 8 +++++ 3 files changed, 42 insertions(+), 2 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc index 783723135..b55730c0a 100644 --- a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc +++ b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc @@ -174,6 +174,8 @@ enableDlq:: By default, messages that result in errors will be forwarded to a topic named `error..`. The DLQ topic name can be configurable via the property `dlqName`. This provides an alternative option to the more common Kafka replay scenario for the case when the number of errors is relatively small and replaying the entire original topic may be too cumbersome. + See <> processing for more information. + Starting with _version 2.0_, messages sent to the DLQ topic are enhanced with the following headers: `x-original-topic`, `x-exception-message` and `x-exception-stacktrace` as `byte[]`. + Default: `false`. configuration:: @@ -559,7 +561,7 @@ The payload of the `ErrorMessage` for a send failure is a `KafkaSendFailureExcep * `failedMessage` - the spring-messaging `Message` that failed to be sent. * `record` - the raw `ProducerRecord` that was created from the `failedMessage` -There is no automatic handling of these exceptions (such as sending to a <>); you can consume these exceptions with your own Spring Integration flow. +There is no automatic handling of producer exceptions (such as sending to a <>); you can consume these exceptions with your own Spring Integration flow. [[kafka-metrics]] == Kafka Metrics diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 1185d0f3e..4ce81797b 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -16,7 +16,10 @@ package org.springframework.cloud.stream.binder.kafka; +import java.io.PrintWriter; +import java.io.StringWriter; import java.nio.ByteBuffer; +import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; @@ -34,6 +37,8 @@ import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.header.Headers; +import org.apache.kafka.common.header.internals.RecordHeader; +import org.apache.kafka.common.header.internals.RecordHeaders; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.serialization.ByteArraySerializer; import org.apache.kafka.common.utils.Utils; @@ -103,6 +108,13 @@ public class KafkaMessageChannelBinder extends AbstractMessageChannelBinder, ExtendedProducerProperties, KafkaTopicProvisioner> implements ExtendedPropertiesBinder { + public static final String X_EXCEPTION_STACKTRACE = "x-exception-stacktrace"; + + public static final String X_EXCEPTION_MESSAGE = "x-exception-message"; + + public static final String X_ORIGINAL_TOPIC = "x-original-topic"; + + private final KafkaBinderConfigurationProperties configurationProperties; private final Map topicsInUse = new HashMap<>(); @@ -425,8 +437,19 @@ public class KafkaMessageChannelBinder extends String dlqName = StringUtils.hasText(extendedConsumerProperties.getExtension().getDlqName()) ? extendedConsumerProperties.getExtension().getDlqName() : "error." + destination.getName() + "." + group; + + Headers kafkaHeaders = new RecordHeaders(record.headers().toArray()); + kafkaHeaders.add(new RecordHeader(X_ORIGINAL_TOPIC, + record.topic().getBytes(StandardCharsets.UTF_8))); + if (message.getPayload() instanceof Throwable) { + Throwable throwable = (Throwable) message.getPayload(); + kafkaHeaders.add(new RecordHeader(X_EXCEPTION_MESSAGE, + throwable.getMessage().getBytes(StandardCharsets.UTF_8))); + kafkaHeaders.add(new RecordHeader(X_EXCEPTION_STACKTRACE, + getStackTraceAsString(throwable).getBytes(StandardCharsets.UTF_8))); + } ProducerRecord producerRecord = new ProducerRecord<>(dlqName, record.partition(), - key, payload, record.headers()); + key, payload, kafkaHeaders); ListenableFuture> sentDlq = kafkaTemplate.send(producerRecord); sentDlq.addCallback(new ListenableFutureCallback>() { StringBuilder sb = new StringBuilder().append(" a message with key='") @@ -509,6 +532,13 @@ public class KafkaMessageChannelBinder extends return original.substring(0, maxCharacters) + "..."; } + private String getStackTraceAsString(Throwable cause) { + StringWriter stringWriter = new StringWriter(); + PrintWriter printWriter = new PrintWriter(stringWriter, true); + cause.printStackTrace(printWriter); + return stringWriter.getBuffer().toString(); + } + private final class ProducerConfigurationMessageHandler extends KafkaProducerMessageHandler implements Lifecycle { diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index ea615c093..088c9894e 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -34,8 +34,10 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import com.fasterxml.jackson.databind.ObjectMapper; + import kafka.utils.ZKStringSerializer$; import kafka.utils.ZkUtils; + import org.I0Itec.zkclient.ZkClient; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; @@ -464,6 +466,12 @@ public class KafkaBinderTests extends assertThat(receivedMessage).isNotNull(); assertThat(receivedMessage.getPayload()).isEqualTo(testMessagePayload.getBytes()); assertThat(handler.getInvocationCount()).isEqualTo(consumerProperties.getMaxAttempts()); + assertThat(receivedMessage.getHeaders().get(KafkaMessageChannelBinder.X_ORIGINAL_TOPIC)) + .isEqualTo(producerName.getBytes(StandardCharsets.UTF_8)); + assertThat(new String((byte[]) receivedMessage.getHeaders().get(KafkaMessageChannelBinder.X_EXCEPTION_MESSAGE))) + .startsWith("failed to send Message to channel 'input'"); + assertThat(receivedMessage.getHeaders().get(KafkaMessageChannelBinder.X_EXCEPTION_STACKTRACE)) + .isNotNull(); binderBindUnbindLatency(); // verify we got a message on the dedicated error channel and the global (via bridge)