diff --git a/pom.xml b/pom.xml
index 77bbc07..c7eb9ae 100644
--- a/pom.xml
+++ b/pom.xml
@@ -6,7 +6,7 @@
org.springframework.cloud
spring-cloud-build
- 1.3.3.BUILD-SNAPSHOT
+ 1.3.5.RELEASE
@@ -87,18 +87,6 @@
maven-antrun-plugin
1.7
-
- org.apache.maven.plugins
- maven-checkstyle-plugin
- 2.17
-
-
- com.puppycrawl.tools
- checkstyle
- 7.1
-
-
-
org.apache.maven.plugins
maven-javadoc-plugin
@@ -126,23 +114,11 @@
${spring-cloud-stream.version}
-
-
- checkstyle-validation
- validate
-
- checkstyle.xml
- checkstyle-header.txt
- UTF-8
- true
- false
- true
-
-
- check
-
-
-
+
+ checkstyle.xml
+ checkstyle-header.txt
+ true
+
diff --git a/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/CheckpointProperties.java b/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/CheckpointProperties.java
new file mode 100644
index 0000000..8d4e103
--- /dev/null
+++ b/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/CheckpointProperties.java
@@ -0,0 +1,56 @@
+/*
+ * Copyright 2017 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
+ *
+ * http://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.cloud.stream.binder.kinesis.properties;
+
+/**
+ *
+ * @author Peter Oates
+ *
+ */
+public class CheckpointProperties {
+
+ private String checkpointTable = "checkpoint";
+
+ private Long dynamoDbReadCapacity = 1L;
+
+ private Long dynamoDbWriteCapacity = 1L;
+
+ public String getCheckpointTable() {
+ return checkpointTable;
+ }
+
+ public void setCheckpointTable(String checkpointTable) {
+ this.checkpointTable = checkpointTable;
+ }
+
+ public Long getDynamoDbReadCapacity() {
+ return dynamoDbReadCapacity;
+ }
+
+ public void setDynamoDbReadCapacity(Long dynamoDbReadCapacity) {
+ this.dynamoDbReadCapacity = dynamoDbReadCapacity;
+ }
+
+ public Long getDynamoDbWriteCapacity() {
+ return dynamoDbWriteCapacity;
+ }
+
+ public void setDynamoDbWriteCapacity(Long dynamoDbWriteCapacity) {
+ this.dynamoDbWriteCapacity = dynamoDbWriteCapacity;
+ }
+
+}
diff --git a/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/KinesisBinderConfigurationProperties.java b/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/KinesisBinderConfigurationProperties.java
index 8402b1c..021b6f9 100644
--- a/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/KinesisBinderConfigurationProperties.java
+++ b/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/KinesisBinderConfigurationProperties.java
@@ -31,6 +31,8 @@ public class KinesisBinderConfigurationProperties {
private int describeStreamBackoff = 1000;
private int describeStreamRetries = 50;
+
+ private CheckpointProperties checkpoint = new CheckpointProperties();
public String[] getHeaders() {
return this.headers;
@@ -55,4 +57,12 @@ public class KinesisBinderConfigurationProperties {
public void setDescribeStreamRetries(int describeStreamRetries) {
this.describeStreamRetries = describeStreamRetries;
}
+
+ public CheckpointProperties getCheckpoint() {
+ return checkpoint;
+ }
+
+ public void setCheckpoint(CheckpointProperties checkpoint) {
+ this.checkpoint = checkpoint;
+ }
}
diff --git a/spring-cloud-stream-binder-kinesis-docs/src/main/asciidoc/overview.adoc b/spring-cloud-stream-binder-kinesis-docs/src/main/asciidoc/overview.adoc
index 9afc092..75a74f4 100644
--- a/spring-cloud-stream-binder-kinesis-docs/src/main/asciidoc/overview.adoc
+++ b/spring-cloud-stream-binder-kinesis-docs/src/main/asciidoc/overview.adoc
@@ -59,6 +59,21 @@ For general binding configuration options and properties, please refer to the ht
[[kinesis-binder-properties]]
=== Kinesis Binder Properties
+Checkpoint properties prefixed with `spring.cloud.strea.kinesis.binder.checkpoint`
+
+checkpointTable::
+ The name to give the DynamoDb table
++
+Default: `checkpoint`
+dynamoDbReadCapacity::
+ The Read capacity of the DynamoDb table. See [AWS docunmentation](http://docs.aws.amazon.com/amazondynamodb/latest/developerguide/HowItWorks.ProvisionedThroughput.html)
++
+Default: `1`
+dynamoDbWriteCapacity::
+ The write capacity of the DynamoDb table. See [AWS docunmentation](http://docs.aws.amazon.com/amazondynamodb/latest/developerguide/HowItWorks.ProvisionedThroughput.html)
++
+Default: `1`
+
=== Kinesis Consumer Properties
The following properties are available for Kinesis consumers only and must be prefixed with `spring.cloud.stream.kinesis.bindings..consumer.`.
diff --git a/spring-cloud-stream-binder-kinesis/pom.xml b/spring-cloud-stream-binder-kinesis/pom.xml
index ea16c42..20d8e1e 100644
--- a/spring-cloud-stream-binder-kinesis/pom.xml
+++ b/spring-cloud-stream-binder-kinesis/pom.xml
@@ -32,6 +32,10 @@
org.springframework.cloud
spring-cloud-starter-aws
+
+ com.amazonaws
+ aws-java-sdk-dynamodb
+
org.springframework.cloud
diff --git a/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/KinesisMessageChannelBinder.java b/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/KinesisMessageChannelBinder.java
index 3e6abf6..cbef70d 100644
--- a/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/KinesisMessageChannelBinder.java
+++ b/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/KinesisMessageChannelBinder.java
@@ -44,6 +44,7 @@ import org.springframework.integration.aws.inbound.kinesis.KinesisShardOffset;
import org.springframework.integration.aws.inbound.kinesis.ListenerMode;
import org.springframework.integration.aws.outbound.KinesisMessageHandler;
import org.springframework.integration.core.MessageProducer;
+import org.springframework.integration.metadata.MetadataStore;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
import org.springframework.util.Assert;
@@ -59,16 +60,14 @@ public class KinesisMessageChannelBinder extends
AbstractMessageChannelBinder, ExtendedProducerProperties, KinesisStreamProvisioner>
implements ExtendedPropertiesBinder {
- // not using currently but expect to include global properties here if
- // required
private final KinesisBinderConfigurationProperties configurationProperties;
- // not currently utilising any extended properties - not sure if I will need
- // these
private KinesisExtendedBindingProperties extendedBindingProperties = new KinesisExtendedBindingProperties();
private final AmazonKinesisAsync amazonKinesis;
+ private MetadataStore metadataStore;
+
public KinesisMessageChannelBinder(AmazonKinesisAsync amazonKinesis,
KinesisBinderConfigurationProperties configurationProperties,
KinesisStreamProvisioner provisioningProvider) {
@@ -114,7 +113,8 @@ public class KinesisMessageChannelBinder extends
kinesisMessageHandler.setSync(producerProperties.getExtension().isSync());
kinesisMessageHandler.setStream(destination.getName());
if (producerProperties.isPartitioned()) {
- kinesisMessageHandler.setPartitionKeyExpressionString("'partitionKey-' + headers." + BinderHeaders.PARTITION_HEADER);
+ kinesisMessageHandler
+ .setPartitionKeyExpressionString("'partitionKey-' + headers." + BinderHeaders.PARTITION_HEADER);
}
kinesisMessageHandler.setBeanFactory(getBeanFactory());
@@ -154,9 +154,10 @@ public class KinesisMessageChannelBinder extends
adapter.setConsumerGroup(consumerGroup);
adapter.setStreamInitialSequence(anonymous ? KinesisShardOffset.latest() : KinesisShardOffset.trimHorizon());
- // need to move these properties to the appropriate properties class
- adapter.setCheckpointMode(CheckpointMode.record);
+
adapter.setListenerMode(ListenerMode.record);
+ adapter.setCheckpointMode(CheckpointMode.record);
+ adapter.setCheckpointStore(metadataStore);
adapter.setConcurrency(properties.getConcurrency());
adapter.setStartTimeout(properties.getExtension().getStartTimeout());
@@ -169,4 +170,8 @@ public class KinesisMessageChannelBinder extends
return adapter;
}
+ public void setMetadataStore(MetadataStore metadataStore) {
+ this.metadataStore = metadataStore;
+ }
+
}
diff --git a/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/config/KinesisBinderConfiguration.java b/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/config/KinesisBinderConfiguration.java
index 1bb4239..afb88b5 100644
--- a/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/config/KinesisBinderConfiguration.java
+++ b/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/config/KinesisBinderConfiguration.java
@@ -17,6 +17,8 @@
package org.springframework.cloud.stream.binder.kinesis.config;
import com.amazonaws.auth.AWSCredentialsProvider;
+import com.amazonaws.services.dynamodbv2.AmazonDynamoDBAsync;
+import com.amazonaws.services.dynamodbv2.AmazonDynamoDBAsyncClientBuilder;
import com.amazonaws.services.kinesis.AmazonKinesisAsync;
import com.amazonaws.services.kinesis.AmazonKinesisAsyncClientBuilder;
@@ -31,7 +33,9 @@ import org.springframework.cloud.stream.binder.kinesis.properties.KinesisExtende
import org.springframework.cloud.stream.binder.kinesis.provisioning.KinesisStreamProvisioner;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
+import org.springframework.integration.aws.metadata.DynamoDbMetaDataStore;
import org.springframework.integration.codec.Codec;
+import org.springframework.integration.metadata.MetadataStore;
/**
*
@@ -48,14 +52,11 @@ public class KinesisBinderConfiguration {
private KinesisBinderConfigurationProperties configurationProperties;
@Bean
- public AmazonKinesisAsync amazonKinesis(
- AWSCredentialsProvider awsCredentialsProvider,
+ public AmazonKinesisAsync amazonKinesis(AWSCredentialsProvider awsCredentialsProvider,
RegionProvider regionProvider) {
- return AmazonKinesisAsyncClientBuilder.standard()
- .withCredentials(awsCredentialsProvider)
- .withRegion(regionProvider.getRegion().getName())
- .build();
+ return AmazonKinesisAsyncClientBuilder.standard().withCredentials(awsCredentialsProvider)
+ .withRegion(regionProvider.getRegion().getName()).build();
}
@Bean
@@ -64,17 +65,37 @@ public class KinesisBinderConfiguration {
}
@Bean
- KinesisMessageChannelBinder kinesisMessageChannelBinder(
- AmazonKinesisAsync amazonKinesis,
- KinesisStreamProvisioner provisioningProvider,
- Codec codec) {
+ KinesisMessageChannelBinder kinesisMessageChannelBinder(AmazonKinesisAsync amazonKinesis,
+ KinesisStreamProvisioner provisioningProvider, Codec codec, MetadataStore kinesisCheckpointStore) {
- KinesisMessageChannelBinder kinesisMessageChannelBinder =
- new KinesisMessageChannelBinder(amazonKinesis, this.configurationProperties, provisioningProvider);
+ KinesisMessageChannelBinder kinesisMessageChannelBinder = new KinesisMessageChannelBinder(amazonKinesis,
+ this.configurationProperties, provisioningProvider);
kinesisMessageChannelBinder.setCodec(codec);
+ kinesisMessageChannelBinder.setMetadataStore(kinesisCheckpointStore);
return kinesisMessageChannelBinder;
}
+ @Bean
+ @ConditionalOnMissingBean
+ MetadataStore kinesisCheckpointStore(AWSCredentialsProvider awsCredentialsProvider, RegionProvider regionProvider) {
+
+ String tableName = this.configurationProperties.getCheckpoint().getCheckpointTable();
+
+ MetadataStore kinesisCheckpointStore;
+ AmazonDynamoDBAsync dynamoDB = AmazonDynamoDBAsyncClientBuilder.standard()
+ .withCredentials(awsCredentialsProvider)
+ .withRegion(regionProvider.getRegion().getName())
+ .build();
+
+ kinesisCheckpointStore = new DynamoDbMetaDataStore(dynamoDB, tableName);
+ ((DynamoDbMetaDataStore) kinesisCheckpointStore)
+ .setReadCapacity(configurationProperties.getCheckpoint().getDynamoDbReadCapacity());
+ ((DynamoDbMetaDataStore) kinesisCheckpointStore)
+ .setWriteCapacity(configurationProperties.getCheckpoint().getDynamoDbWriteCapacity());
+
+
+ return kinesisCheckpointStore;
+ }
}