INT-1268 renamed <publisher> to <scheduled-producer>
This commit is contained in:
@@ -59,7 +59,7 @@ public class IntegrationNamespaceHandler extends AbstractIntegrationNamespaceHan
|
||||
registerBeanDefinitionParser("poller", new PollerParser());
|
||||
registerBeanDefinitionParser("annotation-config", new AnnotationConfigParser());
|
||||
registerBeanDefinitionParser("application-event-multicaster", new ApplicationEventMulticasterParser());
|
||||
registerBeanDefinitionParser("publisher", new PublisherParser());
|
||||
registerBeanDefinitionParser("scheduled-producer", new ScheduledProducerParser());
|
||||
registerBeanDefinitionParser("publishing-interceptor", new PublishingInterceptorParser());
|
||||
registerBeanDefinitionParser("channel-interceptor", new GlobalChannelInterceptorParser());
|
||||
registerBeanDefinitionParser("converter", new ConverterParser());
|
||||
|
||||
@@ -25,16 +25,16 @@ import org.springframework.beans.factory.xml.ParserContext;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
/**
|
||||
* Parser for the <publisher> element.
|
||||
* Parser for the <scheduled-producer> element.
|
||||
*
|
||||
* @author Mark Fisher
|
||||
* @since 2.0
|
||||
*/
|
||||
public class PublisherParser extends AbstractSingleBeanDefinitionParser {
|
||||
public class ScheduledProducerParser extends AbstractSingleBeanDefinitionParser {
|
||||
|
||||
@Override
|
||||
protected String getBeanClassName(Element element) {
|
||||
return IntegrationNamespaceUtils.BASE_PACKAGE + ".endpoint.TriggeredMessagePublisher";
|
||||
return IntegrationNamespaceUtils.BASE_PACKAGE + ".endpoint.ScheduledMessageProducer";
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -81,7 +81,7 @@ public class PublisherParser extends AbstractSingleBeanDefinitionParser {
|
||||
return;
|
||||
}
|
||||
builder.addPropertyReference("outputChannel", element.getAttribute("channel"));
|
||||
builder.addConstructorArgValue(element.getAttribute("payload"));
|
||||
builder.addConstructorArgValue(element.getAttribute("payload-expression"));
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "auto-startup");
|
||||
// TODO: add support for header expression sub-elements
|
||||
}
|
||||
@@ -37,14 +37,14 @@ import org.springframework.util.CollectionUtils;
|
||||
* @author Mark Fisher
|
||||
* @since 2.0
|
||||
*/
|
||||
public class TriggeredMessagePublisher extends MessageProducerSupport {
|
||||
public class ScheduledMessageProducer extends MessageProducerSupport {
|
||||
|
||||
private static final ExpressionParser PARSER = new SpelExpressionParser();
|
||||
|
||||
|
||||
private final Trigger trigger;
|
||||
|
||||
private final MessagePublishingTask task;
|
||||
private final MessageProducingTask task;
|
||||
|
||||
private volatile ScheduledFuture<?> future;
|
||||
|
||||
@@ -53,11 +53,11 @@ public class TriggeredMessagePublisher extends MessageProducerSupport {
|
||||
private final StandardEvaluationContext context = new StandardEvaluationContext();
|
||||
|
||||
|
||||
public TriggeredMessagePublisher(Trigger trigger, String payloadExpression) {
|
||||
public ScheduledMessageProducer(Trigger trigger, String payloadExpression) {
|
||||
Assert.notNull(trigger, "trigger must not be null");
|
||||
Assert.hasText(payloadExpression, "payloadExpression is required");
|
||||
this.trigger = trigger;
|
||||
this.task = new MessagePublishingTask(PARSER.parseExpression(payloadExpression));
|
||||
this.task = new MessageProducingTask(PARSER.parseExpression(payloadExpression));
|
||||
}
|
||||
|
||||
|
||||
@@ -105,12 +105,12 @@ public class TriggeredMessagePublisher extends MessageProducerSupport {
|
||||
}
|
||||
|
||||
|
||||
private class MessagePublishingTask implements Runnable {
|
||||
private class MessageProducingTask implements Runnable {
|
||||
|
||||
private final Expression payloadExpression;
|
||||
|
||||
|
||||
private MessagePublishingTask(Expression payloadExpression) {
|
||||
private MessageProducingTask(Expression payloadExpression) {
|
||||
this.payloadExpression = payloadExpression;
|
||||
}
|
||||
|
||||
@@ -2274,7 +2274,7 @@ Name of the header whose value to use.
|
||||
</xsd:complexContent>
|
||||
</xsd:complexType>
|
||||
|
||||
<xsd:element name="publisher">
|
||||
<xsd:element name="scheduled-producer">
|
||||
<xsd:annotation>
|
||||
<xsd:documentation>
|
||||
Defines a component that evaluates an expression to generate a Message payload
|
||||
@@ -2313,7 +2313,7 @@ Name of the header whose value to use.
|
||||
</xsd:appinfo>
|
||||
</xsd:annotation>
|
||||
</xsd:attribute>
|
||||
<xsd:attribute name="payload" type="xsd:string" use="required">
|
||||
<xsd:attribute name="payload-expression" type="xsd:string" use="required">
|
||||
<xsd:annotation>
|
||||
<xsd:documentation>
|
||||
SpEL expression to be evaluated for each triggered execution.
|
||||
@@ -2325,7 +2325,7 @@ Name of the header whose value to use.
|
||||
<xsd:attribute name="channel" type="xsd:string" use="required">
|
||||
<xsd:annotation>
|
||||
<xsd:documentation>
|
||||
MessageChannel to which this publisher's output should be sent.
|
||||
MessageChannel to which this producer's output should be sent.
|
||||
</xsd:documentation>
|
||||
<xsd:appinfo>
|
||||
<tool:annotation>
|
||||
@@ -2337,11 +2337,12 @@ Name of the header whose value to use.
|
||||
<xsd:attribute name="auto-startup" type="xsd:string">
|
||||
<xsd:annotation>
|
||||
<xsd:documentation>
|
||||
Specify whether this publisher should start automatically.
|
||||
Specify whether this producer should start automatically.
|
||||
By default it will. Set this to 'false' to require a manual start.
|
||||
</xsd:documentation>
|
||||
</xsd:annotation>
|
||||
</xsd:attribute>
|
||||
<!-- TODO: add support for header sub-elements -->
|
||||
</xsd:complexType>
|
||||
</xsd:element>
|
||||
|
||||
|
||||
@@ -17,13 +17,13 @@
|
||||
<queue/>
|
||||
</channel>
|
||||
|
||||
<publisher id="fixedDelayPublisher" fixed-delay="1234" payload="'fixedDelayTest'" channel="fixedDelayChannel" auto-startup="false"/>
|
||||
<scheduled-producer id="fixedDelayProducer" fixed-delay="1234" payload-expression="'fixedDelayTest'" channel="fixedDelayChannel" auto-startup="false"/>
|
||||
|
||||
<publisher id="fixedRatePublisher" fixed-rate="5678" payload="'fixedRateTest'" channel="fixedRateChannel" auto-startup="false"/>
|
||||
<scheduled-producer id="fixedRateProducer" fixed-rate="5678" payload-expression="'fixedRateTest'" channel="fixedRateChannel" auto-startup="false"/>
|
||||
|
||||
<publisher id="cronPublisher" cron="7 6 5 4 3 ?" payload="'cronTest'" channel="cronChannel" auto-startup="false"/>
|
||||
<scheduled-producer id="cronProducer" cron="7 6 5 4 3 ?" payload-expression="'cronTest'" channel="cronChannel" auto-startup="false"/>
|
||||
|
||||
<publisher id="triggerRefPublisher" trigger="customTrigger" payload="'triggerRefTest'" channel="triggerRefChannel"/>
|
||||
<scheduled-producer id="triggerRefProducer" trigger="customTrigger" payload-expression="'triggerRefTest'" channel="triggerRefChannel"/>
|
||||
|
||||
<beans:bean id="customTrigger" class="org.springframework.scheduling.support.PeriodicTrigger">
|
||||
<beans:constructor-arg value="9999"/>
|
||||
@@ -27,7 +27,7 @@ import org.springframework.beans.DirectFieldAccessor;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.context.ApplicationContext;
|
||||
import org.springframework.expression.Expression;
|
||||
import org.springframework.integration.endpoint.TriggeredMessagePublisher;
|
||||
import org.springframework.integration.endpoint.ScheduledMessageProducer;
|
||||
import org.springframework.scheduling.Trigger;
|
||||
import org.springframework.scheduling.support.CronTrigger;
|
||||
import org.springframework.scheduling.support.PeriodicTrigger;
|
||||
@@ -40,7 +40,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
*/
|
||||
@ContextConfiguration
|
||||
@RunWith(SpringJUnit4ClassRunner.class)
|
||||
public class PublisherParserTests {
|
||||
public class ScheduledProducerParserTests {
|
||||
|
||||
@Autowired
|
||||
private ApplicationContext context;
|
||||
@@ -48,61 +48,61 @@ public class PublisherParserTests {
|
||||
|
||||
@Test
|
||||
public void fixedDelay() {
|
||||
TriggeredMessagePublisher publisher = context.getBean("fixedDelayPublisher", TriggeredMessagePublisher.class);
|
||||
assertFalse(publisher.isAutoStartup());
|
||||
DirectFieldAccessor publisherAccessor = new DirectFieldAccessor(publisher);
|
||||
Trigger trigger = (Trigger) publisherAccessor.getPropertyValue("trigger");
|
||||
ScheduledMessageProducer producer = context.getBean("fixedDelayProducer", ScheduledMessageProducer.class);
|
||||
assertFalse(producer.isAutoStartup());
|
||||
DirectFieldAccessor producerAccessor = new DirectFieldAccessor(producer);
|
||||
Trigger trigger = (Trigger) producerAccessor.getPropertyValue("trigger");
|
||||
assertEquals(PeriodicTrigger.class, trigger.getClass());
|
||||
DirectFieldAccessor triggerAccessor = new DirectFieldAccessor(trigger);
|
||||
assertEquals(1234L, triggerAccessor.getPropertyValue("period"));
|
||||
assertEquals(Boolean.FALSE, triggerAccessor.getPropertyValue("fixedRate"));
|
||||
assertEquals(context.getBean("fixedDelayChannel"), publisherAccessor.getPropertyValue("outputChannel"));
|
||||
assertEquals(context.getBean("fixedDelayChannel"), producerAccessor.getPropertyValue("outputChannel"));
|
||||
Expression payloadExpression = (Expression) new DirectFieldAccessor(
|
||||
publisherAccessor.getPropertyValue("task")).getPropertyValue("payloadExpression");
|
||||
producerAccessor.getPropertyValue("task")).getPropertyValue("payloadExpression");
|
||||
assertEquals("'fixedDelayTest'", payloadExpression.getExpressionString());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void fixedRate() {
|
||||
TriggeredMessagePublisher publisher = context.getBean("fixedRatePublisher", TriggeredMessagePublisher.class);
|
||||
assertFalse(publisher.isAutoStartup());
|
||||
DirectFieldAccessor publisherAccessor = new DirectFieldAccessor(publisher);
|
||||
Trigger trigger = (Trigger) publisherAccessor.getPropertyValue("trigger");
|
||||
ScheduledMessageProducer producer = context.getBean("fixedRateProducer", ScheduledMessageProducer.class);
|
||||
assertFalse(producer.isAutoStartup());
|
||||
DirectFieldAccessor producerAccessor = new DirectFieldAccessor(producer);
|
||||
Trigger trigger = (Trigger) producerAccessor.getPropertyValue("trigger");
|
||||
assertEquals(PeriodicTrigger.class, trigger.getClass());
|
||||
DirectFieldAccessor triggerAccessor = new DirectFieldAccessor(trigger);
|
||||
assertEquals(5678L, triggerAccessor.getPropertyValue("period"));
|
||||
assertEquals(Boolean.TRUE, triggerAccessor.getPropertyValue("fixedRate"));
|
||||
assertEquals(context.getBean("fixedRateChannel"), publisherAccessor.getPropertyValue("outputChannel"));
|
||||
assertEquals(context.getBean("fixedRateChannel"), producerAccessor.getPropertyValue("outputChannel"));
|
||||
Expression payloadExpression = (Expression) new DirectFieldAccessor(
|
||||
publisherAccessor.getPropertyValue("task")).getPropertyValue("payloadExpression");
|
||||
producerAccessor.getPropertyValue("task")).getPropertyValue("payloadExpression");
|
||||
assertEquals("'fixedRateTest'", payloadExpression.getExpressionString());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void cron() {
|
||||
TriggeredMessagePublisher publisher = context.getBean("cronPublisher", TriggeredMessagePublisher.class);
|
||||
assertFalse(publisher.isAutoStartup());
|
||||
DirectFieldAccessor publisherAccessor = new DirectFieldAccessor(publisher);
|
||||
Trigger trigger = (Trigger) publisherAccessor.getPropertyValue("trigger");
|
||||
ScheduledMessageProducer producer = context.getBean("cronProducer", ScheduledMessageProducer.class);
|
||||
assertFalse(producer.isAutoStartup());
|
||||
DirectFieldAccessor producerAccessor = new DirectFieldAccessor(producer);
|
||||
Trigger trigger = (Trigger) producerAccessor.getPropertyValue("trigger");
|
||||
assertEquals(CronTrigger.class, trigger.getClass());
|
||||
assertEquals("7 6 5 4 3 ?", new DirectFieldAccessor(new DirectFieldAccessor(
|
||||
trigger).getPropertyValue("sequenceGenerator")).getPropertyValue("expression"));
|
||||
assertEquals(context.getBean("cronChannel"), publisherAccessor.getPropertyValue("outputChannel"));
|
||||
assertEquals(context.getBean("cronChannel"), producerAccessor.getPropertyValue("outputChannel"));
|
||||
Expression payloadExpression = (Expression) new DirectFieldAccessor(
|
||||
publisherAccessor.getPropertyValue("task")).getPropertyValue("payloadExpression");
|
||||
producerAccessor.getPropertyValue("task")).getPropertyValue("payloadExpression");
|
||||
assertEquals("'cronTest'", payloadExpression.getExpressionString());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void triggerRef() {
|
||||
TriggeredMessagePublisher publisher = context.getBean("triggerRefPublisher", TriggeredMessagePublisher.class);
|
||||
assertTrue(publisher.isAutoStartup());
|
||||
DirectFieldAccessor publisherAccessor = new DirectFieldAccessor(publisher);
|
||||
Trigger trigger = (Trigger) publisherAccessor.getPropertyValue("trigger");
|
||||
ScheduledMessageProducer producer = context.getBean("triggerRefProducer", ScheduledMessageProducer.class);
|
||||
assertTrue(producer.isAutoStartup());
|
||||
DirectFieldAccessor producerAccessor = new DirectFieldAccessor(producer);
|
||||
Trigger trigger = (Trigger) producerAccessor.getPropertyValue("trigger");
|
||||
assertEquals(context.getBean("customTrigger"), trigger);
|
||||
assertEquals(context.getBean("triggerRefChannel"), publisherAccessor.getPropertyValue("outputChannel"));
|
||||
assertEquals(context.getBean("triggerRefChannel"), producerAccessor.getPropertyValue("outputChannel"));
|
||||
Expression payloadExpression = (Expression) new DirectFieldAccessor(
|
||||
publisherAccessor.getPropertyValue("task")).getPropertyValue("payloadExpression");
|
||||
producerAccessor.getPropertyValue("task")).getPropertyValue("payloadExpression");
|
||||
assertEquals("'triggerRefTest'", payloadExpression.getExpressionString());
|
||||
}
|
||||
|
||||
@@ -36,7 +36,7 @@ import org.springframework.scheduling.support.PeriodicTrigger;
|
||||
* @author Mark Fisher
|
||||
* @since 2.0
|
||||
*/
|
||||
public class TriggeredMessagePublisherTests {
|
||||
public class ScheduledMessageProducerTests {
|
||||
|
||||
private static final AtomicInteger counter = new AtomicInteger();
|
||||
|
||||
@@ -45,17 +45,17 @@ public class TriggeredMessagePublisherTests {
|
||||
public void test() throws Exception {
|
||||
QueueChannel channel = new QueueChannel();
|
||||
Trigger trigger = new PeriodicTrigger(100);
|
||||
String payloadExpression = "'test-' + T(org.springframework.integration.endpoint.TriggeredMessagePublisherTests).next()";
|
||||
String payloadExpression = "'test-' + T(org.springframework.integration.endpoint.ScheduledMessageProducerTests).next()";
|
||||
ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
|
||||
scheduler.afterPropertiesSet();
|
||||
Map<String, String> headerExpressions = new HashMap<String, String>();
|
||||
headerExpressions.put("foo", "'x'");
|
||||
headerExpressions.put("bar", "7 * 6");
|
||||
TriggeredMessagePublisher publisher = new TriggeredMessagePublisher(trigger, payloadExpression);
|
||||
publisher.setHeaderExpressions(headerExpressions);
|
||||
publisher.setTaskScheduler(scheduler);
|
||||
publisher.setOutputChannel(channel);
|
||||
publisher.start();
|
||||
ScheduledMessageProducer producer = new ScheduledMessageProducer(trigger, payloadExpression);
|
||||
producer.setHeaderExpressions(headerExpressions);
|
||||
producer.setTaskScheduler(scheduler);
|
||||
producer.setOutputChannel(channel);
|
||||
producer.start();
|
||||
List<Message<?>> messages = new ArrayList<Message<?>>();
|
||||
for (int i = 0; i < 3; i++) {
|
||||
messages.add(channel.receive(1000));
|
||||
Reference in New Issue
Block a user