Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 23 additions & 0 deletions docs/src/main/asciidoc/_configprops.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,28 @@
|spring.cloud.aws.dynamodb.table-suffix | | The suffix used to resolve table names.
|spring.cloud.aws.endpoint | | Overrides the default endpoint for all auto-configured AWS clients.
|spring.cloud.aws.fips-enabled | | Configure whether the SDK should use the AWS fips endpoints.
|spring.cloud.aws.kinesis.dualstack-enabled | | Configure whether the AWS client should use the AWS dualstack endpoint. Note that not each AWS service supports dual-stack. For complete list check <a href="https://docs.aws.amazon.com/vpc/latest/userguide/aws-ipv6-support.html">AWS services that support IPv6</a>
|spring.cloud.aws.kinesis.enabled | `+++true+++` | Enables Kinesis integration.
|spring.cloud.aws.kinesis.endpoint | | Overrides the default endpoint.
|spring.cloud.aws.kinesis.listener.auto-startup | |
|spring.cloud.aws.kinesis.listener.billing-mode | |
|spring.cloud.aws.kinesis.listener.checkpoint-interval | |
|spring.cloud.aws.kinesis.listener.checkpoint-mode | |
|spring.cloud.aws.kinesis.listener.checkpoint-record-count | |
|spring.cloud.aws.kinesis.listener.consumer-arn | |
|spring.cloud.aws.kinesis.listener.consumer-name | |
|spring.cloud.aws.kinesis.listener.content-type | |
|spring.cloud.aws.kinesis.listener.graceful-shutdown-timeout | |
|spring.cloud.aws.kinesis.listener.idle-time-between-reads | |
|spring.cloud.aws.kinesis.listener.initial-position | |
|spring.cloud.aws.kinesis.listener.initial-position-timestamp | |
|spring.cloud.aws.kinesis.listener.lease-table-name | |
|spring.cloud.aws.kinesis.listener.max-records | |
|spring.cloud.aws.kinesis.listener.metrics-level | |
|spring.cloud.aws.kinesis.listener.metrics-namespace | |
|spring.cloud.aws.kinesis.listener.phase | |
|spring.cloud.aws.kinesis.listener.retrieval-mode | |
|spring.cloud.aws.kinesis.region | | Overrides the default region.
|spring.cloud.aws.parameterstore.dualstack-enabled | | Configure whether the AWS client should use the AWS dualstack endpoint. Note that not each AWS service supports dual-stack. For complete list check <a href="https://docs.aws.amazon.com/vpc/latest/userguide/aws-ipv6-support.html">AWS services that support IPv6</a>
|spring.cloud.aws.parameterstore.enabled | `+++true+++` | Enables ParameterStore integration.
|spring.cloud.aws.parameterstore.endpoint | | Overrides the default endpoint.
Expand Down Expand Up @@ -106,6 +128,7 @@
|spring.cloud.aws.sqs.dualstack-enabled | | Configure whether the AWS client should use the AWS dualstack endpoint. Note that not each AWS service supports dual-stack. For complete list check <a href="https://docs.aws.amazon.com/vpc/latest/userguide/aws-ipv6-support.html">AWS services that support IPv6</a>
|spring.cloud.aws.sqs.enabled | `+++true+++` | Enables SQS integration.
|spring.cloud.aws.sqs.endpoint | | Overrides the default endpoint.
|spring.cloud.aws.sqs.extended | |
|spring.cloud.aws.sqs.listener.auto-startup | | Defines whether SQS listeners will start automatically or not.
|spring.cloud.aws.sqs.listener.max-concurrent-messages | | The maximum concurrent messages that can be processed simultaneously for each queue. Note that if acknowledgement batching is being used, the actual maximum number of messages inflight might be higher.
|spring.cloud.aws.sqs.listener.max-delay-between-polls | | The maximum amount of time to wait between consecutive polls to SQS.
Expand Down
284 changes: 282 additions & 2 deletions docs/src/main/asciidoc/kinesis.adoc

Large diffs are not rendered by default.

