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
This commit is contained in:
@@ -374,37 +374,34 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport {
|
||||
}
|
||||
|
||||
private void processMultipleRecords(List<Record> 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<Message<Object>> payload =
|
||||
records.stream()
|
||||
.map(this::prepareMessageForRecord)
|
||||
.map(AbstractIntegrationMessageBuilder::build)
|
||||
.collect(Collectors.toList());
|
||||
|
||||
messageBuilder = getMessageBuilderFactory().withPayload(payload);
|
||||
}
|
||||
else if (KclMessageDrivenChannelAdapter.this.converter != null) {
|
||||
final List<String> partitionKeys = new ArrayList<>();
|
||||
final List<String> sequenceNumbers = new ArrayList<>();
|
||||
|
||||
List<Object> 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<String> partitionKeys;
|
||||
final List<String> 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);
|
||||
}
|
||||
|
||||
|
||||
@@ -995,37 +995,34 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i
|
||||
}
|
||||
|
||||
private void processMultipleRecords(List<Record> records) {
|
||||
Object payload = records;
|
||||
|
||||
AbstractIntegrationMessageBuilder<?> messageBuilder = getMessageBuilderFactory().withPayload(records);
|
||||
if (KinesisMessageDrivenChannelAdapter.this.embeddedHeadersMapper != null) {
|
||||
payload = records.stream()
|
||||
.map(this::prepareMessageForRecord)
|
||||
.map(AbstractIntegrationMessageBuilder::build)
|
||||
List<Message<Object>> payload =
|
||||
records.stream()
|
||||
.map(this::prepareMessageForRecord)
|
||||
.map(AbstractIntegrationMessageBuilder::build)
|
||||
.collect(Collectors.toList());
|
||||
|
||||
messageBuilder = getMessageBuilderFactory().withPayload(payload);
|
||||
}
|
||||
else if (KinesisMessageDrivenChannelAdapter.this.converter != null) {
|
||||
final List<String> partitionKeys = new ArrayList<>();
|
||||
final List<String> sequenceNumbers = new ArrayList<>();
|
||||
|
||||
List<Object> 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<String> partitionKeys;
|
||||
final List<String> 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);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user