diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/log4j2/AmqpAppender.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/log4j2/AmqpAppender.java index 4df572a5..9bb8e8dc 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/log4j2/AmqpAppender.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/log4j2/AmqpAppender.java @@ -75,6 +75,7 @@ import com.rabbitmq.client.ConnectionFactory; * @author Stephen Oakey * @author Artem Bilan * @author Nicolas Ristock + * @author Eugene Gusev * * @since 1.6 */ @@ -156,7 +157,8 @@ public class AmqpAppender extends AbstractAppender { @PluginAttribute("contentEncoding") String contentEncoding, @PluginAttribute("clientConnectionProperties") String clientConnectionProperties, @PluginAttribute("async") boolean async, - @PluginAttribute("charset") String charset) { + @PluginAttribute("charset") String charset, + @PluginAttribute(value = "addMdcAsHeaders", defaultBoolean = true) boolean addMdcAsHeaders) { if (name == null) { LOGGER.error("No name for AmqpAppender"); } @@ -187,6 +189,7 @@ public class AmqpAppender extends AbstractAppender { manager.clientConnectionProperties = clientConnectionProperties; manager.charset = charset; manager.async = async; + manager.addMdcAsHeaders = addMdcAsHeaders; AmqpAppender appender = new AmqpAppender(name, filter, theLayout, ignoreExceptions, manager); if (manager.activateOptions()) { appender.startSenders(); @@ -235,7 +238,7 @@ public class AmqpAppender extends AbstractAppender { return message; } - private void sendEvent(final Event event, Map properties) { + protected void sendEvent(Event event, Map properties) { LogEvent logEvent = event.getEvent(); String name = logEvent.getLoggerName(); Level level = logEvent.getLevel(); @@ -264,8 +267,10 @@ public class AmqpAppender extends AbstractAppender { amqpProps.setTimestamp(tstamp.getTime()); // Copy properties in from MDC - for (Entry entry : properties.entrySet()) { - amqpProps.setHeader(entry.getKey().toString(), entry.getValue()); + if (this.manager.addMdcAsHeaders) { + for (Entry entry : properties.entrySet()) { + amqpProps.setHeader(entry.getKey().toString(), entry.getValue()); + } } if (logEvent.getSource() != null) { amqpProps.setHeader( @@ -275,6 +280,10 @@ public class AmqpAppender extends AbstractAppender { logEvent.getSource().getLineNumber())); } + doSend(event, logEvent, amqpProps); + } + + protected void doSend(final Event event, LogEvent logEvent, MessageProperties amqpProps) { StringBuilder msgBody; String routingKey; @@ -491,6 +500,11 @@ public class AmqpAppender extends AbstractAppender { */ private String charset = Charset.defaultCharset().name(); + /** + * Whether or not add MDC properties into message headers. true by default for backward compatibility + */ + private boolean addMdcAsHeaders = true; + private boolean durable = true; private MessageDeliveryMode deliveryMode = MessageDeliveryMode.PERSISTENT; diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/logback/AmqpAppender.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/logback/AmqpAppender.java index 253effa4..a3f02d9d 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/logback/AmqpAppender.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/logback/AmqpAppender.java @@ -77,6 +77,7 @@ import ch.qos.logback.core.Layout; * @author Artem Bilan * @author Gary Russell * @author Nicolas Ristock + * @author Eugene Gusev * * @since 1.4 */ @@ -222,6 +223,11 @@ public class AmqpAppender extends AppenderBase { */ private String charset; + /** + * Whether or not add MDC properties into message headers. true by default for backward compatibility + */ + private boolean addMdcAsHeaders = true; + private boolean durable = true; private MessageDeliveryMode deliveryMode = MessageDeliveryMode.PERSISTENT; @@ -359,6 +365,14 @@ public class AmqpAppender extends AppenderBase { this.maxSenderRetries = maxSenderRetries; } + public boolean isAddMdcAsHeaders() { + return this.addMdcAsHeaders; + } + + public void setAddMdcAsHeaders(boolean addMdcAsHeaders) { + this.addMdcAsHeaders = addMdcAsHeaders; + } + public boolean isDurable() { return this.durable; } @@ -579,10 +593,12 @@ public class AmqpAppender extends AppenderBase { amqpProps.setTimestamp(tstamp.getTime()); // Copy properties in from MDC - Map props = event.getProperties(); - Set> entrySet = props.entrySet(); - for (Entry entry : entrySet) { - amqpProps.setHeader(entry.getKey(), entry.getValue()); + if (AmqpAppender.this.addMdcAsHeaders) { + Map props = event.getProperties(); + Set> entrySet = props.entrySet(); + for (Entry entry : entrySet) { + amqpProps.setHeader(entry.getKey(), entry.getValue()); + } } String[] location = AmqpAppender.this.locationLayout.doLayout(logEvent).split("\\|"); if (!"?".equals(location[0])) { @@ -629,6 +645,7 @@ public class AmqpAppender extends AppenderBase { if (retries < AmqpAppender.this.maxSenderRetries) { // Schedule a retry based on the number of times I've tried to re-send this AmqpAppender.this.retryTimer.schedule(new TimerTask() { + @Override public void run() { AmqpAppender.this.events.add(event); @@ -646,6 +663,7 @@ public class AmqpAppender extends AppenderBase { Thread.currentThread().interrupt(); } } + } /** diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/log4j/AmqpAppenderConfiguration.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/log4j/AmqpAppenderConfiguration.java index e98c5b3a..f62796a7 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/log4j/AmqpAppenderConfiguration.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/log4j/AmqpAppenderConfiguration.java @@ -25,6 +25,7 @@ import org.springframework.amqp.core.Queue; import org.springframework.amqp.core.TopicExchange; import org.springframework.amqp.rabbit.connection.SingleConnectionFactory; import org.springframework.amqp.rabbit.core.RabbitAdmin; +import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.context.annotation.Bean; @@ -33,6 +34,7 @@ import org.springframework.context.annotation.Scope; /** * @author Jon Brisbin + * @author Artem Bilan */ @Configuration public class AmqpAppenderConfiguration { @@ -104,4 +106,12 @@ public class AmqpAppenderConfiguration { public TestListener testListener(int count) { return new TestListener(count); } + + @Bean + public RabbitTemplate rabbitTemplate() { + RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory()); + rabbitTemplate.setReceiveTimeout(10_000); + return rabbitTemplate; + } + } diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/logback/AmqpAppenderIntegrationTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/logback/AmqpAppenderIntegrationTests.java index 0ee3ee6e..0be16942 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/logback/AmqpAppenderIntegrationTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/logback/AmqpAppenderIntegrationTests.java @@ -16,6 +16,7 @@ package org.springframework.amqp.rabbit.logback; +import static org.hamcrest.Matchers.hasEntry; import static org.hamcrest.Matchers.instanceOf; import static org.hamcrest.Matchers.is; import static org.hamcrest.Matchers.startsWith; @@ -37,6 +38,9 @@ import org.slf4j.MDC; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageProperties; +import org.springframework.amqp.core.Queue; +import org.springframework.amqp.rabbit.connection.SingleConnectionFactory; +import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.rabbit.junit.BrokerRunning; import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; import org.springframework.amqp.rabbit.log4j.AmqpAppenderConfiguration; @@ -52,6 +56,7 @@ import ch.qos.logback.classic.Logger; /** * @author Artem Bilan * @author Nicolas Ristock + * @author Eugene Gusev * * @since 1.4 */ @@ -69,15 +74,23 @@ public class AmqpAppenderIntegrationTests { @Autowired private ApplicationContext applicationContext; + @Autowired + private RabbitTemplate template; + + @Autowired + private Queue testQueue; + private SimpleMessageListenerContainer listenerContainer; @Before - public void setUp() throws Exception { - listenerContainer = applicationContext.getBean(SimpleMessageListenerContainer.class); + public void setUp() { + this.listenerContainer = this.applicationContext.getBean(SimpleMessageListenerContainer.class); + MDC.clear(); } @After public void tearDown() { + MDC.clear(); listenerContainer.shutdown(); } @@ -120,7 +133,9 @@ public class AmqpAppenderIntegrationTests { assertNotNull(location); assertThat(location, instanceOf(String.class)); assertThat((String) location, - startsWith("org.springframework.amqp.rabbit.logback.AmqpAppenderIntegrationTests.testAppenderWithProps()")); + startsWith("org.springframework.amqp.rabbit.logback.AmqpAppenderIntegrationTests" + + ".testAppenderWithProps" + + "()")); Object threadName = messageProperties.getHeaders().get("thread"); assertNotNull(threadName); assertThat(threadName, instanceOf(String.class)); @@ -144,6 +159,32 @@ public class AmqpAppenderIntegrationTests { assertEquals(0xbf, body[body.length - 3 - lineSeparatorExtraBytes] & 0xff); } + @Test + public void testAddMdcAsHeaders() { + this.applicationContext.getBean(SingleConnectionFactory.class).createConnection().close(); + + Logger logWithMdc = (Logger) LoggerFactory.getLogger("withMdc"); + Logger logWithoutMdc = (Logger) LoggerFactory.getLogger("withoutMdc"); + MDC.put("mdc1", "test1"); + MDC.put("mdc2", "test2"); + + logWithMdc.info("test message with MDC in headers"); + Message received1 = this.template.receive(this.testQueue.getName()); + + assertNotNull(received1); + assertEquals("test message with MDC in headers", new String(received1.getBody())); + assertThat(received1.getMessageProperties().getHeaders(), hasEntry("mdc1", "test1")); + assertThat(received1.getMessageProperties().getHeaders(), hasEntry("mdc2", "test2")); + + logWithoutMdc.info("test message without MDC in headers"); + Message received2 = this.template.receive(this.testQueue.getName()); + + assertNotNull(received2); + assertEquals("test message without MDC in headers", new String(received2.getBody())); + assertThat(received1.getMessageProperties().getHeaders(), hasEntry("mdc1", "test1")); + assertThat(received1.getMessageProperties().getHeaders(), hasEntry("mdc2", "test2")); + } + public static class EnhancedAppender extends AmqpAppender { private String foo; diff --git a/spring-rabbit/src/test/resources/logback-test.xml b/spring-rabbit/src/test/resources/logback-test.xml index 156e460f..d2791741 100644 --- a/spring-rabbit/src/test/resources/logback-test.xml +++ b/spring-rabbit/src/test/resources/logback-test.xml @@ -25,10 +25,52 @@ bar + + + %m + + localhost:5672 + 36 + true + AmqpAppenderTest + %property{applicationId}.%c.%p + true + UTF-8 + false + NON_PERSISTENT + false + true + + + + + %m + + localhost:5672 + 36 + true + AmqpAppenderTest + %property{applicationId}.%c.%p + true + UTF-8 + false + NON_PERSISTENT + false + false + + + + + + + + + + diff --git a/src/reference/asciidoc/logging.adoc b/src/reference/asciidoc/logging.adoc index 6e0f2114..d13552c5 100644 --- a/src/reference/asciidoc/logging.adoc +++ b/src/reference/asciidoc/logging.adoc @@ -7,8 +7,7 @@ The framework provides logging appenders for several popular logging subsystems: - logback (since Spring AMQP _version 1.4_) - log4j2 (since Spring AMQP _version 1.6_) -The appenders are configured using the normal mechanisms for the logging subsystem, available properties are specified -in the following sections. +The appenders are configured by using the normal mechanisms for the logging subsystem, available properties are specified in the following sections. ==== Common properties @@ -109,6 +108,15 @@ If the charset is unsupported on the current platform, we fall back to using the | null | A comma-delimited list of `key:value` pairs for custom client properties to the RabbitMQ connection. +| addMdcAsHeaders +| true +| MDC properties were always added into RabbitMQ message headers until this property was introduced. +It can lead to issues for big MDC as while RabbitMQ has limited buffer size for all headers and this buffer is pretty small. +This property was introduced to avoid issues in cases of big MDC. +By default this value set to `true` for backward compatibility. +The `false` turns off serialization MDC into headers. +Please note, the `JsonLayout` adds MDC into the message by default. + |=== ==== Log4j Appender @@ -144,7 +152,8 @@ NOTE: This appender is deprecated and will be removed in _version 2.0_. applicationId="myAppId" routingKeyPattern="%X{applicationId}.%c.%p" contentType="text/plain" contentEncoding="UTF-8" generateId="true" deliveryMode="NON_PERSISTENT" charset="UTF-8" - senderPoolSize="3" maxSenderRetries="5"> + senderPoolSize="3" maxSenderRetries="5" + addMdcAsHeaders="false"> ---- @@ -179,6 +188,7 @@ If async publishing is used with the `ReusableLogEventFactory`, events will have false NON_PERSISTENT true + false ----