INTEXT-154: Namespace Support for New Features

JIRA: https://jira.spring.io/browse/INTEXT-154, https://jira.spring.io/browse/INTEXT-138
This commit is contained in:
Gary Russell
2015-03-20 15:56:31 -04:00
committed by Artem Bilan
parent 3ade92ba14
commit 52d218988b
6 changed files with 251 additions and 10 deletions

View File

@@ -20,7 +20,6 @@ import org.w3c.dom.Element;
import org.springframework.beans.factory.support.AbstractBeanDefinition;
import org.springframework.beans.factory.support.BeanDefinitionBuilder;
import org.springframework.beans.factory.xml.AbstractSingleBeanDefinitionParser;
import org.springframework.beans.factory.xml.ParserContext;
import org.springframework.integration.config.xml.AbstractChannelAdapterParser;
import org.springframework.integration.config.xml.IntegrationNamespaceUtils;
@@ -93,6 +92,15 @@ public class KafkaMessageDrivenChannelAdapterParser extends AbstractChannelAdapt
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "key-decoder");
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "payload-decoder");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element,
"auto-commit", "autoCommitOffset");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element,
"use-context-message-builder", "useMessageBuilderFactory");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element,
"set-id", "generateMessageId");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element,
"set-timestamp", "generateTimestamp");
return builder.getBeanDefinition();
}

View File

@@ -73,6 +73,14 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSupport imp
this.payloadDecoder = payloadDecoder;
}
/**
* Automatically commit the offsets when 'true'. When 'false', the
* adapter inserts a 'kafka_acknowledgment` header allowing the user to manually
* commit the offset using the {@link Acknowledgment#acknowledge()} method.
* Default 'true'.
*
* @param autoCommitOffset false to not auto-commit (default true).
*/
public void setAutoCommitOffset(boolean autoCommitOffset) {
this.autoCommitOffset = autoCommitOffset;
}
@@ -80,7 +88,8 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSupport imp
/**
* Generate {@link Message} {@code ids} for produced messages.
* If set to {@code false}, will try to use a default value. By default set to {@code false}.
* Note that this option only works in conjunction with {@link #setUseMessageBuilderFactory(boolean)}.
* Note that this option is only guaranteed to work when
* {@link #setUseMessageBuilderFactory(boolean) useMessageBuilderFactory} is false (default).
* If the latter is set to {@code true}, then some {@link MessageBuilderFactory} implementations such as
* {@link DefaultMessageBuilderFactory} may ignore it.
* @param generateMessageId true if a message id should be generated
@@ -93,7 +102,8 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSupport imp
/**
* Generate {@code timestamp} for produced messages. If set to {@code false}, -1 is used instead.
* By default set to {@code false}.
* Note that this option only works in conjunction with {@link #setUseMessageBuilderFactory(boolean)}.
* Note that this option is only guaranteed to work when
* {@link #setUseMessageBuilderFactory(boolean) useMessageBuilderFactory} is false (default).
* If the latter is set to {@code true}, then some {@link MessageBuilderFactory} implementations such as
* {@link DefaultMessageBuilderFactory} may ignore it.
* @param generateTimestamp true if a timestamp should be generated
@@ -220,7 +230,7 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSupport imp
/**
* Special subclass of {@link Message}. It is used for lower message generation overhead, unless the default
* strategy of the superclass is set via {@link #setUseMessageBuilderFactory(boolean)}
* strategy of the outer class is set via {@link #setUseMessageBuilderFactory(boolean)}
* @since 1.1
*/
private class KafkaMessage implements Message<Object> {
@@ -244,6 +254,20 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSupport imp
return this.messageHeaders;
}
@Override
public String toString() {
StringBuilder sb = new StringBuilder(getClass().getSimpleName());
sb.append(" [payload=");
if (this.payload instanceof byte[]) {
sb.append("byte[").append(((byte[]) this.payload).length).append("]");
}
else {
sb.append(this.payload);
}
sb.append(", headers=").append(this.messageHeaders).append("]");
return sb.toString();
}
}
@SuppressWarnings("serial")

View File

