From 193f3d2439a8d530a01751e3595629fb7033c281 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Fri, 11 Nov 2011 12:14:16 -0500 Subject: [PATCH 1/4] INT-2235 added support for resource-inbound-channel-adapter removed pre-fetching logic changed ResourcePatternResolvingMessageSource to return multiple Resources made filter optional --- .../xml/IntegrationNamespaceHandler.java | 1 + .../ResourceInboundChannelAdapterParser.java | 42 ++++++ .../endpoint/SourcePollingChannelAdapter.java | 5 - .../resource/ResourceMessageSource.java | 104 ++++++++++++++ .../AcceptOnceUntilPurgedElementFilter.java | 81 +++++++++++ .../integration/util/ElementFilter.java | 28 ++++ .../config/xml/spring-integration-2.1.xsd | 72 ++++++++++ .../ResourcePatternResolver-config-custom.xml | 20 +++ .../ResourcePatternResolver-config-fail.xml | 17 +++ .../ResourcePatternResolver-config-usage.xml | 18 +++ ...ResourcePatternResolver-config-usagerf.xml | 21 +++ .../ResourcePatternResolver-config.xml | 17 +++ .../ResourcePatternResolverParserTests.java | 132 ++++++++++++++++++ 13 files changed, 553 insertions(+), 5 deletions(-) create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResourceInboundChannelAdapterParser.java create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/resource/ResourceMessageSource.java create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/util/ElementFilter.java create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-custom.xml create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-fail.xml create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usage.xml create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usagerf.xml create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config.xml create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceHandler.java index 25377b2b8e..ce084278c6 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceHandler.java @@ -53,6 +53,7 @@ public class IntegrationNamespaceHandler extends AbstractIntegrationNamespaceHan registerBeanDefinitionParser("claim-check-in", new ClaimCheckInParser()); registerBeanDefinitionParser("claim-check-out", new ClaimCheckOutParser()); registerBeanDefinitionParser("inbound-channel-adapter", new MethodInvokingInboundChannelAdapterParser()); + registerBeanDefinitionParser("resource-inbound-channel-adapter", new ResourceInboundChannelAdapterParser()); registerBeanDefinitionParser("outbound-channel-adapter", new MethodInvokingOutboundChannelAdapterParser()); registerBeanDefinitionParser("logging-channel-adapter", new LoggingChannelAdapterParser()); registerBeanDefinitionParser("gateway", new GatewayParser()); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResourceInboundChannelAdapterParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResourceInboundChannelAdapterParser.java new file mode 100644 index 0000000000..ad02416e98 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResourceInboundChannelAdapterParser.java @@ -0,0 +1,42 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.integration.config.xml; + +import org.springframework.beans.BeanMetadataElement; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.xml.ParserContext; +import org.springframework.integration.resource.ResourceMessageSource; +import org.w3c.dom.Element; + +/** + * Parser for 'resource-inbound-channel-adapter' + * + * @author Oleg Zhurakousky + * @since 2.1 + */ +public class ResourceInboundChannelAdapterParser extends AbstractPollingInboundChannelAdapterParser { + + + @Override + protected BeanMetadataElement parseSource(Element element, ParserContext parserContext) { + BeanDefinitionBuilder sourceBuilder = BeanDefinitionBuilder.genericBeanDefinition(ResourceMessageSource.class); + IntegrationNamespaceUtils.setValueIfAttributeDefined(sourceBuilder, element, "pattern"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(sourceBuilder, element, "pattern-resolver"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(sourceBuilder, element, "filter"); + return sourceBuilder.getBeanDefinition(); + } + +} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java index ef72728e6f..31f0424f3d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java @@ -16,9 +16,6 @@ package org.springframework.integration.endpoint; -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; - import org.springframework.integration.Message; import org.springframework.integration.MessageChannel; import org.springframework.integration.context.NamedComponent; @@ -36,8 +33,6 @@ import org.springframework.util.Assert; * @author Oleg Zhurakousky */ public class SourcePollingChannelAdapter extends AbstractPollingEndpoint implements TrackableComponent { - - private final Log logger = LogFactory.getLog(this.getClass()); private volatile MessageSource source; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/resource/ResourceMessageSource.java b/spring-integration-core/src/main/java/org/springframework/integration/resource/ResourceMessageSource.java new file mode 100644 index 0000000000..6a9714a187 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/resource/ResourceMessageSource.java @@ -0,0 +1,104 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.resource; + +import java.util.ArrayList; +import java.util.List; + +import org.springframework.beans.factory.InitializingBean; +import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; +import org.springframework.context.MessageSource; +import org.springframework.core.io.Resource; +import org.springframework.core.io.support.ResourcePatternResolver; +import org.springframework.integration.MessagingException; +import org.springframework.integration.endpoint.AbstractMessageSource; +import org.springframework.integration.util.ElementFilter; +import org.springframework.util.Assert; +import org.springframework.util.ObjectUtils; + +/** + * Implementation of {@link MessageSource} based on {@link ResourcePatternResolver} which will + * attempt to resolve {@link Resource}s based on the pattern specified. + * + * @author Oleg Zhurakousky + * @since 2.1 + */ +public class ResourceMessageSource extends AbstractMessageSource implements ApplicationContextAware, InitializingBean { + + private volatile String pattern; + + private volatile ApplicationContext applicationContext; + + private volatile ResourcePatternResolver patternResolver; + + private volatile ElementFilter filter; + + + public void setPatternResolver(ResourcePatternResolver patternResolver) { + this.patternResolver = patternResolver; + } + + public void setPattern(String pattern) { + this.pattern = pattern; + } + + public void setFilter(ElementFilter filter) { + this.filter = filter; + } + + public void setApplicationContext(ApplicationContext applicationContext) { + this.applicationContext = applicationContext; + } + + public void afterPropertiesSet() { + if (this.patternResolver == null) { + if (this.applicationContext instanceof ResourcePatternResolver) { + this.patternResolver = this.applicationContext; + } + } + Assert.notNull(this.patternResolver, "no 'patternResolver' is specified"); + Assert.hasText(this.pattern, "'pattern' must be specified"); + } + + @Override + protected Resource[] doReceive() { + try { + Resource[] resources = this.patternResolver.getResources(this.pattern); + if (this.filter != null && !ObjectUtils.isEmpty(resources)) { + List filteredResources = new ArrayList(); + for (Resource resource : resources) { + Resource filteredResource = this.filter.filter(resource); + if (filteredResource != null) { + filteredResources.add(filteredResource); + } + } + if (filteredResources.size() == 0) { + resources = null; + } + else { + resources = filteredResources.toArray(new Resource[0]); + } + } + return resources; + } + catch (Exception e) { + throw new MessagingException("Attempt to retrieve Resources failed", e); + } + } + +} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java b/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java new file mode 100644 index 0000000000..e34ceeee97 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java @@ -0,0 +1,81 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.integration.util; + +import java.util.Queue; +import java.util.concurrent.LinkedBlockingQueue; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; + +import org.springframework.util.Assert; + +/** + * An implementation of {@link ElementFilter} which will queue all items that's been seen until + * the queue reaches its capacity after which one item from the queue will be purged to make room for a + * new item to be added. Note that however unlikely the removed item will now appear as unprocessed + * so it is highly recommended to move/delete resources which corresponds to the underlying items once processing + * is done to eliminate duplicate processing. + * + * @author Oleg Zhurakousky + * @since 2.1 + */ +public class AcceptOnceUntilPurgedElementFilter implements ElementFilter { + + private final Log logger = LogFactory.getLog(this.getClass()); + + private final Queue seenItems; + + private final Object seenQueueMonitor = new Object(); + + public AcceptOnceUntilPurgedElementFilter(){ + this(Integer.MAX_VALUE); + } + + public AcceptOnceUntilPurgedElementFilter(int maxCapacity){ + seenItems = new LinkedBlockingQueue(maxCapacity); + } + + private boolean accept(T item) { + synchronized (this.seenQueueMonitor) { + boolean accepted = false; + + if (!this.seenItems.contains(item)) { + accepted = this.seenItems.offer(item); + if (!accepted){ + logger.warn("'seenQueueMonitor' queue of AcceptOnceUntilPurgedElementFilter is at the capacity, " + + "evicting one item to make room for another"); + this.seenItems.poll(); + accepted = this.seenItems.offer(item); + } + } + + return accepted; + } + } + + public T filter(T unfilteredElement) { + Assert.notNull(unfilteredElement, "'unfilteredElement' must not be null"); + + if (this.accept(unfilteredElement)){ + return unfilteredElement; + } + else { + return null; + } + } + +} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/ElementFilter.java b/spring-integration-core/src/main/java/org/springframework/integration/util/ElementFilter.java new file mode 100644 index 0000000000..d0aaac9277 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/ElementFilter.java @@ -0,0 +1,28 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.integration.util; + + +/** + * Base strategy for filtering out an element + * + * @author Oleg Zhurakousky + * @since 2.1 + */ +public interface ElementFilter { + + T filter(T unfilteredElement); +} diff --git a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.1.xsd b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.1.xsd index 9d1b2f987f..5ebc415514 100644 --- a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.1.xsd +++ b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.1.xsd @@ -716,6 +716,78 @@ + + + + + Defines a Channel Adapter that receives Resource(s) and sends them to a + MessageChannel identified via 'channel' attribute. + + + + + + + + + + + + Component identifier + + + + + + + Channel where Message will be sent to + + + + + + + Reference to the implementation of org.springframework.integration.util.ElementFilter. + + + + + + + + + + + + Lifecycle attribute signaling if this component should be started during Application Context startup. + + + + + + + Location pattern expression (e.g., "/**/*.txt") + + + + + + + + Reference to a org.springframework.core.io.support.ResourcePatternResolver. + + + + + + + Maximum amount of time in milliseconds to wait when sending a message to the channel if such channel may block. + For example, a Queue Channel can block until space is available if its maximum capacity has been reached. + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-custom.xml b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-custom.xml new file mode 100644 index 0000000000..7c6c3aea2c --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-custom.xml @@ -0,0 +1,20 @@ + + + + + + + + + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-fail.xml b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-fail.xml new file mode 100644 index 0000000000..9671745127 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-fail.xml @@ -0,0 +1,17 @@ + + + + + + + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usage.xml b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usage.xml new file mode 100644 index 0000000000..d24ffcfc45 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usage.xml @@ -0,0 +1,18 @@ + + + + + + + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usagerf.xml b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usagerf.xml new file mode 100644 index 0000000000..5ff664ebf9 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usagerf.xml @@ -0,0 +1,21 @@ + + + + + + + + + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config.xml b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config.xml new file mode 100644 index 0000000000..f5eff57ce0 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config.xml @@ -0,0 +1,17 @@ + + + + + + + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java new file mode 100644 index 0000000000..68da765a4d --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java @@ -0,0 +1,132 @@ +/* + * Copyright 2002-2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.integration.resource; + +import java.io.File; + +import org.junit.Test; + +import org.springframework.beans.factory.BeanCreationException; +import org.springframework.context.ApplicationContext; +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.core.io.Resource; +import org.springframework.integration.Message; +import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.endpoint.SourcePollingChannelAdapter; +import org.springframework.integration.test.util.TestUtils; +import org.springframework.integration.util.ElementFilter; + +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.assertTrue; + +/** + * @author Oleg Zhurakousky + * + */ +public class ResourcePatternResolverParserTests { + + @Test + public void testDefaultConfig(){ + ApplicationContext context = new ClassPathXmlApplicationContext("ResourcePatternResolver-config.xml", this.getClass()); + SourcePollingChannelAdapter resourceAdapter = context.getBean("resourceAdapterDefault", SourcePollingChannelAdapter.class); + ResourceMessageSource source = TestUtils.getPropertyValue(resourceAdapter, "source", ResourceMessageSource.class); + assertNotNull(source); + boolean autoStartup = TestUtils.getPropertyValue(resourceAdapter, "autoStartup", Boolean.class); + assertFalse(autoStartup); + + assertEquals("/**/*", TestUtils.getPropertyValue(source, "pattern")); + assertEquals(context, TestUtils.getPropertyValue(source, "patternResolver")); + } + + @Test(expected=BeanCreationException.class) + public void testDefaultConfigNoLocationPattern(){ + new ClassPathXmlApplicationContext("ResourcePatternResolver-config-fail.xml", this.getClass()); + } + + @Test + public void testCustomPatternResolver(){ + ApplicationContext context = new ClassPathXmlApplicationContext("ResourcePatternResolver-config-custom.xml", this.getClass()); + SourcePollingChannelAdapter resourceAdapter = context.getBean("resourceAdapterDefault", SourcePollingChannelAdapter.class); + ResourceMessageSource source = TestUtils.getPropertyValue(resourceAdapter, "source", ResourceMessageSource.class); + assertNotNull(source); + assertEquals(context.getBean("customResolver"), TestUtils.getPropertyValue(source, "patternResolver")); + } + + @SuppressWarnings("unchecked") + @Test + public void testUsage() throws Exception{ + + File baseDir = new File(System.getProperty("java.io.tmpdir")); + + for (int i = 0; i < 10; i++) { + File f = new File(baseDir, "testUsage"+i); + f.createNewFile(); + } + + ApplicationContext context = new ClassPathXmlApplicationContext("ResourcePatternResolver-config-usage.xml", this.getClass()); + QueueChannel resultChannel = context.getBean("resultChannel", QueueChannel.class); + Message message = (Message) resultChannel.receive(3000); + assertNotNull(message); + Resource[] resources = message.getPayload(); + for (Resource resource : resources) { + assertTrue(resource.getURI().toString().contains("testUsage")); + } + } + + @SuppressWarnings("unchecked") + @Test + public void testUsageWithCustomResourceFilter() throws Exception{ + + File baseDir = new File(System.getProperty("java.io.tmpdir")); + for (int i = 0; i < 10; i++) { + File f = new File(baseDir, "testUsageWithRf"+i); + f.createNewFile(); + } + + ApplicationContext context = new ClassPathXmlApplicationContext("ResourcePatternResolver-config-usagerf.xml", this.getClass()); + SourcePollingChannelAdapter resourceAdapter = context.getBean("resourceAdapterDefault", SourcePollingChannelAdapter.class); + ResourceMessageSource source = TestUtils.getPropertyValue(resourceAdapter, "source", ResourceMessageSource.class); + assertNotNull(source); + assertEquals(context.getBean("rlFilter"), TestUtils.getPropertyValue(source, "filter")); + + QueueChannel resultChannel = context.getBean("resultChannel", QueueChannel.class); + + Message message = (Message) resultChannel.receive(1000); + assertNotNull(message); + + message = (Message) resultChannel.receive(1000); + assertNull(message); + + } + + public static class OneItemAndNeverAgainResourceListFilter implements ElementFilter { + + private volatile boolean once = false; + + public Resource filter(Resource unfilteredElement) { + + if (!once){ + once = true; + return unfilteredElement; + } + return null; + } + + } +} From 32d64b19f42d126e86410a05aa74f67988d514b8 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Tue, 22 Nov 2011 12:57:48 -0500 Subject: [PATCH 2/4] ElementFilter now CollectionFilter Refactored filter to handle a Collection of Resources Made 'pattern' a constructor-arg since it's mandatory --- .../ResourceInboundChannelAdapterParser.java | 3 +- .../resource/ResourceMessageSource.java | 37 +++++++------- .../AcceptOnceUntilPurgedElementFilter.java | 48 +++++++++++-------- ...ementFilter.java => CollectionFilter.java} | 10 ++-- .../ResourcePatternResolverParserTests.java | 32 +++++++------ 5 files changed, 72 insertions(+), 58 deletions(-) rename spring-integration-core/src/main/java/org/springframework/integration/util/{ElementFilter.java => CollectionFilter.java} (76%) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResourceInboundChannelAdapterParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResourceInboundChannelAdapterParser.java index ad02416e98..9284805970 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResourceInboundChannelAdapterParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResourceInboundChannelAdapterParser.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.integration.config.xml; import org.springframework.beans.BeanMetadataElement; @@ -33,7 +34,7 @@ public class ResourceInboundChannelAdapterParser extends AbstractPollingInboundC @Override protected BeanMetadataElement parseSource(Element element, ParserContext parserContext) { BeanDefinitionBuilder sourceBuilder = BeanDefinitionBuilder.genericBeanDefinition(ResourceMessageSource.class); - IntegrationNamespaceUtils.setValueIfAttributeDefined(sourceBuilder, element, "pattern"); + sourceBuilder.addConstructorArgValue(element.getAttribute("pattern")); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(sourceBuilder, element, "pattern-resolver"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(sourceBuilder, element, "filter"); return sourceBuilder.getBeanDefinition(); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/resource/ResourceMessageSource.java b/spring-integration-core/src/main/java/org/springframework/integration/resource/ResourceMessageSource.java index 6a9714a187..7b104b7ca5 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/resource/ResourceMessageSource.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/resource/ResourceMessageSource.java @@ -16,8 +16,8 @@ package org.springframework.integration.resource; -import java.util.ArrayList; -import java.util.List; +import java.util.Arrays; +import java.util.Collection; import org.springframework.beans.factory.InitializingBean; import org.springframework.context.ApplicationContext; @@ -27,8 +27,9 @@ import org.springframework.core.io.Resource; import org.springframework.core.io.support.ResourcePatternResolver; import org.springframework.integration.MessagingException; import org.springframework.integration.endpoint.AbstractMessageSource; -import org.springframework.integration.util.ElementFilter; +import org.springframework.integration.util.CollectionFilter; import org.springframework.util.Assert; +import org.springframework.util.CollectionUtils; import org.springframework.util.ObjectUtils; /** @@ -36,28 +37,31 @@ import org.springframework.util.ObjectUtils; * attempt to resolve {@link Resource}s based on the pattern specified. * * @author Oleg Zhurakousky + * @author Mark Fisher * @since 2.1 */ public class ResourceMessageSource extends AbstractMessageSource implements ApplicationContextAware, InitializingBean { - private volatile String pattern; + private final String pattern; private volatile ApplicationContext applicationContext; private volatile ResourcePatternResolver patternResolver; - private volatile ElementFilter filter; + private volatile CollectionFilter filter; + + + public ResourceMessageSource(String pattern) { + Assert.hasText(pattern, "pattern must not be empty"); + this.pattern = pattern; + } public void setPatternResolver(ResourcePatternResolver patternResolver) { this.patternResolver = patternResolver; } - public void setPattern(String pattern) { - this.pattern = pattern; - } - - public void setFilter(ElementFilter filter) { + public void setFilter(CollectionFilter filter) { this.filter = filter; } @@ -71,8 +75,7 @@ public class ResourceMessageSource extends AbstractMessageSource imp this.patternResolver = this.applicationContext; } } - Assert.notNull(this.patternResolver, "no 'patternResolver' is specified"); - Assert.hasText(this.pattern, "'pattern' must be specified"); + Assert.notNull(this.patternResolver, "no 'patternResolver' available"); } @Override @@ -80,14 +83,8 @@ public class ResourceMessageSource extends AbstractMessageSource imp try { Resource[] resources = this.patternResolver.getResources(this.pattern); if (this.filter != null && !ObjectUtils.isEmpty(resources)) { - List filteredResources = new ArrayList(); - for (Resource resource : resources) { - Resource filteredResource = this.filter.filter(resource); - if (filteredResource != null) { - filteredResources.add(filteredResource); - } - } - if (filteredResources.size() == 0) { + Collection filteredResources = this.filter.filter(Arrays.asList(resources)); + if (CollectionUtils.isEmpty(filteredResources)) { resources = null; } else { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java b/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java index e34ceeee97..2369f23c9d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java @@ -13,40 +13,59 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.integration.util; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; import java.util.Queue; import java.util.concurrent.LinkedBlockingQueue; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; - import org.springframework.util.Assert; /** - * An implementation of {@link ElementFilter} which will queue all items that's been seen until + * An implementation of {@link CollectionFilter} which will queue all items that have been seen until * the queue reaches its capacity after which one item from the queue will be purged to make room for a * new item to be added. Note that however unlikely the removed item will now appear as unprocessed * so it is highly recommended to move/delete resources which corresponds to the underlying items once processing * is done to eliminate duplicate processing. * * @author Oleg Zhurakousky + * @author Mark Fisher * @since 2.1 */ -public class AcceptOnceUntilPurgedElementFilter implements ElementFilter { - +public class AcceptOnceUntilPurgedElementFilter implements CollectionFilter { + private final Log logger = LogFactory.getLog(this.getClass()); private final Queue seenItems; private final Object seenQueueMonitor = new Object(); - - public AcceptOnceUntilPurgedElementFilter(){ + + + public AcceptOnceUntilPurgedElementFilter() { this(Integer.MAX_VALUE); } - - public AcceptOnceUntilPurgedElementFilter(int maxCapacity){ - seenItems = new LinkedBlockingQueue(maxCapacity); + + public AcceptOnceUntilPurgedElementFilter(int maxCapacity) { + this.seenItems = new LinkedBlockingQueue(maxCapacity); + } + + + public Collection filter(Collection unfilteredElements) { + Assert.notNull(unfilteredElements, "'unfilteredElements' must not be null"); + List filteredElements = new ArrayList(); + if (unfilteredElements.size() > 0) { + for (T element : unfilteredElements) { + if (this.accept(element)) { + filteredElements.add(element); + } + } + } + return filteredElements; } private boolean accept(T item) { @@ -67,15 +86,4 @@ public class AcceptOnceUntilPurgedElementFilter implements ElementFilter { } } - public T filter(T unfilteredElement) { - Assert.notNull(unfilteredElement, "'unfilteredElement' must not be null"); - - if (this.accept(unfilteredElement)){ - return unfilteredElement; - } - else { - return null; - } - } - } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/ElementFilter.java b/spring-integration-core/src/main/java/org/springframework/integration/util/CollectionFilter.java similarity index 76% rename from spring-integration-core/src/main/java/org/springframework/integration/util/ElementFilter.java rename to spring-integration-core/src/main/java/org/springframework/integration/util/CollectionFilter.java index d0aaac9277..b1adc56362 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/util/ElementFilter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/CollectionFilter.java @@ -13,16 +13,20 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.integration.util; +import java.util.Collection; /** - * Base strategy for filtering out an element + * Base strategy for filtering out a subset of a Collection of elements. * * @author Oleg Zhurakousky + * @author Mark Fisher * @since 2.1 */ -public interface ElementFilter { +public interface CollectionFilter { + + Collection filter(Collection unfilteredElements); - T filter(T unfilteredElement); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java index 68da765a4d..49acd12c9f 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java @@ -13,9 +13,18 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.integration.resource; +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.assertTrue; + import java.io.File; +import java.util.Collection; +import java.util.Collections; import org.junit.Test; @@ -27,13 +36,8 @@ import org.springframework.integration.Message; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.endpoint.SourcePollingChannelAdapter; import org.springframework.integration.test.util.TestUtils; -import org.springframework.integration.util.ElementFilter; - -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.assertTrue; +import org.springframework.integration.util.CollectionFilter; +import org.springframework.util.CollectionUtils; /** * @author Oleg Zhurakousky @@ -114,19 +118,19 @@ public class ResourcePatternResolverParserTests { assertNull(message); } - - public static class OneItemAndNeverAgainResourceListFilter implements ElementFilter { + + + public static class OneItemAndNeverAgainResourceListFilter implements CollectionFilter { private volatile boolean once = false; - public Resource filter(Resource unfilteredElement) { - - if (!once){ + public Collection filter(Collection unfilteredResources) { + if (!once && !CollectionUtils.isEmpty(unfilteredResources)) { once = true; - return unfilteredElement; + return Collections.singletonList(unfilteredResources.iterator().next()); } return null; } - } + } From b7aabb70e5d33e61148209f9588c552ffaff827d Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Tue, 22 Nov 2011 13:36:31 -0500 Subject: [PATCH 3/4] renamed filter --- ...er.java => AcceptOnceUntilPurgedCollectionFilter.java} | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) rename spring-integration-core/src/main/java/org/springframework/integration/util/{AcceptOnceUntilPurgedElementFilter.java => AcceptOnceUntilPurgedCollectionFilter.java} (91%) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java b/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedCollectionFilter.java similarity index 91% rename from spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java rename to spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedCollectionFilter.java index 2369f23c9d..a76280be43 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedCollectionFilter.java @@ -30,14 +30,14 @@ import org.springframework.util.Assert; * An implementation of {@link CollectionFilter} which will queue all items that have been seen until * the queue reaches its capacity after which one item from the queue will be purged to make room for a * new item to be added. Note that however unlikely the removed item will now appear as unprocessed - * so it is highly recommended to move/delete resources which corresponds to the underlying items once processing + * so it is highly recommended to move/delete resources which correspond to the underlying items once processing * is done to eliminate duplicate processing. * * @author Oleg Zhurakousky * @author Mark Fisher * @since 2.1 */ -public class AcceptOnceUntilPurgedElementFilter implements CollectionFilter { +public class AcceptOnceUntilPurgedCollectionFilter implements CollectionFilter { private final Log logger = LogFactory.getLog(this.getClass()); @@ -46,11 +46,11 @@ public class AcceptOnceUntilPurgedElementFilter implements CollectionFilter(maxCapacity); } From ab6f3f1f270d797fa2b72cfba57d285a507128c2 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Tue, 22 Nov 2011 14:05:35 -0500 Subject: [PATCH 4/4] fixed polling period and renamed test class --- ...sts.java => ResourceInboundChannelAdapterParserTests.java} | 4 ++-- .../resource/ResourcePatternResolver-config-usagerf.xml | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) rename spring-integration-core/src/test/java/org/springframework/integration/resource/{ResourcePatternResolverParserTests.java => ResourceInboundChannelAdapterParserTests.java} (98%) diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourceInboundChannelAdapterParserTests.java similarity index 98% rename from spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java rename to spring-integration-core/src/test/java/org/springframework/integration/resource/ResourceInboundChannelAdapterParserTests.java index 49acd12c9f..2af6712e9a 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourceInboundChannelAdapterParserTests.java @@ -41,9 +41,9 @@ import org.springframework.util.CollectionUtils; /** * @author Oleg Zhurakousky - * + * @since 2.1 */ -public class ResourcePatternResolverParserTests { +public class ResourceInboundChannelAdapterParserTests { @Test public void testDefaultConfig(){ diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usagerf.xml b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usagerf.xml index 5ff664ebf9..3c71c30647 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usagerf.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usagerf.xml @@ -9,10 +9,10 @@ - + - +