diff --git a/README.md b/README.md index ed9e1ec..8014c73 100644 --- a/README.md +++ b/README.md @@ -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 `` & @@ -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` 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` 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 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 dea1d2e..1afd242 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 @@ -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 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 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); break; @@ -1003,7 +995,50 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i } } + private AbstractIntegrationMessageBuilder 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 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 { diff --git a/src/main/java/org/springframework/integration/aws/inbound/kinesis/ListenerMode.java b/src/main/java/org/springframework/integration/aws/inbound/kinesis/ListenerMode.java index ebc20b6..b8d0a44 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/kinesis/ListenerMode.java +++ b/src/main/java/org/springframework/integration/aws/inbound/kinesis/ListenerMode.java @@ -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} if not empty. + * Each {@link Message} will contain {@code List} ( if not empty) + * of converted or raw {@code Record}s. */ batch diff --git a/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java b/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java index 09c0714..44ef72d 100644 --- a/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java +++ b/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java @@ -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 payload = (List) message.getPayload(); + List payload = (List) 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) partitionKeyHeader).contains("partition1"); + Object sequenceNumberHeader = message.getHeaders().get(AwsHeaders.RECEIVED_SEQUENCE_NUMBER); + assertThat(sequenceNumberHeader).isInstanceOf(List.class); + assertThat((List) sequenceNumberHeader).contains("2"); int n = 0; while (n++ < 100) {