From dda290f2716ea3eb5bbe0a5fa832dd5d07ae8756 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 12 Oct 2016 19:20:27 -0400 Subject: [PATCH] GH-22: Add KinesisMessageHandler Fixes GH-22 (https://github.com/spring-projects/spring-integration-aws/issues/22) * Add `KinesisMessageHandler` based on the `AmazonKinesisAsync.putRecordAsync()` operation. The logic is fully similar to the `KafkaProducerMessageHandler` since both protocols pursue the same behavior * Add some Kinesis specific `AwsHeaders` * Upgrade to Gradle 2.14.1 --- build.gradle | 11 +- gradle/wrapper/gradle-wrapper.jar | Bin 53319 -> 53324 bytes gradle/wrapper/gradle-wrapper.properties | 4 +- .../aws/outbound/KinesisMessageHandler.java | 200 ++++++++++++++++++ .../integration/aws/support/AwsHeaders.java | 15 ++ .../outbound/KinesisMessageHandlerTests.java | 160 ++++++++++++++ 6 files changed, 385 insertions(+), 5 deletions(-) create mode 100644 src/main/java/org/springframework/integration/aws/outbound/KinesisMessageHandler.java create mode 100644 src/test/java/org/springframework/integration/aws/outbound/KinesisMessageHandlerTests.java diff --git a/build.gradle b/build.gradle index 7e4254e..9372e43 100644 --- a/build.gradle +++ b/build.gradle @@ -12,7 +12,7 @@ plugins { id 'eclipse' id 'idea' id 'jacoco' - id 'org.sonarqube' version '1.2' + id 'org.sonarqube' version '2.1' id 'checkstyle' } description = 'Spring Integration AWS Support' @@ -52,10 +52,11 @@ if (project.hasProperty('platformVersion')) { ext { assertjVersion = '3.5.2' + awsKinesisVersion = '1.7.0' servletApiVersion = '3.1.0' slf4jVersion = '1.7.21' springCloudAwsVersion = '1.1.1.RELEASE' - springIntegrationVersion = '4.3.2.RELEASE' + springIntegrationVersion = '4.3.4.RELEASE' idPrefix = 'aws' @@ -81,12 +82,15 @@ checkstyle { } dependencies { - compile "org.springframework.integration:spring-integration-core:$springIntegrationVersion" compile "org.springframework.cloud:spring-cloud-aws-core:$springCloudAwsVersion" + compile("org.springframework.cloud:spring-cloud-aws-messaging:$springCloudAwsVersion", optional) compile("org.springframework.integration:spring-integration-file:$springIntegrationVersion", optional) compile("org.springframework.integration:spring-integration-http:$springIntegrationVersion", optional) + + compile("com.amazonaws:amazon-kinesis-client:$awsKinesisVersion", optional) + compile("javax.servlet:javax.servlet-api:$servletApiVersion", provided) testCompile "org.springframework.integration:spring-integration-test:$springIntegrationVersion" @@ -153,6 +157,7 @@ jacocoTestReport { } } +check.dependsOn javadoc build.dependsOn jacocoTestReport task sourcesJar(type: Jar) { diff --git a/gradle/wrapper/gradle-wrapper.jar b/gradle/wrapper/gradle-wrapper.jar index d3b83982b9b1bccad955349d702be9b884c6e049..3baa851b28c65f87dd36a6748e1a85cf360c1301 100644 GIT binary patch delta 2176 zcmZWp3rtg27(TtXg$ixK@+#g!E0`*xQTs#F%>Uf>oed#k~5qmA^$N)myQ%yOlJAT zF%6f|+5AGMOZxb=R>kBp2z)N>M_`+*41ro7xt9wfp{iQ~;ySTcib%iD5dq25gt#J7 zu3P7Y!C1ILhJh;5*ByhYOGe!=n7pc2FqEuy73A)wwQ_!orY0vTA1xg*+V?xkcS1FX zs3TdQZ~sNrtoL5rchVHLwISuH|N1*gNsncZ{*0XxaysPK{`!624=1PsHuae*D?QG> zpQ!%2?rvGvs-D96UNzktY$w_0Wn|9t6*;!3AoUNDhXG+x9opd&O@-G_g>HMY24G(rb`-A%jYg_)VrxjzL4Nj>X{ij%Xp!-z z?TU!yt1bcgwy+BQYL^vNsBvqkWLQA%m);eE*7A!|+}c~k3eI0t4+&elYc>i#w=Tp8 zoPhdk0y5OFM&QID~kd1Hr}kCaszS>)(|pHls6kmW16)zrDYOoKW*U{ z?Gn_UYYE0OZgouDRRSf{sb24-FCA_PLr#W1<4B!^kaa8}ks|*mG{zc-)`eD%Attn* zWl~YQRVd)Lm6NQePpsj3{=IU348R0(b^huK>XT_#@+7LYDfmlIvu&ItjUKWE;=gJe z$Iuiy?C_`eZ3cvttsE1WN++iU)0M4WXudjjLO=x!WE85I@3iJGnOhnV7t+R7!Mq)` zigeWK_qxy(>Av)wU4b40!OD?1s5Ak5!bBebYz2XqhGzIt|2BcV$ZyM7t1BS~;My+` zk&80X{Ys{m=8u!LnZ7mc;MxP(%@vW)=Q>jv%Jp=na{r_5seK^W1}7{M39j!zfnsMM z@2=w;8Bvl${)|bk2l2r7PWEA9|iPq zRS<1&Q{x$hF%emxfZiSS^F(J{qOaQCL0P<=V-6Rh+d+F6LX7||C;}*mzO<_`lVp#D z*Gr+a5GYOjo(+_scC{lAiE&+<2af{r1e7Kw-YFeV*@bG$6runBZt?xtWw^H{1z!0{ z?u0~(ay8Q-M?5CBpq*oORY8ewIVXCt>3r^YaNZdhl6azfs-5RGR{J`qI#p^^_(M<9 z`D#SxWY>T|OJ~?QJnX#=j_JC99`3c=8~IWaeT?tz3-E~&`S{S^JB(;Cy^~`G8qgx6 mGXSB304=(NP+O+~r{pnaF?oZ78Xld~ z6N;x&M;f)+^mdeHuqEwRj6nCY1qM8oLx@i4ZiG+~!%EWZP;li7R67pC#h3XClONGH5 zio@TL!T@2C%fi^8<%&22@5_Eg@TPn(f(D(^&kISo<}4%Lz87RjhWxiNvO#ZEAUR;p z@WY8DS|!JcHa1`^POP);j?vIm$!FcZ>+D{HG-{wyW(2iU(~vReAj=_5S1?EQf^?-2 z2FtTD)m9pn8ZsED5jxYb@aN3T+lq?ceP2$=EBWz{ zft`PK?x?7DdFH(scP8AJsL+J3TT?xwb)e;&pH8nX`sIZkVVi@iNgS80$qtFG^lN(e zT32G(2VZXeblTg8f}+DGwJ$0TUwLAC%2Zq9!z1RHlMiZq9v9P^B5}i8*uM3GsFlfmLRHZhg_AV+d*<2@!GEq*pvEO85sQMC{xb`>6aA`}8lsT8z4zs0q>J~Gfi}guNGvRPQ zBX^F>XPVSwfKhOzIZaBt3St~n#oYzPju^a{S_iM&O)q41p)De9wuBomWi2UO+9d}s zaln0t5tW7(Uh*6Grt2WJ1$zNrZ{azU1ionVhct;6s#;Vitf5%KjbLv`;6g#u%9c#6 zJhxVcOLJR!RIj9HsT|Rs)`=FOwF0j!KJSzha*u8$n(pn<(^#rhqA*Tr0b?7pw6^hF zkOuLR4h1aK_)UJ|9Lc#QuEFgxLK3_QiF{;6$HvloU~ij&IrpuPJet8#Nb{OSnsfeW zt`1OvrJd!BXy>_=dc=0UA1GZa^oYIKbbY^Ahk0j zQ_u)dsG#upuz~(-$-hZi*{s2Ag-de_mPh0i*H#R!Ycp zk*28(2Fl_gv6tnU+sktwm1CZYUe*&^_BHBk;|S@b%_#Pdg~49tpXKxX>1uj%E8U0q zRlaIRGq=YJGF=AT%uA= asyncHandler; + + private Converter converter = new SerializingConverter(); + + private EvaluationContext evaluationContext; + + private volatile Expression streamExpression; + + private volatile Expression partitionKeyExpression; + + private volatile Expression explicitHashKeyExpression; + + private volatile Expression sequenceNumberExpression; + + private boolean sync; + + private Expression sendTimeoutExpression = new ValueExpression<>(DEFAULT_SEND_TIMEOUT); + + public KinesisMessageHandler(AmazonKinesisAsync amazonKinesis) { + this.amazonKinesis = amazonKinesis; + } + + public void setAsyncHandler(AsyncHandler asyncHandler) { + this.asyncHandler = asyncHandler; + } + + public void setConverter(Converter converter) { + Assert.notNull(converter, "'converter' must not be null."); + this.converter = converter; + } + + public void setStream(String stream) { + setStreamExpression(new LiteralExpression(stream)); + } + + public void setStreamExpressionString(String streamExpression) { + setStreamExpression(EXPRESSION_PARSER.parseExpression(streamExpression)); + } + + public void setStreamExpression(Expression streamExpression) { + this.streamExpression = streamExpression; + } + + public void setPartitionKey(String partitionKey) { + setPartitionKeyExpression(new LiteralExpression(partitionKey)); + } + + public void setPartitionKeyExpressionString(String partitionKeyExpression) { + setPartitionKeyExpression(EXPRESSION_PARSER.parseExpression(partitionKeyExpression)); + } + + public void setPartitionKeyExpression(Expression partitionKeyExpression) { + this.partitionKeyExpression = partitionKeyExpression; + } + + public void setExplicitHashKey(String explicitHashKey) { + setExplicitHashKeyExpression(new LiteralExpression(explicitHashKey)); + } + + public void setExplicitHashKeyExpressionString(String explicitHashKeyExpression) { + setExplicitHashKeyExpression(EXPRESSION_PARSER.parseExpression(explicitHashKeyExpression)); + } + + public void setExplicitHashKeyExpression(Expression explicitHashKeyExpression) { + this.explicitHashKeyExpression = explicitHashKeyExpression; + } + + public void setSequenceNumberString(String sequenceNumberExpression) { + setSequenceNumberExpression(EXPRESSION_PARSER.parseExpression(sequenceNumberExpression)); + } + + public void setSequenceNumberExpression(Expression sequenceNumberExpression) { + this.sequenceNumberExpression = sequenceNumberExpression; + } + + public void setSync(boolean sync) { + this.sync = sync; + } + + public void setSendTimeoutExpression(Expression sendTimeoutExpression) { + this.sendTimeoutExpression = sendTimeoutExpression; + } + + @Override + protected void onInit() throws Exception { + super.onInit(); + this.evaluationContext = ExpressionUtils.createStandardEvaluationContext(getBeanFactory()); + } + + @Override + protected void handleMessageInternal(Message message) throws Exception { + String stream = message.getHeaders().get(AwsHeaders.STREAM, String.class); + if (!StringUtils.hasText(stream) && this.streamExpression != null) { + stream = this.streamExpression.getValue(this.evaluationContext, message, String.class); + } + Assert.state(stream != null, "'stream' must not be null for sending a Kinesis record. " + + "Consider configuring this handler with a 'stream'( or 'streamExpression') or supply an " + + "'aws_stream' message header."); + + String partitionKey = message.getHeaders().get(AwsHeaders.PARTITION_KEY, String.class); + if (!StringUtils.hasText(partitionKey) && this.partitionKeyExpression != null) { + partitionKey = this.partitionKeyExpression.getValue(this.evaluationContext, message, String.class); + } + Assert.state(partitionKey != null, "'partitionKey' must not be null for sending a Kinesis record. " + + "Consider configuring this handler with a 'partitionKey'( or 'partitionKeyExpression') or supply an " + + "'aws_partitionKey' message header."); + + String explicitHashKey = + (this.explicitHashKeyExpression != null + ? this.explicitHashKeyExpression.getValue(this.evaluationContext, message, String.class) + : null); + + String sequenceNumber = message.getHeaders().get(AwsHeaders.SEQUENCE_NUMBER, String.class); + if (!StringUtils.hasText(stream) && this.streamExpression != null) { + partitionKey = this.sequenceNumberExpression.getValue(this.evaluationContext, message, String.class); + } + + PutRecordRequest putRecordRequest = new PutRecordRequest() + .withStreamName(stream) + .withPartitionKey(partitionKey) + .withExplicitHashKey(explicitHashKey) + .withSequenceNumberForOrdering(sequenceNumber) + .withData(ByteBuffer.wrap(this.converter.convert(message.getPayload()))); + + Future resultFuture = this.amazonKinesis.putRecordAsync(putRecordRequest, this.asyncHandler); + + if (this.sync) { + Long sendTimeout = this.sendTimeoutExpression.getValue(this.evaluationContext, message, Long.class); + if (sendTimeout == null || sendTimeout < 0) { + resultFuture.get(); + } + else { + try { + resultFuture.get(sendTimeout, TimeUnit.MILLISECONDS); + } + catch (TimeoutException te) { + throw new MessageTimeoutException(message, "Timeout waiting for response from AmazonKinesis", te); + } + } + } + } + +} diff --git a/src/main/java/org/springframework/integration/aws/support/AwsHeaders.java b/src/main/java/org/springframework/integration/aws/support/AwsHeaders.java index d9ca874..61a90d8 100644 --- a/src/main/java/org/springframework/integration/aws/support/AwsHeaders.java +++ b/src/main/java/org/springframework/integration/aws/support/AwsHeaders.java @@ -65,4 +65,19 @@ public abstract class AwsHeaders { */ public static final String SNS_PUBLISHED_MESSAGE_ID = PREFIX + "snsPublishedMessageId"; + /** + * The {@value STREAM} header for sending/receiving data over Kinesis. + */ + public static final String STREAM = PREFIX + "stream"; + + /** + * The {@value PARTITION_KEY} header for sending/receiving data over Kinesis. + */ + public static final String PARTITION_KEY = PREFIX + "partitionKey"; + + /** + * The {@value SEQUENCE_NUMBER} header for sending/receiving data over Kinesis. + */ + public static final String SEQUENCE_NUMBER = PREFIX + "sequenceNumber"; + } diff --git a/src/test/java/org/springframework/integration/aws/outbound/KinesisMessageHandlerTests.java b/src/test/java/org/springframework/integration/aws/outbound/KinesisMessageHandlerTests.java new file mode 100644 index 0000000..f6905e4 --- /dev/null +++ b/src/test/java/org/springframework/integration/aws/outbound/KinesisMessageHandlerTests.java @@ -0,0 +1,160 @@ +/* + * Copyright 2016 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.integration.aws.outbound; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.BDDMockito.given; +import static org.mockito.Matchers.any; +import static org.mockito.Matchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; + +import java.nio.ByteBuffer; +import java.util.concurrent.Future; + +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.core.convert.converter.Converter; +import org.springframework.core.serializer.support.SerializingConverter; +import org.springframework.integration.annotation.ServiceActivator; +import org.springframework.integration.aws.support.AwsHeaders; +import org.springframework.integration.config.EnableIntegration; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.MessageHandler; +import org.springframework.messaging.MessageHandlingException; +import org.springframework.messaging.support.MessageBuilder; +import org.springframework.test.context.junit4.SpringRunner; + +import com.amazonaws.handlers.AsyncHandler; +import com.amazonaws.services.kinesis.AmazonKinesisAsync; +import com.amazonaws.services.kinesis.model.PutRecordRequest; +import com.amazonaws.services.kinesis.model.PutRecordResult; + +/** + * @author Artem Bilan + * @since 1.1 + */ +@RunWith(SpringRunner.class) +public class KinesisMessageHandlerTests { + + @Autowired + protected AmazonKinesisAsync amazonKinesis; + + @Autowired + protected MessageChannel kinesisSendChannel; + + @Autowired + protected KinesisMessageHandler kinesisMessageHandler; + + @Autowired + protected AsyncHandler asyncHandler; + + @Test + public void testKinesisMessageHandler() { + Message message = MessageBuilder.withPayload("message").build(); + try { + this.kinesisSendChannel.send(message); + } + catch (Exception e) { + assertThat(e).isInstanceOf(MessageHandlingException.class); + assertThat(e.getCause()).isInstanceOf(IllegalStateException.class); + assertThat(e.getMessage()).contains("'stream' must not be null for sending a Kinesis record"); + } + + this.kinesisMessageHandler.setStream("foo"); + try { + this.kinesisSendChannel.send(message); + } + catch (Exception e) { + assertThat(e).isInstanceOf(MessageHandlingException.class); + assertThat(e.getCause()).isInstanceOf(IllegalStateException.class); + assertThat(e.getMessage()).contains("'partitionKey' must not be null for sending a Kinesis record"); + } + + message = MessageBuilder.fromMessage(message) + .setHeader(AwsHeaders.PARTITION_KEY, "fooKey") + .setHeader(AwsHeaders.SEQUENCE_NUMBER, "10") + .build(); + + this.kinesisSendChannel.send(message); + + ArgumentCaptor putRecordRequestArgumentCaptor = + ArgumentCaptor.forClass(PutRecordRequest.class); + verify(this.amazonKinesis).putRecordAsync(putRecordRequestArgumentCaptor.capture(), eq(this.asyncHandler)); + + PutRecordRequest putRecordRequest = putRecordRequestArgumentCaptor.getValue(); + + assertThat(putRecordRequest.getStreamName()).isEqualTo("foo"); + assertThat(putRecordRequest.getPartitionKey()).isEqualTo("fooKey"); + assertThat(putRecordRequest.getSequenceNumberForOrdering()).isEqualTo("10"); + assertThat(putRecordRequest.getExplicitHashKey()).isNull(); + assertThat(putRecordRequest.getData()).isEqualTo(ByteBuffer.wrap("message".getBytes())); + } + + + @Configuration + @EnableIntegration + public static class ContextConfiguration { + + @Bean + @SuppressWarnings("unchecked") + public AmazonKinesisAsync amazonKinesis() { + AmazonKinesisAsync mock = mock(AmazonKinesisAsync.class); + given(mock.putRecordAsync(any(PutRecordRequest.class), any(AsyncHandler.class))) + .willReturn(mock(Future.class)); + return mock; + } + + @Bean + @SuppressWarnings("unchecked") + public AsyncHandler asyncHandler() { + return mock(AsyncHandler.class); + } + + @Bean + @ServiceActivator(inputChannel = "kinesisSendChannel") + public MessageHandler kinesisMessageHandler() { + KinesisMessageHandler kinesisMessageHandler = new KinesisMessageHandler(amazonKinesis()); + kinesisMessageHandler.setSync(true); + kinesisMessageHandler.setAsyncHandler(asyncHandler()); + kinesisMessageHandler.setConverter(new Converter() { + + private SerializingConverter serializingConverter = new SerializingConverter(); + + @Override + public byte[] convert(Object source) { + if (source instanceof String) { + return ((String) source).getBytes(); + } + else { + return this.serializingConverter.convert(source); + } + } + + }); + return kinesisMessageHandler; + } + + } + +}