Kinesis-Binder-29: Add RECEIVED headers
Fixes: spring-cloud/spring-cloud-stream-binder-aws-kinesis#29 Previously the same header name has been used for sending and receiving operations (e.g. Kinesis `stream`). This causes collisions in streaming processes when we receive message from the AWS and send it downstream to AWS. The header presence has a precedence over configured property/expression. Therefore we send message to the AWS (e.g. Kinesis) using wrong destination or other correlation properties * Add `AwsHeaders.RECEIVED_*` headers to avoid collisions * Remove Jackson dependency since it is managed now properly by SC-AWS
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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();
|
||||
|
||||
|
||||
@@ -901,10 +901,10 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i
|
||||
}
|
||||
AbstractIntegrationMessageBuilder<Object> 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);
|
||||
|
||||
@@ -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";
|
||||
|
||||
|
||||
@@ -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();
|
||||
|
||||
|
||||
Reference in New Issue
Block a user