From 078dca68c0b3e2312586c67a8b93a72af25f3c2d Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 10 Sep 2019 14:10:03 -0400 Subject: [PATCH] GH-152: Kinesis embeddedHeaders: Build Messages Fixes https://github.com/spring-projects/spring-integration-aws/issues/152 `map()` into `build()`` the result of a stream for the `embeddedHeadersMapper` in the `processMultipleRecords()`. Otherwise our `payload` is going to have `AbstractIntegrationMessageBuilder` instances. The idea is to allow to process a batch of Messages. --- .../KclMessageDrivenChannelAdapter.java | 5 +- .../KinesisMessageDrivenChannelAdapter.java | 110 +++++++++--------- 2 files changed, 61 insertions(+), 54 deletions(-) diff --git a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java index 46036eb..d2c2db7 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java +++ b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java @@ -377,7 +377,10 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport { Object payload = records; if (KclMessageDrivenChannelAdapter.this.embeddedHeadersMapper != null) { - payload = records.stream().map(this::prepareMessageForRecord).collect(Collectors.toList()); + payload = records.stream() + .map(this::prepareMessageForRecord) + .map(AbstractIntegrationMessageBuilder::build) + .collect(Collectors.toList()); } final List partitionKeys; 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 8064f1c..4dbd5f0 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 @@ -822,62 +822,63 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i if (this.task == null) { switch (this.state) { - case NEW: - case EXPIRED: - this.task = () -> { - try { - if (this.shardOffset.isReset()) { - this.checkpointer.remove(); - } - else { - String checkpoint = this.checkpointer.getCheckpoint(); - if (checkpoint != null) { - this.shardOffset.setSequenceNumber(checkpoint); - this.shardOffset.setIteratorType(ShardIteratorType.AFTER_SEQUENCE_NUMBER); + case NEW: + case EXPIRED: + this.task = () -> { + try { + if (this.shardOffset.isReset()) { + this.checkpointer.remove(); + } + else { + String checkpoint = this.checkpointer.getCheckpoint(); + if (checkpoint != null) { + this.shardOffset.setSequenceNumber(checkpoint); + this.shardOffset.setIteratorType(ShardIteratorType.AFTER_SEQUENCE_NUMBER); + } + } + if (logger.isInfoEnabled() && this.state == ConsumerState.NEW) { + logger.info("The [" + this + "] has been started."); + } + GetShardIteratorRequest shardIteratorRequest = + this.shardOffset.toShardIteratorRequest(); + this.shardIterator = KinesisMessageDrivenChannelAdapter.this.amazonKinesis + .getShardIterator(shardIteratorRequest).getShardIterator(); + if (ConsumerState.STOP != this.state) { + this.state = ConsumerState.CONSUME; } } - if (logger.isInfoEnabled() && this.state == ConsumerState.NEW) { - logger.info("The [" + this + "] has been started."); + finally { + this.task = null; } - GetShardIteratorRequest shardIteratorRequest = this.shardOffset.toShardIteratorRequest(); - this.shardIterator = KinesisMessageDrivenChannelAdapter.this.amazonKinesis - .getShardIterator(shardIteratorRequest).getShardIterator(); - if (ConsumerState.STOP != this.state) { - this.state = ConsumerState.CONSUME; + }; + break; + + case CONSUME: + this.task = this.processTask; + break; + + case SLEEP: + if (System.currentTimeMillis() >= this.sleepUntil) { + this.state = ConsumerState.CONSUME; + } + this.task = null; + break; + + case STOP: + if (this.shardIterator == null) { + if (logger.isInfoEnabled()) { + logger.info("Stopping the [" + this + "] on the checkpoint [" + + this.checkpointer.getCheckpoint() + + "] because the shard has been CLOSED and exhausted."); } } - finally { - this.task = null; + else { + if (logger.isInfoEnabled()) { + logger.info("Stopping the [" + this + "]."); + } } - }; - break; - - case CONSUME: - this.task = this.processTask; - break; - - case SLEEP: - if (System.currentTimeMillis() >= this.sleepUntil) { - this.state = ConsumerState.CONSUME; - } - this.task = null; - break; - - case STOP: - if (this.shardIterator == null) { - if (logger.isInfoEnabled()) { - logger.info("Stopping the [" + this + "] on the checkpoint [" - + this.checkpointer.getCheckpoint() - + "] because the shard has been CLOSED and exhausted."); - } - } - else { - if (logger.isInfoEnabled()) { - logger.info("Stopping the [" + this + "]."); - } - } - this.task = null; - break; + this.task = null; + break; } if (this.task != null) { @@ -997,7 +998,10 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i Object payload = records; if (KinesisMessageDrivenChannelAdapter.this.embeddedHeadersMapper != null) { - payload = records.stream().map(this::prepareMessageForRecord).collect(Collectors.toList()); + payload = records.stream() + .map(this::prepareMessageForRecord) + .map(AbstractIntegrationMessageBuilder::build) + .collect(Collectors.toList()); } final List partitionKeys; @@ -1152,7 +1156,7 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i throw new IllegalStateException("ConsumerInvoker thread [" + this + "] has been interrupted", e); } - for (Iterator iterator = this.consumers.iterator(); iterator.hasNext();) { + for (Iterator iterator = this.consumers.iterator(); iterator.hasNext(); ) { ShardConsumer shardConsumer = iterator.next(); if (ConsumerState.STOP == shardConsumer.state) { iterator.remove(); @@ -1267,7 +1271,7 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i } } finally { - for (Iterator iterator = this.locks.values().iterator(); iterator.hasNext();) { + for (Iterator iterator = this.locks.values().iterator(); iterator.hasNext(); ) { Lock lock = iterator.next(); try { lock.unlock();