diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java
index 13368ad2d..acb8f8959 100644
--- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java
+++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java
@@ -22,6 +22,7 @@ import java.util.Map;
/**
* @author Marius Bogoevici
* @author Ilayaperumal Gopinathan
+ * @author Soby Chacko
*
*
* Thanks to Laszlo Szabo for providing the initial patch for generic property support.
@@ -41,6 +42,8 @@ public class KafkaConsumerProperties {
private String dlqName;
+ private KafkaProducerProperties dlqProducerProperties = new KafkaProducerProperties();
+
private int recoveryInterval = 5000;
private String[] trustedPackages;
@@ -133,4 +136,12 @@ public class KafkaConsumerProperties {
public void setTrustedPackages(String[] trustedPackages) {
this.trustedPackages = trustedPackages;
}
+
+ public KafkaProducerProperties getDlqProducerProperties() {
+ return dlqProducerProperties;
+ }
+
+ public void setDlqProducerProperties(KafkaProducerProperties dlqProducerProperties) {
+ this.dlqProducerProperties = dlqProducerProperties;
+ }
}
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 b55730c0a..ea26c845f 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
@@ -186,6 +186,11 @@ dlqName::
The name of the DLQ topic to receive the error messages.
+
Default: null (If not specified, messages that result in errors will be forwarded to a topic named `error..`).
+dlqProducerProperties::
+ Using this, dlq specific producer properties can be set.
+ All the properties available through kafka producer properties can be set through this property.
++
+Default: Default Kafka producer properties.
[[kafka-producer-properties]]
=== Kafka Producer Properties
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 4ce81797b..4c5eef5db 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
@@ -18,7 +18,6 @@ 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;
@@ -28,6 +27,7 @@ import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.UUID;
+import java.util.function.Predicate;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
@@ -41,7 +41,6 @@ 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;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.cloud.stream.binder.AbstractMessageChannelBinder;
@@ -422,21 +421,36 @@ public class KafkaMessageChannelBinder extends
protected MessageHandler getErrorMessageHandler(final ConsumerDestination destination, final String group,
final ExtendedConsumerProperties extendedConsumerProperties) {
if (extendedConsumerProperties.getExtension().isEnableDlq()) {
- ProducerFactory producerFactory = this.transactionManager != null
+ ProducerFactory,?> producerFactory = this.transactionManager != null
? this.transactionManager.getProducerFactory()
- : getProducerFactory(null, new ExtendedProducerProperties<>(new KafkaProducerProperties()));
- final KafkaTemplate kafkaTemplate = new KafkaTemplate<>(producerFactory);
+ : getProducerFactory(null,
+ new ExtendedProducerProperties<>(extendedConsumerProperties.getExtension().getDlqProducerProperties()));
+ final KafkaTemplate,?> kafkaTemplate = new KafkaTemplate<>(producerFactory);
+ String dlqName = StringUtils.hasText(extendedConsumerProperties.getExtension().getDlqName())
+ ? extendedConsumerProperties.getExtension().getDlqName()
+ : "error." + destination.getName() + "." + group;
+
+ @SuppressWarnings({"unchecked", "raw"})
+ DlqSender,?> dlqSender = new DlqSender(kafkaTemplate, dlqName);
+
return message -> {
final ConsumerRecord, ?> record = message.getHeaders()
.get(KafkaHeaders.RAW_DATA, ConsumerRecord.class);
- final byte[] key = record.key() != null ? Utils.toArray(ByteBuffer.wrap((byte[]) record.key()))
- : null;
- final byte[] payload = record.value() != null
- ? Utils.toArray(ByteBuffer.wrap((byte[]) record.value()))
- : null;
- String dlqName = StringUtils.hasText(extendedConsumerProperties.getExtension().getDlqName())
- ? extendedConsumerProperties.getExtension().getDlqName()
- : "error." + destination.getName() + "." + group;
+
+ if (extendedConsumerProperties.isUseNativeDecoding()) {
+ if (record != null) {
+ Map configuration = extendedConsumerProperties.getExtension().getDlqProducerProperties().getConfiguration();
+ if (record.key() != null && !record.key().getClass().isInstance(byte[].class)) {
+ ensureDlqMessageCanBeProperlySerialized(
+ configuration,
+ (Map config) -> !config.containsKey("key.serializer"), "Key");
+ }
+ if (record.value() != null && !record.value().getClass().isInstance(byte[].class)) {
+ ensureDlqMessageCanBeProperlySerialized(configuration,
+ (Map config) -> !config.containsKey("value.serializer"), "Payload");
+ }
+ }
+ }
Headers kafkaHeaders = new RecordHeaders(record.headers().toArray());
kafkaHeaders.add(new RecordHeader(X_ORIGINAL_TOPIC,
@@ -448,37 +462,21 @@ public class KafkaMessageChannelBinder extends
kafkaHeaders.add(new RecordHeader(X_EXCEPTION_STACKTRACE,
getStackTraceAsString(throwable).getBytes(StandardCharsets.UTF_8)));
}
- ProducerRecord producerRecord = new ProducerRecord<>(dlqName, record.partition(),
- key, payload, kafkaHeaders);
- ListenableFuture> sentDlq = kafkaTemplate.send(producerRecord);
- sentDlq.addCallback(new ListenableFutureCallback>() {
- StringBuilder sb = new StringBuilder().append(" a message with key='")
- .append(toDisplayString(ObjectUtils.nullSafeToString(key), 50)).append("'")
- .append(" and payload='")
- .append(toDisplayString(ObjectUtils.nullSafeToString(payload), 50))
- .append("'").append(" received from ")
- .append(record.partition());
-
- @Override
- public void onFailure(Throwable ex) {
- KafkaMessageChannelBinder.this.logger.error(
- "Error sending to DLQ " + sb.toString(), ex);
- }
-
- @Override
- public void onSuccess(SendResult result) {
- if (KafkaMessageChannelBinder.this.logger.isDebugEnabled()) {
- KafkaMessageChannelBinder.this.logger.debug(
- "Sent to DLQ " + sb.toString());
- }
- }
-
- });
+ dlqSender.sendToDlq(record, kafkaHeaders);
};
}
return null;
}
+ private static void ensureDlqMessageCanBeProperlySerialized(Map configuration,
+ Predicate