Addressing PR review comment for gh-58

- Add a check for ignoring dlq producer properties set by the user
  if transaction management is enabled
- Polishing
This commit is contained in:
Soby Chacko
2017-11-14 14:59:23 -05:00
parent 66f194dd93
commit 2f6d5df7bc

View File

@@ -420,26 +420,28 @@ public class KafkaMessageChannelBinder extends
@Override
protected MessageHandler getErrorMessageHandler(final ConsumerDestination destination, final String group,
final ExtendedConsumerProperties<KafkaConsumerProperties> 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<String, String> configuration = extendedConsumerProperties.getExtension().getDlqProducerProperties().getConfiguration();
Map<String, String> configuration = dlqProducerProperties.getConfiguration();
if (record.key() != null && !record.key().getClass().isInstance(byte[].class)) {
ensureDlqMessageCanBeProperlySerialized(
configuration,