From 4a0201f29d0d721d38133cf23efbc2e8e67457ba Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 31 May 2017 13:19:33 -0400 Subject: [PATCH] GH-71: Kinesis Adapters: Handle byte[] directly Fixes spring-projects/spring-integration-aws#71 It isn't correct to call `converter.convert()` in the `KinesisMessageHandler` if `payload` is already `ByteBuffer` or `byte[]` On the other hand an application might be interested in the `byte[]` payload on the consumer side. The `KinesisMessageDrivenChannelAdapter` must be able to produce `byte[]` `payload` without any conversion. * Add logic into `KinesisMessageHandler` to check the `payload` type before calling `converter.convert()` * Allow to configure `converter` to `null` for the `KinesisMessageDrivenChannelAdapter`. The body of the consumed record is presented in the `payload` as is - `byte[]` * Add JavaDoc to the `KinesisMessageHandler.setConverter()` --- .../KinesisMessageDrivenChannelAdapter.java | 13 +++++++--- .../aws/outbound/KinesisMessageHandler.java | 24 ++++++++++++++++++- 2 files changed, 33 insertions(+), 4 deletions(-) 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 5766468..f110482 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 @@ -179,8 +179,12 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i this.streamInitialSequence = streamInitialSequence; } + /** + * Specify a {@link Converter} to deserialize the {@code byte[]} from record's body. + * Can be {@code null} meaning no deserialization. + * @param converter the {@link Converter} to use or null + */ public void setConverter(Converter converter) { - Assert.notNull(converter, "'converter' must not be null"); this.converter = converter; } @@ -840,8 +844,11 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i switch (KinesisMessageDrivenChannelAdapter.this.listenerMode) { case record: for (Record record : records) { - Object payload = - KinesisMessageDrivenChannelAdapter.this.converter.convert(record.getData().array()); + Object payload = record.getData().array(); + + if (KinesisMessageDrivenChannelAdapter.this.converter != null) { + payload = KinesisMessageDrivenChannelAdapter.this.converter.convert((byte[]) payload); + } AbstractIntegrationMessageBuilder messageBuilder = getMessageBuilderFactory() .withPayload(payload) .setHeader(AwsHeaders.STREAM, this.shardOffset.getStream()) diff --git a/src/main/java/org/springframework/integration/aws/outbound/KinesisMessageHandler.java b/src/main/java/org/springframework/integration/aws/outbound/KinesisMessageHandler.java index 6255147..8bbc5f5 100644 --- a/src/main/java/org/springframework/integration/aws/outbound/KinesisMessageHandler.java +++ b/src/main/java/org/springframework/integration/aws/outbound/KinesisMessageHandler.java @@ -85,6 +85,11 @@ public class KinesisMessageHandler extends AbstractMessageHandler { this.asyncHandler = asyncHandler; } + /** + * Specify a {@link Converter} to serialize {@code payload} to the {@code byte[]} + * if that isn't {@code byte[]} already. + * @param converter the {@link Converter} to use; cannot be null. + */ public void setConverter(Converter converter) { Assert.notNull(converter, "'converter' must not be null."); this.converter = converter; @@ -209,12 +214,29 @@ public class KinesisMessageHandler extends AbstractMessageHandler { partitionKey = this.sequenceNumberExpression.getValue(this.evaluationContext, message, String.class); } + Object payload = message.getPayload(); + + ByteBuffer data; + + if (payload instanceof ByteBuffer) { + data = (ByteBuffer) payload; + } + else { + byte[] bytes = + payload instanceof byte[] + ? (byte[]) payload + : this.converter.convert(payload); + + data = ByteBuffer.wrap( + bytes); + } + return new PutRecordRequest() .withStreamName(stream) .withPartitionKey(partitionKey) .withExplicitHashKey(explicitHashKey) .withSequenceNumberForOrdering(sequenceNumber) - .withData(ByteBuffer.wrap(this.converter.convert(message.getPayload()))); + .withData(data); } }