diff --git a/README.md b/README.md index a3f09a0..1b5590e 100644 --- a/README.md +++ b/README.md @@ -398,7 +398,7 @@ Can be specified as `null` with meaning no conversion and the target `Message` i 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 an 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. +The `KinesisMessageDrivenChannelAdapter` can 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. @@ -433,6 +433,9 @@ For example, users may want to fully read any parent shards before starting to r } ``` +Starting with _version 3.0_. the `RequestShardForSequenceException` can be used for flow control to request the shard iterator for specific sequence. +For example, when consumer has failed processing batch at specific record, throwing this exception with a sequence of that record will ensure at-least-once delivery since the shard iterator will move back to the requested record sequence. + Also, the `KclMessageDrivenChannelAdapter` is provided for performing streams consumption by [Kinesis Client Library][]. See its JavaDocs for more information. diff --git a/build.gradle b/build.gradle index ab152ff..ea52b77 100644 --- a/build.gradle +++ b/build.gradle @@ -11,7 +11,7 @@ plugins { id 'checkstyle' id 'org.ajoberstar.grgit' version '4.1.1' id 'io.spring.dependency-management' version '1.1.0' - id 'com.jfrog.artifactory' version '4.31.7' + id 'com.jfrog.artifactory' version '4.31.9' } description = 'Spring Integration AWS Support' @@ -32,15 +32,15 @@ repositories { ext { assertjVersion = '3.24.2' awaitilityVersion = '4.2.0' - awsSdkVersion = '2.20.35' - jacksonVersion = '2.14.2' + awsSdkVersion = '2.20.51' + jacksonVersion = '2.15.0' junitVersion = '5.9.2' - log4jVersion = '2.19.0' + log4jVersion = '2.20.0' servletApiVersion = '6.0.0' springCloudAwsVersion = '3.0.0-RC2' - springIntegrationVersion = '6.0.4' + springIntegrationVersion = '6.0.5' kinesisClientVersion = '2.4.8' - kinesisProducerVersion = '0.15.5' + kinesisProducerVersion = '0.15.6' 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 0505cb8..ae62137 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 @@ -1091,46 +1091,25 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport if (!records.isEmpty()) { processRecords(records); } + this.shardIterator = result.nextShardIterator(); } } + catch (RequestShardForSequenceException requestShardForSequenceException) { + // Something wrong happened and not all records were processed. + // Must start from the provided sequence + KinesisShardOffset newOffset = new KinesisShardOffset(this.shardOffset); + newOffset.setSequenceNumber(requestShardForSequenceException.getSequenceNumber()); + newOffset.setIteratorType(ShardIteratorType.AT_SEQUENCE_NUMBER); + GetShardIteratorRequest shardIteratorRequest = newOffset.toShardIteratorRequest(); + this.shardIterator = + KinesisMessageDrivenChannelAdapter.this.amazonKinesis + .getShardIterator(shardIteratorRequest) + .join() + .shardIterator(); + } finally { attributesHolder.remove(); if (result != null) { - // If using manual checkpointer, we have to make sure we are allowed to use the next shard iterator - // Because if the manual checkpointer was not set to the latest record, it means there are records to be reprocessed - // and if we use the nextShardIterator, we will be skipping records that need to be reprocessed - List records = result.records(); - if (CheckpointMode.manual.equals(KinesisMessageDrivenChannelAdapter.this.checkpointMode) && - !records.isEmpty()) { - logger.info("Manual checkpointer. Must validate if should use getNextShardIterator()"); - String lastRecordSequence = records.get(records.size() - 1).sequenceNumber(); - String lastCheckpointSequence = this.checkpointer.getCheckpoint(); - if (lastCheckpointSequence.equals(lastRecordSequence)) { - logger.info("latestCheckpointSequence is same as latestRecordSequence. " + - "Should getNextShardIterator()"); - // Means the manual checkpointer has processed the last record, Should move forward - this.shardIterator = result.nextShardIterator(); - } - else { - logger.info("latestCheckpointSequence is not the same as latestRecordSequence. " + - "Should Get a new iterator AFTER_SEQUENCE_NUMBER latestCheckpointSequence"); - // Something wrong happened and not all records were processed. - // Must start from the latest known checkpoint - KinesisShardOffset newOffset = new KinesisShardOffset(this.shardOffset); - newOffset.setSequenceNumber(lastCheckpointSequence); - newOffset.setIteratorType(ShardIteratorType.AFTER_SEQUENCE_NUMBER); - GetShardIteratorRequest shardIteratorRequest = newOffset.toShardIteratorRequest(); - this.shardIterator = - KinesisMessageDrivenChannelAdapter.this.amazonKinesis - .getShardIterator(shardIteratorRequest) - .join() - .shardIterator(); - } - } - else { - this.shardIterator = result.nextShardIterator(); - } - if (this.shardIterator == null) { if (KinesisMessageDrivenChannelAdapter.this.lockRegistry != null) { KinesisMessageDrivenChannelAdapter.this.shardConsumerManager.shardOffsetsToConsumer @@ -1138,21 +1117,17 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport } // Shard is closed: nothing to consume anymore. // Checkpoint endingSequenceNumber to ensure shard is marked exhausted. - // If in CheckpointMode.manual, only checkpoint if lastCheckpointValue is also null, as this - // means that no records have ever been read and so the shard was empty - if (!CheckpointMode.manual.equals(KinesisMessageDrivenChannelAdapter.this.checkpointMode) - || this.checkpointer.getLastCheckpointValue() == null) { - for (Shard shard : readShardList(this.shardOffset.getStream())) { - if (shard.shardId().equals(this.shardOffset.getShard())) { - String endingSequenceNumber = - shard.sequenceNumberRange().endingSequenceNumber(); - if (endingSequenceNumber != null) { - checkpointSwallowingProvisioningExceptions(endingSequenceNumber); - } - break; + for (Shard shard : readShardList(this.shardOffset.getStream())) { + if (shard.shardId().equals(this.shardOffset.getShard())) { + String endingSequenceNumber = + shard.sequenceNumberRange().endingSequenceNumber(); + if (endingSequenceNumber != null) { + checkpointSwallowingProvisioningExceptions(endingSequenceNumber); } + break; } } + // Resharding is possible. if (KinesisMessageDrivenChannelAdapter.this.applicationEventPublisher != null) { KinesisMessageDrivenChannelAdapter.this.applicationEventPublisher.publishEvent( @@ -1185,7 +1160,7 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport } catch (ProvisionedThroughputExceededException ignored) { // This exception is ignored to guarantee that an exhausted shard is marked as CLOSED - // even in the case it's not possible to checkpoint. Otherwise the ShardConsumer is + // even in the case it's not possible to checkpoint. Otherwise, the ShardConsumer is // left in an illegal state where the shard iterator is null without any possibility // of recovering from it. logger.debug(ignored, "Exception while checkpointing empty shards"); @@ -1341,16 +1316,39 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport try { sendMessage(messageToSend); } - catch (Exception ex) { - logger.info(ex, () -> - "Got an exception during sending a '" - + messageToSend - + "'" - + "\nfor the '" - + rawRecord - + "'.\n" - + "Consider to use 'errorChannel' flow for the compensation logic."); + catch (RequestShardForSequenceException requestShardForSequenceException) { + // Rethrow + throw requestShardForSequenceException; } + catch (Exception ex) { + RequestShardForSequenceException requestShardForSequenceExceptionInCause = + findRequestShardForSequenceExceptionInCause(ex); + if (requestShardForSequenceExceptionInCause != null) { + throw requestShardForSequenceExceptionInCause; + } + else { + logger.info(ex, () -> + "Got an exception during sending a '" + + messageToSend + + "'" + + "\nfor the '" + + rawRecord + + "'.\n" + + "Consider to use 'errorChannel' flow for the compensation logic."); + } + } + } + + @Nullable + private static RequestShardForSequenceException findRequestShardForSequenceExceptionInCause(Throwable ex) { + if (ex instanceof RequestShardForSequenceException requestShardForSequenceException) { + return requestShardForSequenceException; + } + Throwable cause = ex.getCause(); + if (cause != null && cause != ex) { + return findRequestShardForSequenceExceptionInCause(cause); + } + return null; } private void checkpointIfBatchMode() { diff --git a/src/main/java/org/springframework/integration/aws/inbound/kinesis/RequestShardForSequenceException.java b/src/main/java/org/springframework/integration/aws/inbound/kinesis/RequestShardForSequenceException.java new file mode 100644 index 0000000..f815a2e --- /dev/null +++ b/src/main/java/org/springframework/integration/aws/inbound/kinesis/RequestShardForSequenceException.java @@ -0,0 +1,48 @@ +/* + * Copyright 2023 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.aws.inbound.kinesis; + +import software.amazon.awssdk.services.kinesis.model.ShardIteratorType; + +/** + * The flow control exception to notify the {@link KinesisMessageDrivenChannelAdapter} + * that specific shard iterator ({@link ShardIteratorType#AT_SEQUENCE_NUMBER}) + * must be requested instead of checkpointing. + * + * @author Artem Bilan + * + * @since 3.0 + */ +@SuppressWarnings("serial") +public class RequestShardForSequenceException extends RuntimeException { + + private final String sequenceNumber; + + public RequestShardForSequenceException(String sequenceNumber) { + this.sequenceNumber = sequenceNumber; + } + + public RequestShardForSequenceException(String sequenceNumber, Throwable cause) { + super(cause); + this.sequenceNumber = sequenceNumber; + } + + public String getSequenceNumber() { + return this.sequenceNumber; + } + +} 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 b255670..8c35794 100644 --- a/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java +++ b/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java @@ -148,8 +148,6 @@ public class KinesisMessageDrivenChannelAdapterTests { Checkpointer checkpointer = headers.get(AwsHeaders.CHECKPOINTER, Checkpointer.class); assertThat(checkpointer).isNotNull(); - checkpointer.checkpoint(); - message = this.kinesisChannel.receive(10000); assertThat(message).isNotNull(); assertThat(message.getPayload()).isEqualTo("bar"); @@ -161,6 +159,8 @@ public class KinesisMessageDrivenChannelAdapterTests { assertThat(this.kinesisChannel.receive(10)).isNull(); + checkpointer.checkpoint(); + assertThat(this.checkpointStore.get("SpringIntegration" + ":" + STREAM1 + ":" + "1")).isEqualTo("2"); this.kinesisMessageDrivenChannelAdapter.stop(); diff --git a/src/test/java/org/springframework/integration/aws/kinesis/KinesisIntegrationTests.java b/src/test/java/org/springframework/integration/aws/kinesis/KinesisIntegrationTests.java index 6dbde74..4b24e90 100644 --- a/src/test/java/org/springframework/integration/aws/kinesis/KinesisIntegrationTests.java +++ b/src/test/java/org/springframework/integration/aws/kinesis/KinesisIntegrationTests.java @@ -19,11 +19,13 @@ package org.springframework.integration.aws.kinesis; import java.util.Date; import java.util.HashSet; import java.util.Set; +import java.util.concurrent.atomic.AtomicBoolean; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; import software.amazon.awssdk.services.kinesis.KinesisAsyncClient; +import software.amazon.awssdk.services.kinesis.model.Record; import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.annotation.Autowired; @@ -34,6 +36,7 @@ import org.springframework.integration.annotation.ServiceActivator; import org.springframework.integration.aws.LocalstackContainerTest; import org.springframework.integration.aws.inbound.kinesis.KinesisMessageDrivenChannelAdapter; import org.springframework.integration.aws.inbound.kinesis.KinesisMessageHeaderErrorMessageStrategy; +import org.springframework.integration.aws.inbound.kinesis.RequestShardForSequenceException; import org.springframework.integration.aws.outbound.KinesisMessageHandler; import org.springframework.integration.aws.support.AwsHeaders; import org.springframework.integration.channel.QueueChannel; @@ -114,6 +117,14 @@ public class KinesisIntegrationTests implements LocalstackContainerTest { .contains("Channel 'kinesisReceiveChannel' expected one of the following data types " + "[class java.util.Date], but received [class java.lang.String]"); + String errorSequenceNumber = errorMessage.getHeaders().get(AwsHeaders.RAW_RECORD, Record.class).sequenceNumber(); + + // Second exception for the same record since we have requested via RequestShardForSequenceException + errorMessage = this.errorChannel.receive(30_000); + assertThat(errorMessage).isNotNull(); + assertThat(errorMessage.getHeaders().get(AwsHeaders.RAW_RECORD, Record.class).sequenceNumber()) + .isEqualTo(errorSequenceNumber); + for (int i = 0; i < 2; i++) { this.kinesisSendChannel .send(MessageBuilder.withPayload(new Date()).setHeader(AwsHeaders.STREAM, TEST_STREAM).build()); @@ -208,10 +219,14 @@ public class KinesisIntegrationTests implements LocalstackContainerTest { QueueChannel queueChannel = new QueueChannel(); queueChannel.addInterceptor(new ChannelInterceptor() { + private final AtomicBoolean thrown = new AtomicBoolean(); + @Override public void postSend(Message message, MessageChannel channel, boolean sent) { - if (message instanceof ErrorMessage) { - throw (RuntimeException) ((ErrorMessage) message).getPayload(); + if (message instanceof ErrorMessage errorMessage && this.thrown.compareAndSet(false, true)) { + throw new RequestShardForSequenceException( + errorMessage.getHeaders().get(AwsHeaders.RAW_RECORD, Record.class).sequenceNumber(), + errorMessage.getPayload()); } } diff --git a/src/test/resources/log4j2-test.xml b/src/test/resources/log4j2-test.xml index 2c7a637..5152896 100644 --- a/src/test/resources/log4j2-test.xml +++ b/src/test/resources/log4j2-test.xml @@ -9,7 +9,7 @@ - +