Add records conversion support in batch mode

https://stackoverflow.com/questions/49730808/unable-to-consume-messages-as-batch-mode-in-kinesis-binder
This commit is contained in:
Artem Bilan
2018-04-12 18:11:57 -04:00
parent 6a999b2873
commit 2c58a6925e
4 changed files with 100 additions and 54 deletions

View File

@@ -241,7 +241,7 @@ The S3 Outbound Gateway is represented by the same `S3MessageHandler` with the `
The "request-reply" nature of this gateway is async and the `Transfer` result from the `TransferManager`
operation is sent to the `outputChannel`, assuming the transfer progress observation in the downstream flow.
The `S3ProgressListener can be supplied to track the transfer progress.
The `S3ProgressListener` can be supplied to track the transfer progress.
Also the listener can be populated into the returned `Transfer` afterwards in the downstream flow.
See more information in the `S3MessageHandler` JavaDocs and `<int-aws:s3-outbound-channel-adapter>` &
@@ -530,14 +530,22 @@ This channel adapter can be configured with the `DynamoDbMetaDataStore` mentione
By default this adapter uses `DeserializingConverter` to convert `byte[]` from the `Record` data.
Can be specified as `null` with meaning no conversion and the target `Message` is sent with the `byte[]` payload.
Additional headers like `AwsHeaders.RECEIVED_STREAM`, `AwsHeaders.RECEIVED_PARTITION_KEY` and `AwsHeaders.RECEIVED_SEQUENCE_NUMBER` are populated to the message for downstream logic.
Additional headers like `AwsHeaders.RECEIVED_STREAM`, `AwsHeaders.SHARD`, `AwsHeaders.RECEIVED_PARTITION_KEY` and `AwsHeaders.RECEIVED_SEQUENCE_NUMBER` are populated to the message for downstream logic.
When `CheckpointMode.manual` is used the `Checkpointer` instance is populated to the `AwsHeaders.CHECKPOINTER` header for acknowledgment in the downstream logic manually.
The `KinesisMessageDrivenChannelAdapter` ca be configured with the `ListenerMode` `record` or `batch` to process records one by one or send the whole just polled batch of records.
If `Converter` is configured to `null`, the entire `List<Record>` is sent as a payload.
Otherwise a list of converted `Record.getData().array()` is wrapped to the payload of message to send.
In this case the `AwsHeaders.RECEIVED_PARTITION_KEY` and `AwsHeaders.RECEIVED_SEQUENCE_NUMBER` headers contains values as a `List<String>` of partition keys and sequence numbers of converted records respectively.
The consumer group is included to the metadata store `key`.
When records are consumed, they are filtered by the last stored `lastCheckpoint` under the key as `[CONSUMER_GROUP]:[STREAM]:[SHARD_ID]`.
Starting with _version 2.0_, the `KinesisMessageDrivenChannelAdapter` can be configured with the `InboundMessageMapper` to extract message headers embedded into the record data (if any).
See `EmbeddedJsonHeadersMessageMapper` implementation for more information.
When `InboundMessageMapper` is used together with the `ListenerMode.batch`, each `Record` is converted to the `Message` with extracted embedded headers (if any) and converted `byte[]` payload if any and converter is present.
In this case `AwsHeaders.RECEIVED_PARTITION_KEY` and `AwsHeaders.RECEIVED_SEQUENCE_NUMBER` headers are populated to the particular message for a record.
These messages are wrapped as a list payload to one outbound message.
### Outbound Channel Adapter

View File

