Fix README.md for new state of code in project
This commit is contained in:
226
README.md
226
README.md
@@ -14,19 +14,17 @@ Note the Spring Integration AWS Extension is based on the [Spring Cloud AWS][] p
|
||||
|
||||
## Spring Integration's extensions to AWS
|
||||
|
||||
The current project version is `2.5.x` and it requires minimum Java `8` and Spring Integration `5.5.x`.
|
||||
Can be used with Spring Cloud `2020.0.x` Spring Cloud AWS `2.3.x`.
|
||||
The current project version is `3.0.x` and it requires minimum Java `17` and Spring Integration `6.1.x` & Spring Cloud AWS `3.0.x`.
|
||||
This version is also fully based on AWS Java SDK v2.
|
||||
Therefore, it has a lot of breaking changes since the previous version, for example an XML configuration support was fully removed.
|
||||
|
||||
This guide intends to explain briefly the various adapters available for [Amazon Web Services][] such as:
|
||||
|
||||
* **Amazon Simple Email Service (SES)**
|
||||
* **Amazon Simple Storage Service (S3)**
|
||||
* **Amazon Simple Queue Service (SQS)**
|
||||
* **Amazon Simple Notification Service (SNS)**
|
||||
* **Amazon DynamoDB**
|
||||
|
||||
Sample XML Namespace configurations for each adapter as well as sample code snippets are provided wherever necessary.
|
||||
Of the above libraries, *SES* and *SNS* provide outbound adapters only. All other services have inbound and outbound adapters. The *SQS* inbound adapter is capable of receiving notifications sent out from *SNS* where the topic is an *SQS* Queue.
|
||||
* **Amazon Kinesis**
|
||||
|
||||
## Contributing
|
||||
|
||||
@@ -38,13 +36,16 @@ Additionally, if you are contributing, we recommend following the process for Sp
|
||||
|
||||
These dependencies are optional in the project:
|
||||
|
||||
* `io.awspring.cloud:spring-cloud-aws-messaging` - for SQS and SNS channel adapters
|
||||
* `io.awspring.cloud:spring-cloud-aws-sns` - for SNS channel adapters
|
||||
* `io.awspring.cloud:spring-cloud-aws-sqs` - for SQS channel adapters
|
||||
* `io.awspring.cloud:spring-cloud-aws-s3` - for S3 channel adapters
|
||||
* `org.springframework.integration:spring-integration-file` - for S3 channel adapters
|
||||
* `org.springframework.integration:spring-integration-http` - for SNS inbound channel adapter
|
||||
* `com.amazonaws:aws-java-sdk-kinesis` - for Kinesis channel adapters
|
||||
* `com.amazonaws:amazon-kinesis-client` - for KCL-based inbound channel adapter
|
||||
* `software.amazon.awssdk:kinesis` - for Kinesis channel adapters
|
||||
* `software.amazon.kinesis:amazon-kinesis-client` - for KCL-based inbound channel adapter
|
||||
* `com.amazonaws:amazon-kinesis-producer` - for KPL-based `MessageHandler`
|
||||
* `com.amazonaws:aws-java-sdk-dynamodb` - for `DynamoDbMetadataStore` and `DynamoDbLockRegistry`
|
||||
* `software.amazon.awssdk:dynamodb` - for `DynamoDbMetadataStore` and `DynamoDbLockRegistry`
|
||||
* `software.amazon.awssdk:s3-transfer-manager` - for `S3MessageHandler`
|
||||
|
||||
Consider to include an appropriate dependency into your project when you use particular component from this project.
|
||||
|
||||
@@ -59,7 +60,7 @@ See their specification and JavaDocs for more information.
|
||||
|
||||
### Inbound Channel Adapter
|
||||
|
||||
The S3 Inbound Channel Adapter is represented by the `S3InboundFileSynchronizingMessageSource` (`<int-aws:s3-inbound-channel-adapter>`) and allows pulling S3 objects as files from the S3 bucket to the local directory for synchronization.
|
||||
The S3 Inbound Channel Adapter is represented by the `S3InboundFileSynchronizingMessageSource` and allows pulling S3 objects as files from the S3 bucket to the local directory for synchronization.
|
||||
This adapter is fully similar to the Inbound Channel Adapters in the FTP and SFTP Spring Integration modules.
|
||||
See more information in the [FTP/FTPS Adapters Chapter][] for common options or `SessionFactory`, `RemoteFileTemplate` and `FileListFilter` abstractions.
|
||||
|
||||
@@ -70,7 +71,7 @@ The Java Configuration is:
|
||||
public static class MyConfiguration {
|
||||
|
||||
@Autowired
|
||||
private AmazonS3 amazonS3;
|
||||
private S3Client amazonS3;
|
||||
|
||||
@Bean
|
||||
public S3InboundFileSynchronizer s3InboundFileSynchronizer() {
|
||||
@@ -104,26 +105,6 @@ public static class MyConfiguration {
|
||||
|
||||
With this config you receive messages with `java.io.File` `payload` from the `s3FilesChannel` after periodic synchronization of content from the Amazon S3 bucket into the local directory.
|
||||
|
||||
An XML variant may look like:
|
||||
|
||||
````xml
|
||||
<bean id="s3SessionFactory" class="org.springframework.integration.aws.support.S3SessionFactory"/>
|
||||
|
||||
<int-aws:s3-inbound-channel-adapter channel="s3Channel"
|
||||
session-factory="s3SessionFactory"
|
||||
auto-create-local-directory="true"
|
||||
delete-remote-files="true"
|
||||
preserve-timestamp="true"
|
||||
filename-pattern="*.txt"
|
||||
local-directory="."
|
||||
local-filename-generator-expression="#this.toUpperCase() + '.a' + @fooString"
|
||||
temporary-file-suffix=".foo"
|
||||
local-filter="acceptAllFilter"
|
||||
remote-directory-expression="'my_bucket'">
|
||||
<int:poller fixed-rate="1000"/>
|
||||
</int-aws:s3-inbound-channel-adapter>
|
||||
````
|
||||
|
||||
### Streaming Inbound Channel Adapter
|
||||
|
||||
This adapter produces message with payloads of type `InputStream`, allowing S3 objects to be fetched without writing to the local file system.
|
||||
@@ -144,7 +125,7 @@ public class S3JavaApplication {
|
||||
}
|
||||
|
||||
@Autowired
|
||||
private AmazonS3 amazonS3;
|
||||
private S3Client amazonS3;
|
||||
|
||||
@Bean
|
||||
@InboundChannelAdapter(value = "s3Channel", poller = @Poller(fixedDelay = "100"))
|
||||
@@ -174,28 +155,6 @@ public class S3JavaApplication {
|
||||
}
|
||||
````
|
||||
|
||||
An XML variant may look like:
|
||||
|
||||
````xml
|
||||
<bean id="metadataStore" class="org.springframework.integration.metadata.SimpleMetadataStore"/>
|
||||
|
||||
<bean id="acceptOnceFilter" class="org.springframework.integration.aws.support.filters.S3PersistentAcceptOnceFileListFilter">
|
||||
<constructor-arg index="0" ref="metadataStore"/>
|
||||
<constructor-arg index="1" value="streaming"/>
|
||||
</bean>
|
||||
|
||||
<bean id="s3SessionFactory" class="org.springframework.integration.aws.support.S3SessionFactory"/>
|
||||
|
||||
<int-aws:s3-inbound-streaming-channel-adapter channel="s3Channel"
|
||||
session-factory="s3SessionFactory"
|
||||
filter="acceptOnceFilter"
|
||||
remote-directory-expression="'my_bucket'">
|
||||
<int:poller fixed-rate="1000"/>
|
||||
</int-aws:s3-inbound-streaming-channel-adapter>
|
||||
````
|
||||
|
||||
Only one of `filename-pattern`, `filename-regex` or `filter` is allowed.
|
||||
|
||||
> NOTE: Unlike the non-streaming inbound channel adapter, this adapter does not prevent duplicates by default.
|
||||
> If you do not delete the remote file and wish to prevent the file being processed again, you can configure an `S3PersistentFileListFilter` in the `filter` attribute.
|
||||
> If you don’t actually want to persist the state, an in-memory `SimpleMetadataStore` can be used with the filter.
|
||||
@@ -203,7 +162,7 @@ Only one of `filename-pattern`, `filename-regex` or `filter` is allowed.
|
||||
|
||||
### Outbound Channel Adapter
|
||||
|
||||
The S3 Outbound Channel Adapter is represented by the `S3MessageHandler` (`<int-aws:s3-outbound-channel-adapter>` and `<int-aws:s3-outbound-gateway>`) and allows performing `upload`, `download` and `copy` (see `S3MessageHandler.Command` enum) operations in the provided S3 bucket.
|
||||
The S3 Outbound Channel Adapter is represented by the `S3MessageHandler` and allows performing `upload`, `download` and `copy` (see `S3MessageHandler.Command` enum) operations in the provided S3 bucket.
|
||||
|
||||
The Java Configuration is:
|
||||
|
||||
@@ -212,12 +171,12 @@ The Java Configuration is:
|
||||
public static class MyConfiguration {
|
||||
|
||||
@Autowired
|
||||
private AmazonS3 amazonS3;
|
||||
private S3AsyncClient amazonS3;
|
||||
|
||||
@Bean
|
||||
@ServiceActivator(inputChannel = "s3UploadChannel")
|
||||
public MessageHandler s3MessageHandler() {
|
||||
return new S3MessageHandler(amazonS3(), "myBuck");
|
||||
return new S3MessageHandler(amazonS3(), "my-bucket");
|
||||
}
|
||||
|
||||
}
|
||||
@@ -225,34 +184,21 @@ public static class MyConfiguration {
|
||||
|
||||
With this config you can send a message with the `java.io.File` as `payload` and the `transferManager.upload()` operation will be performed, where the file name is used as a S3 Object key.
|
||||
|
||||
An XML variant may look like:
|
||||
|
||||
````xml
|
||||
<bean id="transferManager" class="com.amazonaws.services.s3.transfer.TransferManager"/>
|
||||
|
||||
<int-aws:s3-outbound-channel-adapter transfer-manager="transferManager"
|
||||
channel="s3SendChannel"
|
||||
bucket="foo"
|
||||
command="DOWNLOAD"
|
||||
key="myDirectory"/>
|
||||
````
|
||||
|
||||
See more information in the `S3MessageHandler` JavaDocs and `<int-aws:s3-outbound-channel-adapter>` & `<int-aws:s3-outbound-gateway>` descriptions.
|
||||
See more information in the `S3MessageHandler` JavaDocs.
|
||||
|
||||
### Outbound Gateway
|
||||
|
||||
The S3 Outbound Gateway is represented by the same `S3MessageHandler` with the `produceReply = true` constructor argument for Java Configuration and `<int-aws:s3-outbound-gateway>` for xml definitions.
|
||||
The S3 Outbound Gateway is represented by the same `S3MessageHandler` with the `produceReply = true` constructor argument for Java Configuration.
|
||||
|
||||
The "request-reply" nature of this gateway is async and the `Transfer` result from the `TransferManager` operation is sent to the `outputChannel`, assuming the transfer progress observation in the downstream flow.
|
||||
|
||||
The `S3ProgressListener` can be supplied to track the transfer progress.
|
||||
Also the listener can be populated into the returned `Transfer` afterwards in the downstream flow.
|
||||
The `TransferListener` can be supplied via `AwsHeaders.TRANSFER_LISTENER` header of the request message to track the transfer progress.
|
||||
|
||||
See more information in the `S3MessageHandler` JavaDocs and `<int-aws:s3-outbound-channel-adapter>` & `<int-aws:s3-outbound-gateway>` descriptions.
|
||||
See more information in the `S3MessageHandler` JavaDocs.
|
||||
|
||||
## Simple Email Service (SES)
|
||||
|
||||
There is no adapter for SES, since [Spring Cloud AWS][] provides implementations for `org.springframework.mail.MailSender` - `SimpleEmailServiceMailSender` and `SimpleEmailServiceJavaMailSender`, which can be injected to the `<int-mail:outbound-channel-adapter>`.
|
||||
There is no adapter for SES, since [Spring Cloud AWS][] provides implementations for `org.springframework.mail.MailSender` - `SimpleEmailServiceMailSender` and `SimpleEmailServiceJavaMailSender`, which can be injected to the `MailSendingMessageHandler`.
|
||||
|
||||
## Amazon Simple Queue Service (SQS)
|
||||
|
||||
@@ -260,7 +206,7 @@ The `SQS` adapters are fully based on the [Spring Cloud AWS][] foundation, so fo
|
||||
|
||||
### Outbound Channel Adapter
|
||||
|
||||
The SQS Outbound Channel Adapter is presented by the `SqsMessageHandler` implementation (`<int-aws:sqs-outbound-channel-adapter>`) and allows sending messages to the SQS `queue` with provided `AmazonSQS` client.
|
||||
The SQS Outbound Channel Adapter is presented by the `SqsMessageHandler` implementation and allows sending messages to the SQS `queue` with provided `SqsAsyncClient` client.
|
||||
An SQS queue can be configured explicitly on the adapter (using `org.springframework.integration.expression.ValueExpression`) or as a SpEL `Expression`, which is evaluated against request message as a root object of evaluation context.
|
||||
In addition, the `queue` can be extracted from the message headers under `AwsHeaders.QUEUE`.
|
||||
|
||||
@@ -272,31 +218,20 @@ public static class MyConfiguration {
|
||||
|
||||
@Bean
|
||||
@ServiceActivator(inputChannel = "sqsSendChannel")
|
||||
public MessageHandler sqsMessageHandler() {
|
||||
return new SqsMessageHandler(AmazonSQSAsync amazonSqs);
|
||||
public MessageHandler sqsMessageHandler(SqsAsyncClient amazonSqs) {
|
||||
return new SqsMessageHandler(amazonSqs);
|
||||
}
|
||||
|
||||
}
|
||||
````
|
||||
|
||||
An XML variant may look like:
|
||||
|
||||
````xml
|
||||
<aws-messaging:sqs-async-client id="sqs"/>
|
||||
|
||||
<int-aws:sqs-outbound-channel-adapter sqs="sqs"
|
||||
channel="sqsSendChannel"
|
||||
queue="foo"/>
|
||||
````
|
||||
|
||||
Starting with _version 2.0_, the `SqsMessageHandler` can be configured with the `HeaderMapper` to map message headers to the SQS message attributes.
|
||||
See `SqsHeaderMapper` implementation for more information and also consult with [Amazon SQS Message Attributes][] about value types and restrictions.
|
||||
|
||||
### Inbound Channel Adapter
|
||||
|
||||
The SQS Inbound Channel Adapter is a `message-driven` implementation for the `MessageProducer` and is represented with `SqsMessageDrivenChannelAdapter`.
|
||||
This channel adapter is based on the `io.awspring.cloud.messaging.listener.SimpleMessageListenerContainer` to receive messages from the provided `queues` in async manner and send an enhanced Spring Integration Message to the provided `MessageChannel`.
|
||||
The enhancements include `AwsHeaders.MESSAGE_ID`, `AwsHeaders.RECEIPT_HANDLE` and `AwsHeaders.RECEIVED_QUEUE` message headers.
|
||||
This channel adapter is based on the `io.awspring.cloud.sqs.listener.SqsMessageListenerContainer` to receive messages from the provided `queues` in async manner and send an enhanced Spring Integration Message to the provided `MessageChannel`.
|
||||
|
||||
The Java Configuration is pretty simple:
|
||||
|
||||
@@ -305,7 +240,7 @@ The Java Configuration is pretty simple:
|
||||
public static class MyConfiguration {
|
||||
|
||||
@Autowired
|
||||
private AmazonSQSAsync amazonSqs;
|
||||
private SqsAsyncClient amazonSqs;
|
||||
|
||||
@Bean
|
||||
public PollableChannel inputChannel() {
|
||||
@@ -321,53 +256,25 @@ public static class MyConfiguration {
|
||||
}
|
||||
````
|
||||
|
||||
An XML variant may look like:
|
||||
|
||||
````xml
|
||||
<aws-messaging:sqs-async-client id="sqs"/>
|
||||
|
||||
<int-aws:sqs-message-driven-channel-adapter id="sqsInboundChannel"
|
||||
sqs="sqs"
|
||||
error-channel="myErrorChannel"
|
||||
queues="foo, bar"
|
||||
delete-message-on-exception="false"
|
||||
max-number-of-messages="5"
|
||||
visibility-timeout="200"
|
||||
wait-time-out="40"
|
||||
send-timeout="2000"/>
|
||||
````
|
||||
|
||||
The `SqsMessageDrivenChannelAdapter` exposes all `SimpleMessageListenerContainer` attributes to configure, and an important one of them is `messageDeletionPolicy`, which is set to `NO_REDRIVE` by default.
|
||||
|
||||
Possible values are:
|
||||
|
||||
- `ALWAYS` - Always deletes a message automatically.
|
||||
- `NEVER` - Never deletes a message automatically.
|
||||
- `NO_REDRIVE` - Deletes a message if no redrive policy is defined.
|
||||
- `ON_SUCCESS` - Deletes a message when successfully executed by the listener method (no exception is thrown).
|
||||
|
||||
|
||||
Having that to `NEVER`, it is a responsibility of end-application to delete message. For this purpose a `AwsHeaders.RECEIPT_HANDLE` message header must be used for the message deletion:
|
||||
|
||||
````java
|
||||
MessageHeaders headers = message.getHeaders();
|
||||
this.amazonSqs.deleteMessageAsync(
|
||||
new DeleteMessageRequest(headers.get(AwsHeaders.QUEUE), headers.get(AwsHeaders.RECEIPT_HANDLE)));
|
||||
````
|
||||
The target listener container can be configured via `SqsMessageDrivenChannelAdapter.setSqsContainerOptions(SqsContainerOptions)` option.
|
||||
|
||||
## Amazon Simple Notification Service (SNS)
|
||||
|
||||
Amazon SNS is a publish-subscribe messaging system that allows clients to publish notification to a particular topic.
|
||||
Other interested clients may subscribe using different protocols like HTTP/HTTPS, e-mail, or an Amazon SQS queue to receive the messages. Plus mobile devices can be registered as subscribers from the AWS Management Console.
|
||||
Other interested clients may subscribe using different protocols like HTTP/HTTPS, e-mail, or an Amazon SQS queue to receive the messages.
|
||||
Plus mobile devices can be registered as subscribers from the AWS Management Console.
|
||||
|
||||
Unfortunately [Spring Cloud AWS][] doesn't provide flexible components which can be used from the channel adapter implementations, but Amazon SNS API is pretty simple, on the other hand.
|
||||
Hence, Spring Integration AWS SNS Support is straightforward and just allows to provide channel adapter foundation for Spring Integration applications.
|
||||
|
||||
Since e-mail, SMS and mobile device subscription/unsubscription confirmation is out of the Spring Integration application scope and can be done only from the AWS Management Console, we provide only HTTP/HTTPS SNS endpoint in face of `SnsInboundChannelAdapter`. The SQS-to-SNS subscription can be done with the simple usage of `com.amazonaws.services.sns.util.Topics#subscribeQueue()`, which confirms subscription automatically.
|
||||
Since e-mail, SMS and mobile device subscription/unsubscription confirmation is out of the Spring Integration application scope and can be done only from the AWS Management Console, we provide only HTTP/HTTPS SNS endpoint in face of `SnsInboundChannelAdapter`.
|
||||
The SQS-to-SNS subscription can also be done through account configuration: https://docs.aws.amazon.com/sns/latest/dg/subscribe-sqs-queue-to-sns-topic.html.
|
||||
|
||||
### Inbound Channel Adapter
|
||||
|
||||
The `SnsInboundChannelAdapter` (`<int-aws:sns-inbound-channel-adapter>`) is an extension of `HttpRequestHandlingMessagingGateway` and must be as a part of Spring MVC application. Its URL must be used from the AWS Management Console to add this endpoint as a subscriber to the SNS Topic. However before receiving any notification itself this HTTP endpoint must confirm the subscription.
|
||||
The `SnsInboundChannelAdapter` is an extension of `HttpRequestHandlingMessagingGateway` and must be as a part of Spring MVC application.
|
||||
Its URL must be used from the AWS Management Console to add this endpoint as a subscriber to the SNS Topic.
|
||||
However. before receiving any notification itself this HTTP endpoint must confirm the subscription.
|
||||
|
||||
See `SnsInboundChannelAdapter` JavaDocs for more information.
|
||||
|
||||
@@ -385,7 +292,7 @@ The Java Configuration is pretty simple:
|
||||
public static class MyConfiguration {
|
||||
|
||||
@Autowired
|
||||
private AmazonSNS amazonSns;
|
||||
private SnsClient amazonSns;
|
||||
|
||||
@Bean
|
||||
public PollableChannel inputChannel() {
|
||||
@@ -402,30 +309,19 @@ public static class MyConfiguration {
|
||||
}
|
||||
````
|
||||
|
||||
An XML variant may look like:
|
||||
|
||||
````xml
|
||||
<int-aws:sns-inbound-channel-adapter sns="amazonSns"
|
||||
path="/foo"
|
||||
channel="snsChannel"
|
||||
error-channel="errorChannel"
|
||||
handle-notification-status="true"
|
||||
payload-expression="payload.Message"/>
|
||||
````
|
||||
|
||||
Note: by default the message `payload` is a `Map` converted from the received Topic JSON message.
|
||||
For the convenience a `payload-expression` is provided with the `Message` as a root object of the evaluation context.
|
||||
Hence, even some HTTP headers, populated by the `DefaultHttpHeaderMapper`, are available for the evaluation context.
|
||||
|
||||
### Outbound Channel Adapter
|
||||
|
||||
The `SnsMessageHandler` (`<int-aws:sns-outbound-channel-adapter>`) is a simple one-way Outbound Channel Adapter to send Topic Notification using `AmazonSNS` service.
|
||||
The `SnsMessageHandler` is a simple one-way Outbound Channel Adapter to send Topic Notification using `SnsAsyncClient` service.
|
||||
|
||||
This Channel Adapter (`MessageHandler`) accepts these options:
|
||||
|
||||
- `topic-arn` (`topic-arn-expression`) - the SNS Topic to send notification for.
|
||||
- `subject` (`subject-expression`) - the SNS Notification Subject;
|
||||
- `body-expression` - the SpEL expression to evaluate the `message` property for the `com.amazonaws.services.sns.model.PublishRequest`.
|
||||
- `body-expression` - the SpEL expression to evaluate the `message` property for the `software.amazon.awssdk.services.sns.model.PublishRequest`.
|
||||
- `resource-id-resolver` - a `ResourceIdResolver` bean reference to resolve logical topic names to physical resource ids;
|
||||
|
||||
See `SnsMessageHandler` JavaDocs for more information.
|
||||
@@ -450,21 +346,7 @@ public MessageHandler snsMessageHandler() {
|
||||
|
||||
NOTE: the `bodyExpression` can be evaluated to a `org.springframework.integration.aws.support.SnsBodyBuilder` allowing the configuration of a `json` `messageStructure` for the `PublishRequest` and provide separate messages for different protocols.
|
||||
The same `SnsBodyBuilder` rule is applied for the raw `payload` if the `bodyExpression` hasn't been configured.
|
||||
NOTE: if the `payload` of `requestMessage` is a `com.amazonaws.services.sns.model.PublishRequest` already, the `SnsMessageHandler` doesn't do anything with it and it is sent as-is.
|
||||
|
||||
The XML variant may look like:
|
||||
|
||||
````xml
|
||||
<int-aws:sns-outbound-channel-adapter
|
||||
id="snsAdapter"
|
||||
sns="amazonSns"
|
||||
channel="notificationChannel"
|
||||
topic-arn="foo"
|
||||
subject="bar"
|
||||
message-group-id="foo-messages"
|
||||
message-deduplication-id-expression="headers.id"
|
||||
body-expression="payload.toUpperCase()"/>
|
||||
````
|
||||
NOTE: if the `payload` of `requestMessage` is a `software.amazon.awssdk.services.sns.model.PublishRequest` already, the `SnsMessageHandler` doesn't do anything with it, and it is sent as-is.
|
||||
|
||||
Starting with _version 2.0_, the `SnsMessageHandler` can be configured with the `HeaderMapper` to map message headers to the SNS message attributes.
|
||||
See `SnsHeaderMapper` implementation for more information and also consult with [Amazon SNS Message Attributes][] about value types and restrictions.
|
||||
@@ -475,12 +357,12 @@ and `messageDeduplicationIdExpression` properties.
|
||||
## Metadata Store for Amazon DynamoDB
|
||||
|
||||
The `DynamoDbMetadataStore`, a `ConcurrentMetadataStore` implementation, is provided to keep the metadata for Spring Integration components in the distributed Amazon DynamoDB store.
|
||||
The implementation is based on a simple table with `KEY` and `VALUE` attributes, both are string types and the `KEY` is primary key of the table.
|
||||
The implementation is based on a simple table with `metadataKey` and `metadataValue` attributes, both are string types and the `metadataKey` is partition key of the table.
|
||||
By default, the `SpringIntegrationMetadataStore` table is used, and it is created during `DynamoDbMetaDataStore` initialization if that doesn't exist yet.
|
||||
The `DynamoDbMetadataStore` can be used for the `KinesisMessageDrivenChannelAdapter` as a cloud-based `cehckpointStore`.
|
||||
|
||||
Starting with _version 2.0_, the `DynamoDbMetadataStore` can be configured with the `timeToLive` option to enable the [DynamoDB TTL][] feature.
|
||||
The `TTL` attribute is added to each item with the value based on the sum of current time and provided `timeToLive` in seconds.
|
||||
The `expireAt` attribute is added to each item with the value based on the sum of current time and provided `timeToLive` in seconds.
|
||||
If the provided `timeToLive` value is non-positive, the TTL functionality is disabled on the table.
|
||||
|
||||
## Amazon Kinesis
|
||||
@@ -500,7 +382,7 @@ The Java Configuration is pretty simple:
|
||||
public static class MyConfiguration {
|
||||
|
||||
@Bean
|
||||
public KinesisMessageDrivenChannelAdapter kinesisInboundChannelChannel(AmazonKinesis amazonKinesis) {
|
||||
public KinesisMessageDrivenChannelAdapter kinesisInboundChannelChannel(KinesisAsyncClient amazonKinesis) {
|
||||
KinesisMessageDrivenChannelAdapter adapter =
|
||||
new KinesisMessageDrivenChannelAdapter(amazonKinesis, "MY_STREAM");
|
||||
adapter.setOutputChannel(kinesisReceiveChannel());
|
||||
@@ -558,22 +440,19 @@ See its JavaDocs for more information.
|
||||
|
||||
The `KinesisMessageHandler` is an `AbstractMessageHandler` to perform put record to the Kinesis stream.
|
||||
The stream, partition key (or explicit hash key) and sequence number can be determined against request message via evaluation provided expressions or can be specified statically.
|
||||
They also can specified as `AwsHeaders.STREAM`, `AwsHeaders.PARTITION_KEY` and `AwsHeaders.SEQUENCE_NUMBER` respectively.
|
||||
They also can be specified as `AwsHeaders.STREAM`, `AwsHeaders.PARTITION_KEY` and `AwsHeaders.SEQUENCE_NUMBER` respectively.
|
||||
|
||||
The `KinesisMessageHandler` can be configured with the `outputChannel` for sending a `Message` on successful put operation.
|
||||
The payload is the original request and additional `AwsHeaders.SHARD` and `AwsHeaders.SEQUENCE_NUMBER` headers are populated from the `PutRecordResult`.
|
||||
If the request payload is a `PutRecordsRequest`, the full `PutRecordsResult` is populated in the `AwsHeaders.SERVICE_RESULT` header instead.
|
||||
|
||||
When an async failure is happened on the put operation, the `ErrorMessage` is send to the `failureChannel`.
|
||||
When an async failure is happened on the put operation, the `ErrorMessage` is sent to the `errorChannel` header or global one.
|
||||
The payload is an `AwsRequestFailureException`.
|
||||
|
||||
An `com.amazonaws.handlers.AsyncHandler` can also be provided to the `KinesisMessageHandler` for custom handling after putting record(s) to the stream.
|
||||
This is called independently if `outputChannel` and/or `failureChannel` are provided.
|
||||
|
||||
The `payload` of request message can be:
|
||||
|
||||
- `PutRecordsRequest` to perform `AmazonKinesisAsync.putRecordsAsync`
|
||||
- `PutRecordRequest` to perform `AmazonKinesisAsync.putRecordAsync`
|
||||
- `PutRecordsRequest` to perform `KinesisAsyncClient.putRecords`
|
||||
- `PutRecordRequest` to perform `KinesisAsyncClient.putRecord`
|
||||
- `ByteBuffer` to represent a data of the `PutRecordRequest`
|
||||
- `byte[]` which is wrapped to the `ByteBuffer`
|
||||
- any other type which is converted to the `byte[]` by the provided `Converter`; the `SerializingConverter` is used by default.
|
||||
@@ -583,13 +462,11 @@ The Java Configuration for the message handler:
|
||||
````java
|
||||
@Bean
|
||||
@ServiceActivator(inputChannel = "kinesisSendChannel")
|
||||
public MessageHandler kinesisMessageHandler(AmazonKinesis amazonKinesis,
|
||||
MessageChannel channel,
|
||||
MessageChannel errorChannel) {
|
||||
public MessageHandler kinesisMessageHandler(KinesisAsyncClient amazonKinesis,
|
||||
MessageChannel channel) {
|
||||
KinesisMessageHandler kinesisMessageHandler = new KinesisMessageHandler(amazonKinesis);
|
||||
kinesisMessageHandler.setPartitionKey("1");
|
||||
kinesisMessageHandler.setOutputChannel(channel);
|
||||
kinesisMessageHandler.setFailureChannel(errorChannel);
|
||||
return kinesisMessageHandler;
|
||||
}
|
||||
````
|
||||
@@ -607,6 +484,11 @@ The `DefaultLockRegistry` performs this function within a single component; you
|
||||
When used with a shared `MessageGroupStore`, the `DynamoDbLockRegistry` can be used to provide this functionality across multiple application instances, such that only one instance can manipulate the group at a time.
|
||||
This implementation can also be used for the distributed leader elections using a [LockRegistryLeaderInitiator][].
|
||||
|
||||
## Testing
|
||||
|
||||
The tests in the project are performed via Testcontainers and [Local Stack][] image.
|
||||
See `LocalstackContainerTest` interface Javadocs for more information.
|
||||
|
||||
[Spring Cloud AWS]: https://awspring.io/
|
||||
[AWS SDK for Java]: https://aws.amazon.com/sdkforjava/
|
||||
[Amazon Web Services]: https://aws.amazon.com/
|
||||
@@ -617,10 +499,8 @@ This implementation can also be used for the distributed leader elections using
|
||||
[Pull requests]: https://help.github.com/en/articles/creating-a-pull-request
|
||||
[contributor guidelines]: https://github.com/spring-projects/spring-integration/blob/main/CONTRIBUTING.adoc
|
||||
[administrator guidelines]: https://github.com/spring-projects/spring-integration/wiki/Administrator-Guidelines
|
||||
[Dynalite]: https://github.com/mhart/dynalite
|
||||
[DynamoDB TTL]: https://docs.aws.amazon.com/amazondynamodb/latest/developerguide/TTL.html
|
||||
[Kinesis Client Library]: https://docs.aws.amazon.com/streams/latest/dev/developing-consumers-with-kcl.html
|
||||
[Kinesalite]: https://github.com/mhart/kinesalite
|
||||
[Amazon SQS Message Attributes]: https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/sqs-message-attributes.html
|
||||
[Amazon SNS Message Attributes]: https://docs.aws.amazon.com/sns/latest/dg/SNSMessageAttributes.html
|
||||
[Leader Election]: https://docs.spring.io/spring-integration/docs/current/reference/html/messaging-endpoints-chapter.html#leadership-event-handling
|
||||
|
||||
@@ -82,8 +82,7 @@ public class DynamoDbMetadataStore implements ConcurrentMetadataStore, Initializ
|
||||
*/
|
||||
public static final String TTL = "expireAt";
|
||||
|
||||
private static final String KEY_NOT_EXISTS_EXPRESSION =
|
||||
String.format("attribute_not_exists(%s)", KEY);
|
||||
private static final String KEY_NOT_EXISTS_EXPRESSION = String.format("attribute_not_exists(%s)", KEY);
|
||||
|
||||
private final DynamoDbAsyncClient dynamoDB;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user