@@ -29,9 +29,9 @@ public abstract class KafkaHeaders {
public static final String MESSAGE_KEY = PREFIX + "messageKey";
public static final String PARTITION_ID = PREFIX + "_partitionId";
public static final String PARTITION_ID = PREFIX + "partitionId";
public static final String OFFSET = PREFIX + "_offset";
public static final String OFFSET = PREFIX + "offset";
public static final String NEXT_OFFSET = PREFIX + "nextOffset";

View File

@@ -651,6 +651,48 @@
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="auto-commit" type="xsd:string">
<xsd:annotation>
<xsd:documentation>
When 'true', the adapter automatically commits the offset. When 'false', the adapter
inserts a 'kafka_acknowledgment` header allowing the user to manually commit the offset
using the `Acknowledgment.acknowledge()` method. Default 'true'.
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="use-context-message-builder" type="xsd:string">
<xsd:annotation>
<xsd:documentation>
When 'true', use the configured message builder from the integration application context
(bean:"messageBuilderFactory", which defaults to a "DefaultMessageBuilderFactory").
When 'false', use an internal, optimized, message creation algorithm. Default 'false'.
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="set-id" type="xsd:string">
<xsd:annotation>
<xsd:documentation>
When 'true', generate a unique ID header for each new message.
Default 'false'. Set to `true` if messages contain collections and you use a splitter
later in the flow, or a router where a message is sent to multiple recipients with
'apply-sequence' set to 'true'. In this cases, the message id is used as the correlation
id of the resulting messages. If 'use-context-message-builder' is 'true', this
setting may be ignored by the configured message builder factory (this is the case
for the 'DefaultMessageBuilderFactory', where an ID is always generated).
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="set-timestamp" type="xsd:string">
<xsd:annotation>
<xsd:documentation>
When 'true', generate a TIMESTAMP header for each new message.
Default 'false' (TIMESTAMP header is set to -1).
If 'use-context-message-builder' is 'true', this
setting may be ignored by the configured message builder factory (this is the case
for the 'DefaultMessageBuilderFactory', where a timestamp header is always generated).
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
</xsd:complexType>
</xsd:element>

View File

@@ -49,8 +49,71 @@
queue-size="${queue.size:1024}"
concurrency="${concurrency:10}"
max-fetch="${max.fetch:1000}"
topics="${topics:foo,bar}"/>
topics="${topics:foo,bar}" />
<int-kafka:message-driven-channel-adapter
id="withMBFactoryOverrideAndId"
auto-startup="false"
phase="100"
send-timeout="5000"
channel="nullChannel"
error-channel="errorChannel"
error-handler="errorHandler"
connection-factory="connectionFactory"
key-decoder="keyDecoder"
payload-decoder="payloadDecoder"
offset-manager="offsetManager"
task-executor="executor"
queue-size="${queue.size:1024}"
concurrency="${concurrency:10}"
max-fetch="${max.fetch:1000}"
topics="${topics:foo,bar}"
auto-commit="false"
use-context-message-builder="true"
set-id="true"
set-timestamp="false" />
<int-kafka:message-driven-channel-adapter
id="withMBFactoryOverrideAndTS"
auto-startup="false"
phase="100"
send-timeout="5000"
channel="nullChannel"
error-channel="errorChannel"
error-handler="errorHandler"
connection-factory="connectionFactory"
key-decoder="keyDecoder"
payload-decoder="payloadDecoder"
offset-manager="offsetManager"
task-executor="executor"
queue-size="${queue.size:1024}"
concurrency="${concurrency:10}"
max-fetch="${max.fetch:1000}"
topics="${topics:foo,bar}"
use-context-message-builder="true"
set-id="false"
set-timestamp="true" />
<int-kafka:message-driven-channel-adapter
id="withOverrideIdTS"
auto-startup="false"
phase="100"
send-timeout="5000"
channel="nullChannel"
error-channel="errorChannel"
error-handler="errorHandler"
connection-factory="connectionFactory"
key-decoder="keyDecoder"
payload-decoder="payloadDecoder"
offset-manager="offsetManager"
task-executor="executor"
queue-size="${queue.size:1024}"
concurrency="${concurrency:10}"
max-fetch="${max.fetch:1000}"
topics="${topics:foo,bar}"
use-context-message-builder="false"
set-id="true"
set-timestamp="true" />
<!-- Invalid config

View File

@@ -16,34 +16,51 @@
package org.springframework.integration.kafka.config.xml;
import static org.hamcrest.Matchers.equalTo;
import static org.junit.Assert.assertArrayEquals;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertSame;
import static org.junit.Assert.assertThat;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
import java.lang.reflect.Method;
import java.util.concurrent.Executor;
import java.util.concurrent.atomic.AtomicReference;
import org.hamcrest.Matchers;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.integration.channel.NullChannel;
import org.springframework.integration.channel.PublishSubscribeChannel;
import org.springframework.integration.kafka.core.ConnectionFactory;
import org.springframework.integration.kafka.core.KafkaMessageMetadata;
import org.springframework.integration.kafka.core.Partition;
import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter;
import org.springframework.integration.kafka.listener.Acknowledgment;
import org.springframework.integration.kafka.listener.ErrorHandler;
import org.springframework.integration.kafka.listener.KafkaMessageListenerContainer;
import org.springframework.integration.kafka.listener.OffsetManager;
import org.springframework.integration.kafka.support.KafkaHeaders;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.messaging.PollableChannel;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import org.springframework.util.ReflectionUtils;
import org.springframework.util.ReflectionUtils.MethodCallback;
import org.springframework.util.ReflectionUtils.MethodFilter;
import kafka.serializer.Decoder;
/**
* @author Artem Bilan.
* @author Gary Russell
*/
@RunWith(SpringJUnit4ClassRunner.class)
@ContextConfiguration
@@ -76,8 +93,17 @@ public class KafkaMessageDrivenChannelAdapterParserTests {
@Autowired
private KafkaMessageDrivenChannelAdapter kafkaListener;
@Autowired
private KafkaMessageDrivenChannelAdapter withMBFactoryOverrideAndId;
@Autowired
private KafkaMessageDrivenChannelAdapter withMBFactoryOverrideAndTS;
@Autowired
private KafkaMessageDrivenChannelAdapter withOverrideIdTS;
@Test
public void testKafkaMessageDrivenChannelAdapterParser() {
public void testKafkaMessageDrivenChannelAdapterParser() throws Exception {
assertFalse(this.kafkaListener.isAutoStartup());
assertFalse(this.kafkaListener.isRunning());
assertEquals(100, this.kafkaListener.getPhase());
@@ -97,6 +123,84 @@ public class KafkaMessageDrivenChannelAdapterParserTests {
assertEquals(1024, container.getQueueSize());
assertEquals(1024, container.getQueueSize());
assertArrayEquals(new String[] {"foo", "bar"}, TestUtils.getPropertyValue(container, "topics", String[].class));
assertOverrides(this.kafkaListener, false, false, false, true);
assertOverrides(this.withMBFactoryOverrideAndId, true, true, false, false);
assertOverrides(this.withMBFactoryOverrideAndTS, true, false, true, true);
assertOverrides(this.withOverrideIdTS, false, true, true, true);
final AtomicReference<Method> toMessage = new AtomicReference<Method>();
ReflectionUtils.doWithMethods(KafkaMessageDrivenChannelAdapter.class, new MethodCallback() {
@Override
public void doWith(Method method) throws IllegalArgumentException, IllegalAccessException {
method.setAccessible(true);
toMessage.set(method);
; }
},
new MethodFilter() {
@Override
public boolean matches(Method method) {
return method.getName().equals("toMessage");
}
});
Message<?> m = getAMessageFrom(this.kafkaListener, toMessage.get());
assertThat(m.getClass().getSimpleName(), equalTo("KafkaMessage"));
assertNull(m.getHeaders().getId());
assertNull(m.getHeaders().getTimestamp());
assertNull(m.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT));
assertRest(m);
m = getAMessageFrom(this.withMBFactoryOverrideAndId, toMessage.get());
assertThat(m, Matchers.instanceOf(GenericMessage.class));
assertNotNull(m.getHeaders().getId());
assertNotNull(m.getHeaders().getTimestamp());
assertNotNull(m.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT));
assertRest(m);
m = getAMessageFrom(this.withMBFactoryOverrideAndTS, toMessage.get());
assertThat(m, Matchers.instanceOf(GenericMessage.class));
assertNotNull(m.getHeaders().getId());
assertNotNull(m.getHeaders().getTimestamp());
assertNull(m.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT));
assertRest(m);
m = getAMessageFrom(this.withOverrideIdTS, toMessage.get());
assertThat(m.getClass().getSimpleName(), equalTo("KafkaMessage"));
assertNotNull(m.getHeaders().getId());
assertNotNull(m.getHeaders().getTimestamp());
assertNull(m.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT));
assertRest(m);
}
private void assertOverrides(KafkaMessageDrivenChannelAdapter kafkaListener, boolean mbf, boolean id, boolean ts,
boolean ac) {
assertThat(TestUtils.getPropertyValue(kafkaListener, "autoCommitOffset", Boolean.class), equalTo(ac));
assertThat(TestUtils.getPropertyValue(kafkaListener, "useMessageBuilderFactory", Boolean.class), equalTo(mbf));
assertThat(TestUtils.getPropertyValue(kafkaListener, "generateMessageId", Boolean.class), equalTo(id));
assertThat(TestUtils.getPropertyValue(kafkaListener, "generateTimestamp", Boolean.class), equalTo(ts));
}
private void assertRest(Message<?> m) {
assertEquals("bar", m.getPayload());
assertEquals("foo", m.getHeaders().get(KafkaHeaders.MESSAGE_KEY));
assertEquals("topic", m.getHeaders().get(KafkaHeaders.TOPIC));
assertEquals(42, m.getHeaders().get(KafkaHeaders.PARTITION_ID));
assertEquals(1L, m.getHeaders().get(KafkaHeaders.OFFSET));
assertEquals(2L, m.getHeaders().get(KafkaHeaders.NEXT_OFFSET));
}
private Message<?> getAMessageFrom(KafkaMessageDrivenChannelAdapter adapter, Method toMessage) throws Exception {
KafkaMessageMetadata meta = mock(KafkaMessageMetadata.class);
Partition partition = mock(Partition.class);
when(partition.getTopic()).thenReturn("topic");
when(partition.getId()).thenReturn(42);
when(meta.getPartition()).thenReturn(partition);
when(meta.getOffset()).thenReturn(1L);
when(meta.getNextOffset()).thenReturn(2L);
Acknowledgment ack = mock(Acknowledgment.class);
return (Message<?>) toMessage.invoke(adapter, "foo", "bar", meta, ack);
}
}