diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AnnotationConfigUtils.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AnnotationConfigUtils.java index 5cda107d26..49af70609c 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AnnotationConfigUtils.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AnnotationConfigUtils.java @@ -79,7 +79,7 @@ public abstract class AnnotationConfigUtils { public static Trigger parseTriggerFromPollerAnnotation(Poller pollerAnnotation) { IntervalTrigger trigger = new IntervalTrigger( pollerAnnotation.interval(), pollerAnnotation.timeUnit()); - trigger.setInitialDelay(pollerAnnotation.initialDelay(), pollerAnnotation.timeUnit()); + trigger.setInitialDelay(pollerAnnotation.initialDelay()); trigger.setFixedRate(pollerAnnotation.fixedRate()); return trigger; } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AbstractConsumerEndpointParser.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AbstractConsumerEndpointParser.java index 74f0976893..62fd65a7b9 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AbstractConsumerEndpointParser.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AbstractConsumerEndpointParser.java @@ -104,7 +104,7 @@ public abstract class AbstractConsumerEndpointParser extends AbstractSingleBeanD builder.addPropertyValue("inputChannelName", inputChannelName); Element pollerElement = DomUtils.getChildElementByTagName(element, POLLER_ELEMENT); if (pollerElement != null) { - IntegrationNamespaceUtils.configureTrigger(pollerElement, builder); + IntegrationNamespaceUtils.configureTrigger(pollerElement, builder, parserContext); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, pollerElement, "max-messages-per-poll"); Element txElement = DomUtils.getChildElementByTagName(pollerElement, "transactional"); if (txElement != null) { diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AbstractOutboundChannelAdapterParser.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AbstractOutboundChannelAdapterParser.java index 8d82207fd9..20b614e190 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AbstractOutboundChannelAdapterParser.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AbstractOutboundChannelAdapterParser.java @@ -40,7 +40,7 @@ public abstract class AbstractOutboundChannelAdapterParser extends AbstractChann builder.addConstructorArgReference(this.parseAndRegisterConsumer(element, parserContext)); if (pollerElement != null) { Assert.hasText(channelName, "outbound channel adapter with a 'poller' requires a 'channel' to poll"); - IntegrationNamespaceUtils.configureTrigger(pollerElement, builder); + IntegrationNamespaceUtils.configureTrigger(pollerElement, builder, parserContext); Element txElement = DomUtils.getChildElementByTagName(pollerElement, "transactional"); if (txElement != null) { IntegrationNamespaceUtils.configureTransactionAttributes(txElement, builder); diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AbstractPollingInboundChannelAdapterParser.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AbstractPollingInboundChannelAdapterParser.java index e8e2c85870..de8d7fe9a9 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AbstractPollingInboundChannelAdapterParser.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AbstractPollingInboundChannelAdapterParser.java @@ -42,7 +42,7 @@ public abstract class AbstractPollingInboundChannelAdapterParser extends Abstrac adapterBuilder.addPropertyReference("source", source); adapterBuilder.addPropertyReference("outputChannel", channelName); if (pollerElement != null) { - IntegrationNamespaceUtils.configureTrigger(pollerElement, adapterBuilder); + IntegrationNamespaceUtils.configureTrigger(pollerElement, adapterBuilder, parserContext); IntegrationNamespaceUtils.setValueIfAttributeDefined(adapterBuilder, pollerElement, "max-messages-per-poll"); Element txElement = DomUtils.getChildElementByTagName(pollerElement, "transactional"); if (txElement != null) { diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceUtils.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceUtils.java index 61f66d1172..7fdb4107ae 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceUtils.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceUtils.java @@ -23,12 +23,12 @@ import org.w3c.dom.Element; import org.springframework.beans.factory.config.BeanDefinitionHolder; import org.springframework.beans.factory.parsing.BeanComponentDefinition; import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.support.BeanDefinitionReaderUtils; import org.springframework.beans.factory.xml.BeanDefinitionParserDelegate; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.core.Conventions; import org.springframework.integration.scheduling.CronTrigger; import org.springframework.integration.scheduling.IntervalTrigger; -import org.springframework.integration.scheduling.Trigger; import org.springframework.transaction.support.DefaultTransactionDefinition; import org.springframework.util.Assert; import org.springframework.util.StringUtils; @@ -142,38 +142,39 @@ public abstract class IntegrationNamespaceUtils { * @param pollerElement the "poller" element to parse * @param targetBuilder the builder that expects the "trigger" property */ - public static void configureTrigger(Element pollerElement, BeanDefinitionBuilder targetBuilder) { - Trigger trigger = null; + public static void configureTrigger(Element pollerElement, BeanDefinitionBuilder targetBuilder, ParserContext parserContext) { + String triggerBeanName = null; Element intervalElement = DomUtils.getChildElementByTagName(pollerElement, "interval-trigger"); if (intervalElement != null) { - trigger = createIntervalTrigger(intervalElement); + triggerBeanName = parseIntervalTrigger(intervalElement, parserContext); } else { Element cronElement = DomUtils.getChildElementByTagName(pollerElement, "cron-trigger"); Assert.notNull(cronElement, "A element must include either an or child element."); - trigger = createCronTrigger(cronElement); + triggerBeanName = parseCronTrigger(cronElement, parserContext); } - targetBuilder.addPropertyValue("trigger", trigger); + targetBuilder.addPropertyReference("trigger", triggerBeanName); } - private static Trigger createIntervalTrigger(Element element) { + private static String parseIntervalTrigger(Element element, ParserContext parserContext) { String interval = element.getAttribute("interval"); Assert.hasText(interval, "the 'interval' attribute is required for an "); TimeUnit timeUnit = TimeUnit.valueOf(element.getAttribute("time-unit")); - IntervalTrigger trigger = new IntervalTrigger(Long.valueOf(interval), timeUnit); - String initialDelay = element.getAttribute("initial-delay"); - if (StringUtils.hasText(initialDelay)) { - trigger.setInitialDelay(Long.valueOf(initialDelay), timeUnit); - } - trigger.setFixedRate("true".equals(element.getAttribute("fixed-rate").toLowerCase())); - return trigger; + BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(IntervalTrigger.class); + builder.addConstructorArgValue(interval); + builder.addConstructorArgValue(timeUnit); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "initial-delay"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "fixed-rate"); + return BeanDefinitionReaderUtils.registerWithGeneratedName(builder.getBeanDefinition(), parserContext.getRegistry()); } - private static Trigger createCronTrigger(Element element) { + private static String parseCronTrigger(Element element, ParserContext parserContext) { String cronExpression = element.getAttribute("expression"); Assert.hasText(cronExpression, "the 'expression' attribute is required for a "); - return new CronTrigger(cronExpression); + BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(CronTrigger.class); + builder.addConstructorArgValue(cronExpression); + return BeanDefinitionReaderUtils.registerWithGeneratedName(builder.getBeanDefinition(), parserContext.getRegistry()); } /** diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/scheduling/IntervalTrigger.java b/org.springframework.integration/src/main/java/org/springframework/integration/scheduling/IntervalTrigger.java index 4c6c8dede1..98b6160a68 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/scheduling/IntervalTrigger.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/scheduling/IntervalTrigger.java @@ -35,6 +35,8 @@ public class IntervalTrigger implements Trigger { private final long interval; + private final TimeUnit timeUnit; + private volatile long initialDelay = 0; private volatile boolean fixedRate = false; @@ -44,15 +46,16 @@ public class IntervalTrigger implements Trigger { * Create a trigger with the given interval in milliseconds. */ public IntervalTrigger(long interval) { - Assert.isTrue(interval >= 0, "interval must not be negative"); - this.interval = interval; + this(interval, null); } /** * Create a trigger with the given interval and time unit. */ - public IntervalTrigger(long interval, TimeUnit unit) { - this(unit.toMillis(interval)); + public IntervalTrigger(long interval, TimeUnit timeUnit) { + Assert.isTrue(interval >= 0, "interval must not be negative"); + this.timeUnit = (timeUnit != null) ? timeUnit : TimeUnit.MILLISECONDS; + this.interval = this.timeUnit.toMillis(interval); } @@ -60,14 +63,7 @@ public class IntervalTrigger implements Trigger { * Specify the delay for the initial execution. */ public void setInitialDelay(long initialDelay) { - this.initialDelay = initialDelay; - } - - /** - * Specify the delay for the initial execution using the given time unit. - */ - public void setInitialDelay(long initialDelay, TimeUnit unit) { - this.initialDelay = unit.toMillis(initialDelay); + this.initialDelay = this.timeUnit.toMillis(initialDelay); } /**