@@ -37,6 +37,7 @@ import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.core.AttributeAccessor;
@@ -941,43 +942,7 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i
switch (KinesisMessageDrivenChannelAdapter.this.listenerMode) {
case record:
for (Record record : records) {
Object payload = record.getData().array();
Message<?> messageToUse = null;
if (KinesisMessageDrivenChannelAdapter.this.embeddedHeadersMapper != null) {
try {
messageToUse =
KinesisMessageDrivenChannelAdapter.this.embeddedHeadersMapper
.toMessage((byte[]) payload);
payload = messageToUse.getPayload();
}
catch (Exception e) {
logger.warn("Could not parse embedded headers. Remain payload untouched.", e);
}
}
if (payload instanceof byte[] &&
KinesisMessageDrivenChannelAdapter.this.converter != null) {
payload = KinesisMessageDrivenChannelAdapter.this.converter.convert((byte[]) payload);
}
AbstractIntegrationMessageBuilder<Object> messageBuilder = getMessageBuilderFactory()
.withPayload(payload)
.setHeader(AwsHeaders.RECEIVED_STREAM, this.shardOffset.getStream())
.setHeader(AwsHeaders.SHARD, this.shardOffset.getShard())
.setHeader(AwsHeaders.RECEIVED_PARTITION_KEY, record.getPartitionKey())
.setHeader(AwsHeaders.RECEIVED_SEQUENCE_NUMBER, record.getSequenceNumber());
if (CheckpointMode.manual.equals(KinesisMessageDrivenChannelAdapter.this.checkpointMode)) {
messageBuilder.setHeader(AwsHeaders.CHECKPOINTER, this.checkpointer);
}
if (messageToUse != null) {
messageBuilder.copyHeadersIfAbsent(messageToUse.getHeaders());
}
performSend(messageBuilder, record);
performSend(prepareMessageForRecord(record), record);
if (CheckpointMode.record.equals(KinesisMessageDrivenChannelAdapter.this.checkpointMode)) {
this.checkpointer.checkpoint(record.getSequenceNumber());
@@ -987,14 +952,41 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i
break;
case batch:
AbstractIntegrationMessageBuilder<?> messageBuilder = getMessageBuilderFactory()
.withPayload(records)
.setHeader(AwsHeaders.RECEIVED_STREAM, this.shardOffset.getStream())
.setHeader(AwsHeaders.SHARD, this.shardOffset.getShard());
if (CheckpointMode.manual.equals(KinesisMessageDrivenChannelAdapter.this.checkpointMode)) {
messageBuilder.setHeader(AwsHeaders.CHECKPOINTER, this.checkpointer);
Object payload = records;
if (KinesisMessageDrivenChannelAdapter.this.embeddedHeadersMapper != null) {
payload = records.stream()
.map(this::prepareMessageForRecord)
.collect(Collectors.toList());
}
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);
break;
@@ -1003,7 +995,50 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i
}
}
private AbstractIntegrationMessageBuilder<Object> prepareMessageForRecord(Record record) {
Object payload = record.getData().array();
Message<?> messageToUse = null;
if (KinesisMessageDrivenChannelAdapter.this.embeddedHeadersMapper != null) {
try {
messageToUse =
KinesisMessageDrivenChannelAdapter.this.embeddedHeadersMapper
.toMessage((byte[]) payload);
payload = messageToUse.getPayload();
}
catch (Exception e) {
logger.warn("Could not parse embedded headers. Remain payload untouched.", e);
}
}
if (payload instanceof byte[] &&
KinesisMessageDrivenChannelAdapter.this.converter != null) {
payload = KinesisMessageDrivenChannelAdapter.this.converter.convert((byte[]) payload);
}
AbstractIntegrationMessageBuilder<Object> messageBuilder =
getMessageBuilderFactory()
.withPayload(payload)
.setHeader(AwsHeaders.RECEIVED_PARTITION_KEY, record.getPartitionKey())
.setHeader(AwsHeaders.RECEIVED_SEQUENCE_NUMBER, record.getSequenceNumber());
if (messageToUse != null) {
messageBuilder.copyHeadersIfAbsent(messageToUse.getHeaders());
}
return messageBuilder;
}
private void performSend(AbstractIntegrationMessageBuilder<?> messageBuilder, Object rawRecord) {
messageBuilder.setHeader(AwsHeaders.RECEIVED_STREAM, this.shardOffset.getStream())
.setHeader(AwsHeaders.SHARD, this.shardOffset.getShard());
if (CheckpointMode.manual.equals(KinesisMessageDrivenChannelAdapter.this.checkpointMode)) {
messageBuilder.setHeader(AwsHeaders.CHECKPOINTER, this.checkpointer);
}
Message<?> messageToSend = messageBuilder.build();
setAttributesIfNecessary(rawRecord, messageToSend);
try {

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2017 the original author or authors.
* Copyright 2017-2018 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -33,7 +33,8 @@ public enum ListenerMode {
record,
/**
* Each {@link Message} will contains {@code List<Record>} if not empty.
* Each {@link Message} will contain {@code List} ( if not empty)
* of converted or raw {@code Record}s.
*/
batch

View File

@@ -36,7 +36,6 @@ import org.springframework.beans.DirectFieldAccessor;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.serializer.support.DeserializingConverter;
import org.springframework.core.serializer.support.SerializingConverter;
import org.springframework.integration.aws.inbound.kinesis.CheckpointMode;
import org.springframework.integration.aws.inbound.kinesis.Checkpointer;
@@ -158,15 +157,18 @@ public class KinesisMessageDrivenChannelAdapterTests {
message = this.kinesisChannel.receive(10000);
assertThat(message).isNotNull();
assertThat(message.getPayload()).isInstanceOf(List.class);
List<Record> payload = (List<Record>) message.getPayload();
List<String> payload = (List<String>) message.getPayload();
assertThat(payload).size().isEqualTo(1);
Record record = payload.get(0);
assertThat(record.getPartitionKey()).isEqualTo("partition1");
assertThat(record.getSequenceNumber()).isEqualTo("2");
String record = payload.get(0);
assertThat(record).isEqualTo("bar");
DeserializingConverter deserializingConverter = new DeserializingConverter();
assertThat(deserializingConverter.convert(record.getData().array())).isEqualTo("bar");
Object partitionKeyHeader = message.getHeaders().get(AwsHeaders.RECEIVED_PARTITION_KEY);
assertThat(partitionKeyHeader).isInstanceOf(List.class);
assertThat((List<String>) partitionKeyHeader).contains("partition1");
Object sequenceNumberHeader = message.getHeaders().get(AwsHeaders.RECEIVED_SEQUENCE_NUMBER);
assertThat(sequenceNumberHeader).isInstanceOf(List.class);
assertThat((List<String>) sequenceNumberHeader).contains("2");
int n = 0;
while (n++ < 100) {