1 change: 1 addition & 0 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@
<module>spring-cloud-aws-app-config</module>
<module>spring-cloud-aws-kinesis</module>
<module>spring-cloud-aws-kinesis-stream-binder</module>
<module>spring-cloud-aws-starters/spring-cloud-aws-starter-kinesis</module>
<module>spring-cloud-aws-starters/spring-cloud-aws-starter-integration-kinesis</module>
<module>spring-cloud-aws-starters/spring-cloud-aws-starter-integration-kinesis-producer</module>
<module>spring-cloud-aws-starters/spring-cloud-aws-starter-integration-kinesis-client</module>
Expand Down
5 changes: 5 additions & 0 deletions spring-cloud-aws-autoconfigure/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,11 @@
<artifactId>spring-cloud-aws-dynamodb</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>io.awspring.cloud</groupId>
<artifactId>spring-cloud-aws-kinesis</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>io.awspring.cloud</groupId>
<artifactId>spring-cloud-aws-s3</artifactId>
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
/*
* Copyright 2013-2022 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 io.awspring.cloud.autoconfigure.kinesis;

import io.awspring.cloud.autoconfigure.AwsClientCustomizer;
import software.amazon.awssdk.services.cloudwatch.CloudWatchAsyncClientBuilder;

/**
* Callback interface that can be used to customize a {@link CloudWatchAsyncClientBuilder}.
*
* @author Matej Nedic
* @since 4.3.0
*/
@FunctionalInterface
public interface CloudwatchAsyncClientCustomizer extends AwsClientCustomizer<CloudWatchAsyncClientBuilder> {
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
/*
* Copyright 2013-2022 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 io.awspring.cloud.autoconfigure.kinesis;

import io.awspring.cloud.autoconfigure.AwsClientCustomizer;
import software.amazon.awssdk.services.dynamodb.DynamoDbAsyncClientBuilder;

/**
* Callback interface that can be used to customize a {@link DynamoDbAsyncClientBuilder}.
*
* @author Matej Nedic
* @since 4.3.0
*/
@FunctionalInterface
public interface DynamoDbAsyncClientCustomizer extends AwsClientCustomizer<DynamoDbAsyncClientBuilder> {
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
/*
* Copyright 2013-2022 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 io.awspring.cloud.autoconfigure.kinesis;

import io.awspring.cloud.autoconfigure.AwsClientCustomizer;
import software.amazon.awssdk.services.kinesis.KinesisAsyncClientBuilder;

/**
* Callback interface that can be used to customize a {@link KinesisAsyncClientBuilder}.
*
* @author Matej Nedic
* @since 4.3.0
*/
@FunctionalInterface
public interface KinesisAsyncClientCustomizer extends AwsClientCustomizer<KinesisAsyncClientBuilder> {
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,137 @@
/*
* Copyright 2013-2026 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 io.awspring.cloud.autoconfigure.kinesis;

import io.awspring.cloud.autoconfigure.AwsAsyncClientCustomizer;
import io.awspring.cloud.autoconfigure.core.AwsClientBuilderConfigurer;
import io.awspring.cloud.autoconfigure.core.AwsConnectionDetails;
import io.awspring.cloud.autoconfigure.core.CredentialsProviderAutoConfiguration;
import io.awspring.cloud.autoconfigure.core.RegionProviderAutoConfiguration;
import io.awspring.cloud.kinesis.config.KclBootstrapConfiguration;
import io.awspring.cloud.kinesis.config.KclMessageListenerContainerFactory;
import io.awspring.cloud.kinesis.listener.KclContainerOptions;
import io.awspring.cloud.kinesis.listener.errorhandler.ErrorHandler;
import io.awspring.cloud.kinesis.operations.KinesisOperations;
import io.awspring.cloud.kinesis.operations.KinesisTemplate;
import org.springframework.beans.factory.ObjectProvider;
import org.springframework.boot.autoconfigure.AutoConfiguration;
import org.springframework.boot.autoconfigure.AutoConfigureAfter;
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.boot.context.properties.PropertyMapper;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Import;
import org.springframework.util.MimeTypeUtils;
import software.amazon.awssdk.services.cloudwatch.CloudWatchAsyncClient;
import software.amazon.awssdk.services.dynamodb.DynamoDbAsyncClient;
import software.amazon.awssdk.services.kinesis.KinesisAsyncClient;
import tools.jackson.databind.json.JsonMapper;

/**
* @author Matej Nedic
* @since 4.2.0
*/
@AutoConfiguration
@ConditionalOnClass({KinesisAsyncClient.class, KclBootstrapConfiguration.class})
@EnableConfigurationProperties(KinesisProperties.class)
@Import(KclBootstrapConfiguration.class)
@AutoConfigureAfter({CredentialsProviderAutoConfiguration.class, RegionProviderAutoConfiguration.class})
@ConditionalOnProperty(name = "spring.cloud.aws.kinesis.enabled", havingValue = "true", matchIfMissing = true)
public class KinesisAutoConfiguration {

private final KinesisProperties properties;

public KinesisAutoConfiguration(KinesisProperties properties) {
this.properties = properties;
}

@ConditionalOnMissingBean
@Bean
public KinesisAsyncClient kinesisAsyncClient(AwsClientBuilderConfigurer awsClientBuilderConfigurer,
ObjectProvider<KinesisAsyncClientCustomizer> kinesisAsyncClientCustomizers,
ObjectProvider<AwsAsyncClientCustomizer> awsAsyncClientCustomizers,
ObjectProvider<AwsConnectionDetails> connectionDetails) {
return awsClientBuilderConfigurer
.configureAsyncClient(KinesisAsyncClient.builder(), this.properties, connectionDetails.getIfAvailable(), kinesisAsyncClientCustomizers.orderedStream(),
awsAsyncClientCustomizers.orderedStream()).build();
}

@ConditionalOnMissingBean
@Bean
public DynamoDbAsyncClient kinesisDynamoDbAsyncClient(AwsClientBuilderConfigurer awsClientBuilderConfigurer,
ObjectProvider<DynamoDbAsyncClientCustomizer> dynamoDbAsyncClientCustomizers,
ObjectProvider<AwsAsyncClientCustomizer> awsAsyncClientCustomizers,
ObjectProvider<AwsConnectionDetails> connectionDetails) {
return awsClientBuilderConfigurer
.configureAsyncClient(DynamoDbAsyncClient.builder(), this.properties, connectionDetails.getIfAvailable(), dynamoDbAsyncClientCustomizers.orderedStream(),
awsAsyncClientCustomizers.orderedStream()).build();
}

@ConditionalOnMissingBean
@Bean
public CloudWatchAsyncClient kinesisCloudWatchAsyncClient(AwsClientBuilderConfigurer awsClientBuilderConfigurer,
ObjectProvider<CloudwatchAsyncClientCustomizer> cloudwatchAsyncClientCustomizers,
ObjectProvider<AwsAsyncClientCustomizer> awsAsyncClientCustomizers,
ObjectProvider<AwsConnectionDetails> connectionDetails) {
return awsClientBuilderConfigurer
.configureAsyncClient(CloudWatchAsyncClient.builder(), this.properties, connectionDetails.getIfAvailable(), cloudwatchAsyncClientCustomizers.orderedStream(),
awsAsyncClientCustomizers.orderedStream()).build();
}

@ConditionalOnMissingBean
@Bean
public KinesisTemplate kinesisTemplate(KinesisAsyncClient kinesisAsyncClient,
ObjectProvider<JsonMapper> jsonMapperProvider) {
return new KinesisTemplate(kinesisAsyncClient, jsonMapperProvider.getIfAvailable(JsonMapper::new));
}

@ConditionalOnMissingBean
@Bean
public KclMessageListenerContainerFactory defaultKclListenerContainerFactory(KinesisAsyncClient kinesisAsyncClient,
DynamoDbAsyncClient dynamoDbAsyncClient, CloudWatchAsyncClient cloudWatchAsyncClient,
ObjectProvider<ErrorHandler> errorHandler, ObjectProvider<KinesisOperations> kinesisOperations) {
KclMessageListenerContainerFactory factory = new KclMessageListenerContainerFactory(kinesisAsyncClient,
dynamoDbAsyncClient, cloudWatchAsyncClient);
factory.configure(this::configureContainerOptions);
errorHandler.ifUnique(factory::setErrorHandler);
kinesisOperations.ifUnique(factory::setKinesisOperations);
return factory;
}

private void configureContainerOptions(KclContainerOptions.Builder options) {
PropertyMapper mapper = PropertyMapper.get();
KinesisProperties.Listener listener = this.properties.getListener();
mapper.from(listener.getMaxRecords()).to(options::maxRecords);
mapper.from(listener.getIdleTimeBetweenReads())
.to(duration -> options.idleTimeBetweenReadsInMillis(duration.toMillis()));
mapper.from(listener.getRetrievalMode()).to(options::retrievalMode);
mapper.from(listener.getCheckpointMode()).to(options::checkpointMode);
mapper.from(listener.getInitialPosition()).to(options::initialPositionInStream);
mapper.from(listener.getGracefulShutdownTimeout()).to(options::gracefulShutdownTimeout);
mapper.from(listener.getCheckpointRecordCount()).to(options::checkpointRecordCount);
mapper.from(listener.getCheckpointInterval()).to(options::checkpointInterval);
mapper.from(listener.getMetricsLevel()).to(options::metricsLevel);
mapper.from(listener.getMetricsNamespace()).to(options::metricsNamespace);
mapper.from(listener.getAutoStartup()).to(options::autoStartup);
mapper.from(listener.getPhase()).to(options::phase);
mapper.from(listener.getInitialPositionTimestamp()).to(options::initialPositionTimestamp);
mapper.from(listener.getBillingMode()).to(options::billingMode);
mapper.from(listener.getContentType()).as(MimeTypeUtils::parseMimeType).to(options::payloadContentType);
}

}
Loading
Loading