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 4c5eef5db..0f44e3c1c 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 @@ -420,26 +420,28 @@ public class KafkaMessageChannelBinder extends @Override protected MessageHandler getErrorMessageHandler(final ConsumerDestination destination, final String group, final ExtendedConsumerProperties extendedConsumerProperties) { - if (extendedConsumerProperties.getExtension().isEnableDlq()) { + KafkaConsumerProperties kafkaConsumerProperties = extendedConsumerProperties.getExtension(); + if (kafkaConsumerProperties.isEnableDlq()) { + KafkaProducerProperties dlqProducerProperties = kafkaConsumerProperties.getDlqProducerProperties(); ProducerFactory producerFactory = this.transactionManager != null ? this.transactionManager.getProducerFactory() : getProducerFactory(null, - new ExtendedProducerProperties<>(extendedConsumerProperties.getExtension().getDlqProducerProperties())); + new ExtendedProducerProperties<>(dlqProducerProperties)); final KafkaTemplate kafkaTemplate = new KafkaTemplate<>(producerFactory); - String dlqName = StringUtils.hasText(extendedConsumerProperties.getExtension().getDlqName()) - ? extendedConsumerProperties.getExtension().getDlqName() + String dlqName = StringUtils.hasText(kafkaConsumerProperties.getDlqName()) + ? kafkaConsumerProperties.getDlqName() : "error." + destination.getName() + "." + group; - @SuppressWarnings({"unchecked", "raw"}) + @SuppressWarnings({"unchecked", "rawtypes"}) DlqSender dlqSender = new DlqSender(kafkaTemplate, dlqName); return message -> { final ConsumerRecord record = message.getHeaders() .get(KafkaHeaders.RAW_DATA, ConsumerRecord.class); - if (extendedConsumerProperties.isUseNativeDecoding()) { + if (this.transactionManager == null && extendedConsumerProperties.isUseNativeDecoding()) { if (record != null) { - Map configuration = extendedConsumerProperties.getExtension().getDlqProducerProperties().getConfiguration(); + Map configuration = dlqProducerProperties.getConfiguration(); if (record.key() != null && !record.key().getClass().isInstance(byte[].class)) { ensureDlqMessageCanBeProperlySerialized( configuration,