diff --git a/eclipse-code-formatter.xml b/eclipse-code-formatter.xml
new file mode 100644
index 0000000..c72937b
--- /dev/null
+++ b/eclipse-code-formatter.xml
@@ -0,0 +1,397 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/eclipse.importorder b/eclipse.importorder
new file mode 100644
index 0000000..080e73d
--- /dev/null
+++ b/eclipse.importorder
@@ -0,0 +1,7 @@
+#Organize Import Order
+#Wed Apr 26 12:53:22 EDT 2017
+4=\#
+3=org.springframework
+2=
+1=javax
+0=java
diff --git a/spring-cloud-stream-binder-kinesis-core/pom.xml b/spring-cloud-stream-binder-kinesis-core/pom.xml
index 033917b..e4742c0 100644
--- a/spring-cloud-stream-binder-kinesis-core/pom.xml
+++ b/spring-cloud-stream-binder-kinesis-core/pom.xml
@@ -1,6 +1,6 @@
+ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
4.0.0
@@ -26,6 +26,12 @@
spring-boot-configuration-processor
true
+
+
+ org.springframework.integration
+ spring-integration-test
+ test
+
diff --git a/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/KinesisConsumerProperties.java b/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/KinesisConsumerProperties.java
index 8cf4749..1aa5892 100644
--- a/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/KinesisConsumerProperties.java
+++ b/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/KinesisConsumerProperties.java
@@ -19,8 +19,29 @@ package org.springframework.cloud.stream.binder.kinesis.properties;
/**
*
* @author Peter Oates
+ * @author Jacob Severson
*
*/
public class KinesisConsumerProperties {
+ private int startTimeout = 60000;
+
+ private int describeStreamRetries = 50;
+
+ public int getStartTimeout() {
+ return this.startTimeout;
+ }
+
+ public void setStartTimeout(int startTimeout) {
+ this.startTimeout = startTimeout;
+ }
+
+ public int getDescribeStreamRetries() {
+ return this.describeStreamRetries;
+ }
+
+ public void setDescribeStreamRetries(int describeStreamRetries) {
+ this.describeStreamRetries = describeStreamRetries;
+ }
+
}
diff --git a/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/provisioning/KinesisStreamProvisioner.java b/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/provisioning/KinesisStreamProvisioner.java
index 416b373..c52b209 100644
--- a/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/provisioning/KinesisStreamProvisioner.java
+++ b/spring-cloud-stream-binder-kinesis-core/src/main/java/org/springframework/cloud/stream/binder/kinesis/provisioning/KinesisStreamProvisioner.java
@@ -17,6 +17,11 @@
package org.springframework.cloud.stream.binder.kinesis.provisioning;
import com.amazonaws.services.kinesis.AmazonKinesis;
+import com.amazonaws.services.kinesis.model.DescribeStreamResult;
+import com.amazonaws.services.kinesis.model.ResourceNotFoundException;
+
+import org.apache.commons.logging.Log;
+import org.apache.commons.logging.LogFactory;
import org.springframework.cloud.stream.binder.ExtendedConsumerProperties;
import org.springframework.cloud.stream.binder.ExtendedProducerProperties;
@@ -34,12 +39,15 @@ import org.springframework.util.Assert;
*
* @author Peter Oates
* @author Artem Bilan
+ * @author Jacob Severson
*
*/
public class KinesisStreamProvisioner
implements
ProvisioningProvider, ExtendedProducerProperties> {
+ private final Log logger = LogFactory.getLog(getClass());
+
private final AmazonKinesis amazonKinesis;
private final KinesisBinderConfigurationProperties configurationProperties;
@@ -56,18 +64,51 @@ public class KinesisStreamProvisioner
public ProducerDestination provisionProducerDestination(String name,
ExtendedProducerProperties properties) throws ProvisioningException {
- KinesisProducerDestination producer = new KinesisProducerDestination(name);
+ if (logger.isInfoEnabled()) {
+ logger.info("Using Kinesis stream for outbound: " + name);
+ }
- return producer;
+ return new KinesisProducerDestination(name, createOrUpdate(name, properties.getPartitionCount()));
}
@Override
public ConsumerDestination provisionConsumerDestination(String name, String group,
ExtendedConsumerProperties properties) throws ProvisioningException {
- KinesisConsumerDestination consumer = new KinesisConsumerDestination(name);
+ if (logger.isInfoEnabled()) {
+ logger.info("Using Kinesis stream for inbound: " + name);
+ }
- return consumer;
+ int shardCount = properties.getInstanceCount() * properties.getConcurrency();
+
+ return new KinesisConsumerDestination(name, createOrUpdate(name, shardCount));
+ }
+
+ private Integer createOrUpdate(String name, Integer shards) {
+
+ try {
+ DescribeStreamResult streamResult = amazonKinesis.describeStream(name);
+
+ if (logger.isInfoEnabled()) {
+ logger.info("Stream found, using existing stream");
+ }
+
+ return streamResult.getStreamDescription().getShards().size();
+
+ }
+ catch (ResourceNotFoundException e) {
+ if (logger.isInfoEnabled()) {
+ logger.info("Stream not found");
+ }
+ }
+
+ if (logger.isInfoEnabled()) {
+ logger.info("Attempting to create stream");
+ }
+
+ amazonKinesis.createStream(name, shards);
+
+ return shards;
}
private static final class KinesisProducerDestination implements ProducerDestination {
@@ -76,10 +117,6 @@ public class KinesisStreamProvisioner
private final int shards;
- KinesisProducerDestination(String streamName) {
- this(streamName, 0);
- }
-
KinesisProducerDestination(String streamName, Integer shards) {
this.streamName = streamName;
this.shards = shards;
@@ -113,10 +150,6 @@ public class KinesisStreamProvisioner
private final String dlqName;
- KinesisConsumerDestination(String streamName) {
- this(streamName, 0, null);
- }
-
KinesisConsumerDestination(String streamName, int shards) {
this(streamName, shards, null);
}
@@ -140,7 +173,6 @@ public class KinesisStreamProvisioner
", dlqName='" + dlqName + '\'' +
'}';
}
-
}
}
diff --git a/spring-cloud-stream-binder-kinesis-core/src/test/java/org/springframework/cloud/stream/binder/kinesis/provisioning/KinesisStreamProvisionerTests.java b/spring-cloud-stream-binder-kinesis-core/src/test/java/org/springframework/cloud/stream/binder/kinesis/provisioning/KinesisStreamProvisionerTests.java
new file mode 100644
index 0000000..a915c49
--- /dev/null
+++ b/spring-cloud-stream-binder-kinesis-core/src/test/java/org/springframework/cloud/stream/binder/kinesis/provisioning/KinesisStreamProvisionerTests.java
@@ -0,0 +1,149 @@
+/*
+ * 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.provisioning;
+
+import java.util.Collections;
+import java.util.List;
+
+import com.amazonaws.services.kinesis.AmazonKinesis;
+import com.amazonaws.services.kinesis.model.CreateStreamResult;
+import com.amazonaws.services.kinesis.model.DescribeStreamResult;
+import com.amazonaws.services.kinesis.model.ResourceNotFoundException;
+import com.amazonaws.services.kinesis.model.Shard;
+import com.amazonaws.services.kinesis.model.StreamDescription;
+
+import org.junit.Test;
+
+import org.springframework.cloud.stream.binder.ExtendedConsumerProperties;
+import org.springframework.cloud.stream.binder.ExtendedProducerProperties;
+import org.springframework.cloud.stream.binder.kinesis.properties.KinesisBinderConfigurationProperties;
+import org.springframework.cloud.stream.binder.kinesis.properties.KinesisConsumerProperties;
+import org.springframework.cloud.stream.binder.kinesis.properties.KinesisProducerProperties;
+import org.springframework.cloud.stream.provisioning.ConsumerDestination;
+import org.springframework.cloud.stream.provisioning.ProducerDestination;
+
+import static org.hamcrest.Matchers.is;
+import static org.junit.Assert.assertThat;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * @author Jacob Severson
+ */
+public class KinesisStreamProvisionerTests {
+
+ @Test
+ public void testProvisionProducerSuccessfulWithExistingStream() {
+ AmazonKinesis amazonKinesisMock = mock(AmazonKinesis.class);
+ KinesisBinderConfigurationProperties binderProperties = new KinesisBinderConfigurationProperties();
+ KinesisStreamProvisioner provisioner = new KinesisStreamProvisioner(amazonKinesisMock, binderProperties);
+ ExtendedProducerProperties extendedProducerProperties = new ExtendedProducerProperties<>(
+ new KinesisProducerProperties());
+ String name = "test-stream";
+
+ DescribeStreamResult describeStreamResult = describeStreamResultWithShards(
+ Collections.singletonList(new Shard()));
+
+ when(amazonKinesisMock.describeStream(name)).thenReturn(describeStreamResult);
+
+ ProducerDestination destination = provisioner.provisionProducerDestination(name, extendedProducerProperties);
+
+ verify(amazonKinesisMock, times(1)).describeStream(name);
+ assertThat(destination.getName(), is(name));
+ }
+
+ @Test
+ public void testProvisionConsumerSuccessfulWithExistingStream() {
+ AmazonKinesis amazonKinesisMock = mock(AmazonKinesis.class);
+ KinesisBinderConfigurationProperties binderProperties = new KinesisBinderConfigurationProperties();
+ KinesisStreamProvisioner provisioner = new KinesisStreamProvisioner(amazonKinesisMock, binderProperties);
+
+ ExtendedConsumerProperties extendedConsumerProperties = new ExtendedConsumerProperties<>(
+ new KinesisConsumerProperties());
+
+ String name = "test-stream";
+ String group = "test-group";
+
+ DescribeStreamResult describeStreamResult = describeStreamResultWithShards(
+ Collections.singletonList(new Shard()));
+
+ when(amazonKinesisMock.describeStream(name)).thenReturn(describeStreamResult);
+
+ ConsumerDestination destination = provisioner.provisionConsumerDestination(name, group,
+ extendedConsumerProperties);
+
+ verify(amazonKinesisMock, times(1)).describeStream(name);
+ assertThat(destination.getName(), is(name));
+ }
+
+ @Test
+ public void testProvisionProducerSuccessfulWithNewStream() {
+ AmazonKinesis amazonKinesisMock = mock(AmazonKinesis.class);
+ KinesisBinderConfigurationProperties binderProperties = new KinesisBinderConfigurationProperties();
+ KinesisStreamProvisioner provisioner = new KinesisStreamProvisioner(amazonKinesisMock, binderProperties);
+ ExtendedProducerProperties extendedProducerProperties = new ExtendedProducerProperties<>(
+ new KinesisProducerProperties());
+ String name = "test-stream";
+ Integer shards = 1;
+
+ when(amazonKinesisMock.describeStream(name)).thenThrow(new ResourceNotFoundException("I got nothing"));
+ when(amazonKinesisMock.createStream(name, shards)).thenReturn(new CreateStreamResult());
+
+ ProducerDestination destination = provisioner.provisionProducerDestination(name, extendedProducerProperties);
+
+ verify(amazonKinesisMock, times(1)).describeStream(name);
+ verify(amazonKinesisMock, times(1)).createStream(name, shards);
+ assertThat(destination.getName(), is(name));
+ }
+
+ @Test
+ public void testProvisionConsumerSuccessfulWithNewStream() {
+ AmazonKinesis amazonKinesisMock = mock(AmazonKinesis.class);
+ KinesisBinderConfigurationProperties binderProperties = new KinesisBinderConfigurationProperties();
+ KinesisStreamProvisioner provisioner = new KinesisStreamProvisioner(amazonKinesisMock, binderProperties);
+ int instanceCount = 1;
+ int concurrency = 1;
+
+ ExtendedConsumerProperties extendedConsumerProperties = new ExtendedConsumerProperties<>(
+ new KinesisConsumerProperties());
+ extendedConsumerProperties.setInstanceCount(instanceCount);
+ extendedConsumerProperties.setConcurrency(concurrency);
+
+ String name = "test-stream";
+ String group = "test-group";
+
+ when(amazonKinesisMock.describeStream(name)).thenThrow(new ResourceNotFoundException("I got nothing"));
+ when(amazonKinesisMock.createStream(name, instanceCount * concurrency)).thenReturn(new CreateStreamResult());
+
+ ConsumerDestination destination = provisioner.provisionConsumerDestination(name, group,
+ extendedConsumerProperties);
+
+ verify(amazonKinesisMock, times(1)).describeStream(name);
+ verify(amazonKinesisMock, times(1)).createStream(name, instanceCount * concurrency);
+ assertThat(destination.getName(), is(name));
+ }
+
+ private static DescribeStreamResult describeStreamResultWithShards(List shards) {
+ return new DescribeStreamResult()
+ .withStreamDescription(
+ new StreamDescription()
+ .withShards(shards));
+ }
+
+}
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 7d83ebe..009a138 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
@@ -48,6 +48,15 @@ For general binding configuration options and properties, please refer to the ht
The following properties are available for Kinesis consumers only and must be prefixed with `spring.cloud.stream.kinesis.bindings..consumer.`.
+startTimeout::
+ The amount of time to wait for the consumer to start, in milliseconds.
++
+Default: `60000`.
+describeStreamRetries::
+ The amount of times the consumer will retry a `DescribeStream` operation waiting for the stream to be in `ACTIVE` state.
++
+Default: `50`.
+
=== Kinesis Producer Properties
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 de639dc..31ee588 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
@@ -108,8 +108,8 @@ public class KinesisMessageChannelBinder extends
// need to move these properties to the appropriate properties class
adapter.setCheckpointMode(CheckpointMode.record);
adapter.setListenerMode(ListenerMode.record);
- adapter.setStartTimeout(10000);
- adapter.setDescribeStreamRetries(1);
+ adapter.setStartTimeout(properties.getExtension().getStartTimeout());
+ adapter.setDescribeStreamRetries(properties.getExtension().getDescribeStreamRetries());
adapter.setConcurrency(10);
// Deffer byte[] conversion to the ReceivingHandler
diff --git a/spring-cloud-stream-binder-kinesis/src/test/java/org/springframework/cloud/stream/binder/kinesis/KinesisBinderTests.java b/spring-cloud-stream-binder-kinesis/src/test/java/org/springframework/cloud/stream/binder/kinesis/KinesisBinderTests.java
index 5d07739..1a66faf 100644
--- a/spring-cloud-stream-binder-kinesis/src/test/java/org/springframework/cloud/stream/binder/kinesis/KinesisBinderTests.java
+++ b/spring-cloud-stream-binder-kinesis/src/test/java/org/springframework/cloud/stream/binder/kinesis/KinesisBinderTests.java
@@ -16,8 +16,13 @@
package org.springframework.cloud.stream.binder.kinesis;
-import org.junit.ClassRule;
+import com.amazonaws.services.kinesis.model.DescribeStreamResult;
+import org.junit.ClassRule;
+import org.junit.Test;
+
+import org.springframework.cloud.stream.binder.Binder;
+import org.springframework.cloud.stream.binder.Binding;
import org.springframework.cloud.stream.binder.ExtendedConsumerProperties;
import org.springframework.cloud.stream.binder.ExtendedProducerProperties;
import org.springframework.cloud.stream.binder.PartitionCapableBinderTests;
@@ -25,9 +30,14 @@ import org.springframework.cloud.stream.binder.Spy;
import org.springframework.cloud.stream.binder.kinesis.properties.KinesisBinderConfigurationProperties;
import org.springframework.cloud.stream.binder.kinesis.properties.KinesisConsumerProperties;
import org.springframework.cloud.stream.binder.kinesis.properties.KinesisProducerProperties;
+import org.springframework.integration.channel.DirectChannel;
+
+import static org.hamcrest.Matchers.is;
+import static org.junit.Assert.assertThat;
/**
* @author Artem Bilan
+ * @author Jacob Severson
*
*/
public class KinesisBinderTests
@@ -38,6 +48,26 @@ public class KinesisBinderTests
@ClassRule
public static LocalKinesisResource localKinesisResource = new LocalKinesisResource();
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testAutoCreateStreamForNonExistingStream() throws Exception {
+ Binder binder = getBinder();
+ DirectChannel output = new DirectChannel();
+ ExtendedConsumerProperties consumerProperties = createConsumerProperties();
+ String testStreamName = "nonexisting" + System.currentTimeMillis();
+ Binding> binding = binder.bindConsumer(testStreamName, "test", output, consumerProperties);
+ binding.unbind();
+
+ DescribeStreamResult streamResult = localKinesisResource.getResource().describeStream(testStreamName);
+ String createdStreamName = streamResult.getStreamDescription().getStreamName();
+ int createdShards = streamResult.getStreamDescription().getShards().size();
+ String createdStreamStatus = streamResult.getStreamDescription().getStreamStatus();
+
+ assertThat(createdStreamName, is(testStreamName));
+ assertThat(createdShards, is(consumerProperties.getInstanceCount() * consumerProperties.getConcurrency()));
+ assertThat(createdStreamStatus, is("ACTIVE"));
+ }
+
@Override
protected boolean usesExplicitRouting() {
return false;
@@ -51,16 +81,15 @@ public class KinesisBinderTests
@Override
protected KinesisTestBinder getBinder() throws Exception {
if (this.testBinder == null) {
- this.testBinder =
- new KinesisTestBinder(localKinesisResource.getResource(),
- new KinesisBinderConfigurationProperties());
+ this.testBinder = new KinesisTestBinder(localKinesisResource.getResource(),
+ new KinesisBinderConfigurationProperties());
}
return this.testBinder;
}
@Override
protected ExtendedConsumerProperties createConsumerProperties() {
- final ExtendedConsumerProperties kafkaConsumerProperties =
+ ExtendedConsumerProperties kafkaConsumerProperties =
new ExtendedConsumerProperties<>(new KinesisConsumerProperties());
// set the default values that would normally be propagated by Spring Cloud Stream
kafkaConsumerProperties.setInstanceCount(1);
diff --git a/spring-cloud-stream-binder-kinesis/src/test/java/org/springframework/cloud/stream/binder/kinesis/KinesisTestBinder.java b/spring-cloud-stream-binder-kinesis/src/test/java/org/springframework/cloud/stream/binder/kinesis/KinesisTestBinder.java
index ab8f484..1ad044e 100644
--- a/spring-cloud-stream-binder-kinesis/src/test/java/org/springframework/cloud/stream/binder/kinesis/KinesisTestBinder.java
+++ b/spring-cloud-stream-binder-kinesis/src/test/java/org/springframework/cloud/stream/binder/kinesis/KinesisTestBinder.java
@@ -37,6 +37,7 @@ public class KinesisTestBinder
public KinesisTestBinder(AmazonKinesisAsync amazonKinesis,
KinesisBinderConfigurationProperties kinesisBinderConfigurationProperties) {
+
KinesisStreamProvisioner provisioningProvider =
new KinesisStreamProvisioner(amazonKinesis, kinesisBinderConfigurationProperties);
diff --git a/spring-cloud-stream-binder-kinesis/src/test/java/org/springframework/cloud/stream/binder/kinesis/LocalKinesisResource.java b/spring-cloud-stream-binder-kinesis/src/test/java/org/springframework/cloud/stream/binder/kinesis/LocalKinesisResource.java
index 79dc19c..109b0ac 100644
--- a/spring-cloud-stream-binder-kinesis/src/test/java/org/springframework/cloud/stream/binder/kinesis/LocalKinesisResource.java
+++ b/spring-cloud-stream-binder-kinesis/src/test/java/org/springframework/cloud/stream/binder/kinesis/LocalKinesisResource.java
@@ -18,7 +18,8 @@ package org.springframework.cloud.stream.binder.kinesis;
import com.amazonaws.ClientConfiguration;
import com.amazonaws.SDKGlobalConfiguration;
-import com.amazonaws.auth.AWSCredentialsProvider;
+import com.amazonaws.auth.AWSStaticCredentialsProvider;
+import com.amazonaws.auth.BasicAWSCredentials;
import com.amazonaws.client.builder.AwsClientBuilder;
import com.amazonaws.regions.Regions;
import com.amazonaws.services.kinesis.AmazonKinesisAsync;
@@ -26,10 +27,9 @@ import com.amazonaws.services.kinesis.AmazonKinesisAsyncClientBuilder;
import org.springframework.cloud.stream.test.junit.AbstractExternalResourceTestSupport;
-import static org.mockito.Mockito.mock;
-
/**
* @author Artem Bilan
+ * @author Jacob Severson
*
*/
public class LocalKinesisResource extends AbstractExternalResourceTestSupport {
@@ -53,7 +53,6 @@ public class LocalKinesisResource extends AbstractExternalResourceTestSupport