diff --git a/build.gradle b/build.gradle index ab249a5..90ceba6 100644 --- a/build.gradle +++ b/build.gradle @@ -32,7 +32,6 @@ repositories { ext { assertjVersion = '3.8.0' - jacksonVersion = '2.9.3' servletApiVersion = '3.1.0' slf4jVersion = '1.7.25' springCloudAwsVersion = '2.0.0.BUILD-SNAPSHOT' @@ -99,7 +98,6 @@ dependencies { compile 'org.springframework.integration:spring-integration-core' compile 'org.springframework.cloud:spring-cloud-aws-core' - compile("com.fasterxml.jackson.core:jackson-databind:$jacksonVersion", optional) compile('org.springframework.cloud:spring-cloud-aws-messaging', optional) compile('org.springframework.integration:spring-integration-file', optional) compile('org.springframework.integration:spring-integration-http', optional) diff --git a/src/main/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapter.java b/src/main/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapter.java index 4ed5dac..0165d69 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapter.java +++ b/src/main/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapter.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * Copyright 2016-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. @@ -161,7 +161,7 @@ public class SqsMessageDrivenChannelAdapter extends MessageProducerSupport imple } @Override - public void destroy() throws Exception { + public void destroy() { this.listenerContainer.destroy(); } @@ -187,7 +187,7 @@ public class SqsMessageDrivenChannelAdapter extends MessageProducerSupport imple "Acknowledgment") .setHeader(AwsHeaders.MESSAGE_ID, headers.get("MessageId")) .setHeader(AwsHeaders.RECEIPT_HANDLE, headers.get("ReceiptHandle")) - .setHeader(AwsHeaders.QUEUE, headers.get("LogicalResourceId")) + .setHeader(AwsHeaders.RECEIVED_QUEUE, headers.get("LogicalResourceId")) .setHeader(AwsHeaders.ACKNOWLEDGMENT, headers.get("Acknowledgment")) .build(); diff --git a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java index 02d77ac..edbd1fa 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java +++ b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java @@ -901,10 +901,10 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i } AbstractIntegrationMessageBuilder messageBuilder = getMessageBuilderFactory() .withPayload(payload) - .setHeader(AwsHeaders.STREAM, this.shardOffset.getStream()) + .setHeader(AwsHeaders.RECEIVED_STREAM, this.shardOffset.getStream()) .setHeader(AwsHeaders.SHARD, this.shardOffset.getShard()) - .setHeader(AwsHeaders.PARTITION_KEY, record.getPartitionKey()) - .setHeader(AwsHeaders.SEQUENCE_NUMBER, record.getSequenceNumber()); + .setHeader(AwsHeaders.RECEIVED_PARTITION_KEY, record.getPartitionKey()) + .setHeader(AwsHeaders.RECEIVED_SEQUENCE_NUMBER, record.getSequenceNumber()); if (CheckpointMode.manual.equals(KinesisMessageDrivenChannelAdapter.this.checkpointMode)) { messageBuilder.setHeader(AwsHeaders.CHECKPOINTER, this.checkpointer); } @@ -921,7 +921,7 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i case batch: AbstractIntegrationMessageBuilder messageBuilder = getMessageBuilderFactory() .withPayload(records) - .setHeader(AwsHeaders.STREAM, this.shardOffset.getStream()) + .setHeader(AwsHeaders.RECEIVED_STREAM, this.shardOffset.getStream()) .setHeader(AwsHeaders.SHARD, this.shardOffset.getShard()); if (CheckpointMode.manual.equals(KinesisMessageDrivenChannelAdapter.this.checkpointMode)) { messageBuilder.setHeader(AwsHeaders.CHECKPOINTER, this.checkpointer); 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 364825d..b541cd0 100644 --- a/src/main/java/org/springframework/integration/aws/support/AwsHeaders.java +++ b/src/main/java/org/springframework/integration/aws/support/AwsHeaders.java @@ -26,10 +26,15 @@ public abstract class AwsHeaders { private static final String PREFIX = "aws_"; /** - * The {@value QUEUE} header for sending/receiving data over SQS. + * The {@value QUEUE} header for sending data to SQS. */ public static final String QUEUE = PREFIX + "queue"; + /** + * The {@value RECEIVED_QUEUE} header for receiving data from SQS. + */ + public static final String RECEIVED_QUEUE = PREFIX + "receivedQueue"; + /** * The {@value TOPIC} header for sending/receiving data over SNS. */ @@ -65,23 +70,38 @@ 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 SHARD} header to represent Kinesis shardId. */ public static final String SHARD = PREFIX + "shard"; /** - * The {@value PARTITION_KEY} header for sending/receiving data over Kinesis. + * The {@value RECEIVED_STREAM} header for receiving data from Kinesis. + */ + public static final String RECEIVED_STREAM = PREFIX + "receivedStream"; + + /** + * The {@value RECEIVED_PARTITION_KEY} header for receiving data from Kinesis. + */ + public static final String RECEIVED_PARTITION_KEY = PREFIX + "receivedPartitionKey"; + + /** + * The {@value RECEIVED_SEQUENCE_NUMBER} header for receiving data from Kinesis. + */ + public static final String RECEIVED_SEQUENCE_NUMBER = PREFIX + "receivedSequenceNumber"; + + /** + * The {@value STREAM} header for sending data to Kinesis. + */ + public static final String STREAM = PREFIX + "stream"; + + /** + * The {@value PARTITION_KEY} header for sending data to Kinesis. */ public static final String PARTITION_KEY = PREFIX + "partitionKey"; /** - * The {@value SEQUENCE_NUMBER} header for sending/receiving data over Kinesis. + * The {@value SEQUENCE_NUMBER} header for sending data to Kinesis. */ public static final String SEQUENCE_NUMBER = PREFIX + "sequenceNumber"; diff --git a/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java b/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java index fda329b..4412887 100644 --- a/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java +++ b/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java @@ -126,10 +126,10 @@ public class KinesisMessageDrivenChannelAdapterTests { assertThat(message).isNotNull(); assertThat(message.getPayload()).isEqualTo("foo"); MessageHeaders headers = message.getHeaders(); - assertThat(headers.get(AwsHeaders.PARTITION_KEY)).isEqualTo("partition1"); + assertThat(headers.get(AwsHeaders.RECEIVED_PARTITION_KEY)).isEqualTo("partition1"); assertThat(headers.get(AwsHeaders.SHARD)).isEqualTo("1"); - assertThat(headers.get(AwsHeaders.SEQUENCE_NUMBER)).isEqualTo("1"); - assertThat(headers.get(AwsHeaders.STREAM)).isEqualTo(STREAM1); + assertThat(headers.get(AwsHeaders.RECEIVED_SEQUENCE_NUMBER)).isEqualTo("1"); + assertThat(headers.get(AwsHeaders.RECEIVED_STREAM)).isEqualTo(STREAM1); Checkpointer checkpointer = headers.get(AwsHeaders.CHECKPOINTER, Checkpointer.class); assertThat(checkpointer).isNotNull(); @@ -139,10 +139,10 @@ public class KinesisMessageDrivenChannelAdapterTests { assertThat(message).isNotNull(); assertThat(message.getPayload()).isEqualTo("bar"); headers = message.getHeaders(); - assertThat(headers.get(AwsHeaders.PARTITION_KEY)).isEqualTo("partition1"); + assertThat(headers.get(AwsHeaders.RECEIVED_PARTITION_KEY)).isEqualTo("partition1"); assertThat(headers.get(AwsHeaders.SHARD)).isEqualTo("1"); - assertThat(headers.get(AwsHeaders.SEQUENCE_NUMBER)).isEqualTo("2"); - assertThat(headers.get(AwsHeaders.STREAM)).isEqualTo(STREAM1); + assertThat(headers.get(AwsHeaders.RECEIVED_SEQUENCE_NUMBER)).isEqualTo("2"); + assertThat(headers.get(AwsHeaders.RECEIVED_STREAM)).isEqualTo(STREAM1); assertThat(this.kinesisChannel.receive(10)).isNull();