From bfde46ad41348aa7786ae06d7231a0df9d02382e Mon Sep 17 00:00:00 2001 From: abilan Date: Tue, 2 May 2023 17:00:05 -0400 Subject: [PATCH] GH-223: Fix errors in KinesisMDChA.getRecords() Fixes https://github.com/spring-projects/spring-integration-aws/issues/223 A `amazonKinesis.getRecords(getRecordsRequest).join()` throws `CompletionException` where the target reason is in a `cause` * Fix `KinesisMessageDrivenChannelAdapter.getRecords()` to `catch (CompletionException)` and then process its `cause` respectively * Upgrade to Spring Cloud AWS `3.0.0` GA * Upgrade other deps to the latest versions --- build.gradle | 6 +-- .../KinesisMessageDrivenChannelAdapter.java | 46 +++++++++++-------- ...nesisMessageDrivenChannelAdapterTests.java | 12 ++++- 3 files changed, 39 insertions(+), 25 deletions(-) diff --git a/build.gradle b/build.gradle index ea52b77..b67701e 100644 --- a/build.gradle +++ b/build.gradle @@ -32,15 +32,15 @@ repositories { ext { assertjVersion = '3.24.2' awaitilityVersion = '4.2.0' - awsSdkVersion = '2.20.51' + awsSdkVersion = '2.20.57' jacksonVersion = '2.15.0' junitVersion = '5.9.2' log4jVersion = '2.20.0' servletApiVersion = '6.0.0' - springCloudAwsVersion = '3.0.0-RC2' + springCloudAwsVersion = '3.0.0' springIntegrationVersion = '6.0.5' kinesisClientVersion = '2.4.8' - kinesisProducerVersion = '0.15.6' + kinesisProducerVersion = '0.15.7' testcontainersVersion = '1.18.0' idPrefix = 'aws' 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 9aba2d2..63ef44f 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 @@ -1195,26 +1195,32 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport try { return KinesisMessageDrivenChannelAdapter.this.amazonKinesis.getRecords(getRecordsRequest).join(); } - catch (ExpiredIteratorException e) { - // Iterator expired, but this does not mean that shard no longer contains - // records. - // Let's acquire iterator again (using checkpointer for iterator start - // sequence number). - logger.info(() -> - "Shard iterator for [" - + ShardConsumer.this - + "] expired.\n" - + "A new one will be started from the check pointed sequence number."); - this.state = ConsumerState.EXPIRED; - } - catch (ProvisionedThroughputExceededException ex) { - logger.warn(() -> - "GetRecords request throttled for [" - + ShardConsumer.this - + "] with the reason: " - + ex.getMessage()); - // We are throttled, so let's sleep - prepareSleepState(); + catch (CompletionException ex) { + Throwable cause = ex.getCause(); + if (cause instanceof ExpiredIteratorException) { + // Iterator expired, but this does not mean that shard no longer contains + // records. + // Let's acquire iterator again (using checkpointer for iterator start + // sequence number). + logger.info(() -> + "Shard iterator for [" + + ShardConsumer.this + + "] expired.\n" + + "A new one will be started from the check pointed sequence number."); + this.state = ConsumerState.EXPIRED; + } + else if (cause instanceof ProvisionedThroughputExceededException) { + logger.warn(() -> + "GetRecords request throttled for [" + + ShardConsumer.this + + "] with the reason: " + + cause.getMessage()); + // We are throttled, so let's sleep + prepareSleepState(); + } + else { + throw ex; + } } return null; 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 8c35794..a2b42ab 100644 --- a/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java +++ b/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java @@ -335,8 +335,16 @@ public class KinesisMessageDrivenChannelAdapterTests { .shardIterator(shard1Iterator1) .limit(25) .build())) - .willThrow(ProvisionedThroughputExceededException.builder().message("Iterator throttled").build()) - .willThrow(ExpiredIteratorException.builder().message("Iterator expired").build()); + .willReturn( + CompletableFuture.failedFuture( + ProvisionedThroughputExceededException.builder() + .message("Iterator throttled") + .build())) + .willReturn( + CompletableFuture.failedFuture( + ExpiredIteratorException.builder() + .message("Iterator expired") + .build())); SerializingConverter serializingConverter = new SerializingConverter();