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()`
This commit is contained in:
@@ -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<byte[], Object> 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<Object> messageBuilder = getMessageBuilderFactory()
|
||||
.withPayload(payload)
|
||||
.setHeader(AwsHeaders.STREAM, this.shardOffset.getStream())
|
||||
|
||||
@@ -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<Object, byte[]> 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);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user