INT-3586: Map Consumer Metadata Headers

JIRA: https://jira.spring.io/browse/INT-3586

Map the consumer tag and queue `MessageProperties` to
Spring Integration headers.

Use reflection to detect presence of Spring AMQP 1.4.2, to avoid forcing a new Spring AMQP version.

Polishing
This commit is contained in:
Gary Russell
2015-01-08 19:50:14 +02:00
committed by Artem Bilan
parent b741b86f7e
commit c24a12d19b
4 changed files with 63 additions and 3 deletions

View File

@@ -121,14 +121,14 @@ subprojects { subproject ->
tomcatVersion = "7.0.55"
smack3Version = '3.2.1'
smackVersion = '4.0.0'
springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '1.4.1.RELEASE'
springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '1.4.2.RELEASE'
springDataMongoVersion = '1.6.0.RELEASE'
springDataRedisVersion = '1.4.0.RELEASE'
springGemfireVersion = '1.5.0.RELEASE'
springSecurityVersion = '3.2.5.RELEASE'
springSocialTwitterVersion = '1.1.0.RELEASE'
springRetryVersion = '1.1.1.RELEASE'
springVersion = project.hasProperty('springVersion') ? project.springVersion : '4.1.3.RELEASE'
springVersion = project.hasProperty('springVersion') ? project.springVersion : '4.1.4.RELEASE'
springWsVersion = '2.2.0.RELEASE'
xmlUnitVersion = '1.5'
xstreamVersion = '1.4.7'

View File

@@ -16,11 +16,13 @@
package org.springframework.integration.amqp.support;
import java.lang.reflect.Field;
import java.util.ArrayList;
import java.util.Date;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicBoolean;
import org.springframework.amqp.core.MessageDeliveryMode;
import org.springframework.amqp.core.MessageProperties;
@@ -28,6 +30,9 @@ import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.integration.IntegrationMessageHeaderAccessor;
import org.springframework.integration.mapping.AbstractHeaderMapper;
import org.springframework.integration.mapping.support.JsonHeaders;
import org.springframework.util.ReflectionUtils;
import org.springframework.util.ReflectionUtils.FieldCallback;
import org.springframework.util.ReflectionUtils.FieldFilter;
import org.springframework.util.StringUtils;
/**
@@ -53,6 +58,8 @@ import org.springframework.util.StringUtils;
*/
public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessageProperties> implements AmqpHeaderMapper {
static final boolean CONSUMER_METADATA_PRESENT;
private static final List<String> STANDARD_HEADER_NAMES = new ArrayList<String>();
static {
@@ -79,6 +86,27 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessagePropert
STANDARD_HEADER_NAMES.add(JsonHeaders.KEY_TYPE_ID);
STANDARD_HEADER_NAMES.add(AmqpHeaders.SPRING_REPLY_CORRELATION);
STANDARD_HEADER_NAMES.add(AmqpHeaders.SPRING_REPLY_TO_STACK);
final AtomicBoolean consumerTagHeader = new AtomicBoolean();
try {
ReflectionUtils.doWithFields(AmqpHeaders.class, new FieldCallback() {
@Override
public void doWith(Field field) throws IllegalArgumentException, IllegalAccessException {
consumerTagHeader.set(true);
}
},
new FieldFilter() {
@Override
public boolean matches(Field field) {
return field.getName().equals("CONSUMER_TAG") && field.getType().equals(String.class);
}
});
}
catch (Exception e) {
}
CONSUMER_METADATA_PRESENT = consumerTagHeader.get();
}
public DefaultAmqpHeaderMapper() {
@@ -347,4 +375,24 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessagePropert
return contentTypeStringValue;
}
@Override
public Map<String, Object> toHeadersFromRequest(MessageProperties source) {
Map<String, Object> headersFromRequest = super.toHeadersFromRequest(source);
if (CONSUMER_METADATA_PRESENT) {
addConsumerMetadata(source, headersFromRequest);
}
return headersFromRequest;
}
private void addConsumerMetadata(MessageProperties messageProperties, Map<String, Object> headers) {
String consumerTag = messageProperties.getConsumerTag();
if (consumerTag != null) {
headers.put(AmqpHeaders.CONSUMER_TAG, consumerTag);
}
String consumerQueue = messageProperties.getConsumerQueue();
if (consumerQueue != null) {
headers.put(AmqpHeaders.CONSUMER_QUEUE, consumerQueue);
}
}
}

View File

@@ -171,6 +171,17 @@ public class DefaultAmqpHeaderMapperTests {
assertEquals("test.replyTo2", headerMap.get(AmqpHeaders.SPRING_REPLY_TO_STACK));
}
@Test // INT-3586 requires Spring AMQP 1.4.2
public void testToHeadersConsumerMetadata() {
DefaultAmqpHeaderMapper headerMapper = new DefaultAmqpHeaderMapper();
MessageProperties amqpProperties = new MessageProperties();
amqpProperties.setConsumerTag("consumerTag");
amqpProperties.setConsumerQueue("consumerQueue");
Map<String, Object> headerMap = headerMapper.toHeadersFromRequest(amqpProperties);
assertEquals("consumerTag", headerMap.get(AmqpHeaders.CONSUMER_TAG));
assertEquals("consumerQueue", headerMap.get(AmqpHeaders.CONSUMER_QUEUE));
}
@Test
public void messageIdNotMappedToAmqpProperties() {
DefaultAmqpHeaderMapper headerMapper = new DefaultAmqpHeaderMapper();

View File

@@ -55,8 +55,9 @@
<year>2012</year>
<year>2013</year>
<year>2014</year>
<year>2015</year>
<holder>
GoPivotal, Inc. All Rights Reserved.
Pivotal, Inc. All Rights Reserved.
</holder>
</copyright>
</bookinfo>