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); }