From d9196c32e50605c5e84a4b7ca3b610e6a641a398 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 10 Sep 2019 14:30:47 -0400 Subject: [PATCH] GH-153: Kinesis: Fix embeddedHeaders and converter Fixes https://github.com/spring-projects/spring-integration-aws/issues/153 Right now we ignore an `embeddedHeadersMapper` result when we also have a `converter`. * Move `converter` logic for `batch` mode into `else` of the `if (KclMessageDrivenChannelAdapter.this.embeddedHeadersMapper != null) {`. This way we cover all the possible use-case, and don't override results of each other and also don't perform extra logic for nothing --- .../KclMessageDrivenChannelAdapter.java | 51 +++++++++---------- .../KinesisMessageDrivenChannelAdapter.java | 51 +++++++++---------- 2 files changed, 48 insertions(+), 54 deletions(-) diff --git a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java index d2c2db7..726e9b4 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java +++ b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java @@ -374,37 +374,34 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport { } private void processMultipleRecords(List records, IRecordProcessorCheckpointer checkpointer) { - Object payload = records; - + AbstractIntegrationMessageBuilder messageBuilder = getMessageBuilderFactory().withPayload(records); if (KclMessageDrivenChannelAdapter.this.embeddedHeadersMapper != null) { - payload = records.stream() - .map(this::prepareMessageForRecord) - .map(AbstractIntegrationMessageBuilder::build) + List> payload = + records.stream() + .map(this::prepareMessageForRecord) + .map(AbstractIntegrationMessageBuilder::build) + .collect(Collectors.toList()); + + messageBuilder = getMessageBuilderFactory().withPayload(payload); + } + else if (KclMessageDrivenChannelAdapter.this.converter != null) { + final List partitionKeys = new ArrayList<>(); + final List sequenceNumbers = new ArrayList<>(); + + List payload = records.stream() + .map(r -> { + partitionKeys.add(r.getPartitionKey()); + sequenceNumbers.add(r.getSequenceNumber()); + + return KclMessageDrivenChannelAdapter.this.converter.convert(r.getData().array()); + }) .collect(Collectors.toList()); + + messageBuilder = getMessageBuilderFactory().withPayload(payload) + .setHeader(AwsHeaders.RECEIVED_PARTITION_KEY, partitionKeys) + .setHeader(AwsHeaders.RECEIVED_SEQUENCE_NUMBER, sequenceNumbers); } - final List partitionKeys; - final List sequenceNumbers; - if (KclMessageDrivenChannelAdapter.this.converter != null) { - partitionKeys = new ArrayList<>(); - sequenceNumbers = new ArrayList<>(); - - payload = records.stream().map(r -> { - partitionKeys.add(r.getPartitionKey()); - sequenceNumbers.add(r.getSequenceNumber()); - - return KclMessageDrivenChannelAdapter.this.converter.convert(r.getData().array()); - }).collect(Collectors.toList()); - } - else { - partitionKeys = null; - sequenceNumbers = null; - } - - AbstractIntegrationMessageBuilder messageBuilder = getMessageBuilderFactory().withPayload(payload) - .setHeader(AwsHeaders.RECEIVED_PARTITION_KEY, partitionKeys) - .setHeader(AwsHeaders.RECEIVED_SEQUENCE_NUMBER, sequenceNumbers); - performSend(messageBuilder, records, checkpointer); } diff --git a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java index 4dbd5f0..684bd2f 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java +++ b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java @@ -995,37 +995,34 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i } private void processMultipleRecords(List records) { - Object payload = records; - + AbstractIntegrationMessageBuilder messageBuilder = getMessageBuilderFactory().withPayload(records); if (KinesisMessageDrivenChannelAdapter.this.embeddedHeadersMapper != null) { - payload = records.stream() - .map(this::prepareMessageForRecord) - .map(AbstractIntegrationMessageBuilder::build) + List> payload = + records.stream() + .map(this::prepareMessageForRecord) + .map(AbstractIntegrationMessageBuilder::build) + .collect(Collectors.toList()); + + messageBuilder = getMessageBuilderFactory().withPayload(payload); + } + else if (KinesisMessageDrivenChannelAdapter.this.converter != null) { + final List partitionKeys = new ArrayList<>(); + final List sequenceNumbers = new ArrayList<>(); + + List payload = records.stream() + .map(r -> { + partitionKeys.add(r.getPartitionKey()); + sequenceNumbers.add(r.getSequenceNumber()); + + return KinesisMessageDrivenChannelAdapter.this.converter.convert(r.getData().array()); + }) .collect(Collectors.toList()); + + messageBuilder = getMessageBuilderFactory().withPayload(payload) + .setHeader(AwsHeaders.RECEIVED_PARTITION_KEY, partitionKeys) + .setHeader(AwsHeaders.RECEIVED_SEQUENCE_NUMBER, sequenceNumbers); } - final List partitionKeys; - final List sequenceNumbers; - if (KinesisMessageDrivenChannelAdapter.this.converter != null) { - partitionKeys = new ArrayList<>(); - sequenceNumbers = new ArrayList<>(); - - payload = records.stream().map(r -> { - partitionKeys.add(r.getPartitionKey()); - sequenceNumbers.add(r.getSequenceNumber()); - - return KinesisMessageDrivenChannelAdapter.this.converter.convert(r.getData().array()); - }).collect(Collectors.toList()); - } - else { - partitionKeys = null; - sequenceNumbers = null; - } - - AbstractIntegrationMessageBuilder messageBuilder = getMessageBuilderFactory().withPayload(payload) - .setHeader(AwsHeaders.RECEIVED_PARTITION_KEY, partitionKeys) - .setHeader(AwsHeaders.RECEIVED_SEQUENCE_NUMBER, sequenceNumbers); - performSend(messageBuilder, records); }