BATCH-1882: Add support for AMQP backed ItemReader / ItemWriter implementations

Add Spring AMQP and Spring Rabbit dependencies. Update
EasyMock from 2.3 to 2.4 to allow for concrete class mocking.
This commit is contained in:
Chris Schaefer
2012-07-30 22:41:51 -04:00
committed by Dave Syer
parent 7476ac1628
commit 1607014886
8 changed files with 316 additions and 3 deletions

View File

@@ -25,7 +25,7 @@
<dependency org="org.hsqldb" name="com.springsource.org.hsqldb" rev="1.8.0.9" conf="test->runtime,provided"/>
<dependency org="org.apache.commons" name="com.springsource.org.apache.commons.io" rev="1.4.0" conf="test->runtime,provided"/>
<dependency org="org.apache.commons" name="com.springsource.org.apache.commons.lang" rev="2.1.0" conf="compile->runtime,provided"/>
<dependency org="org.easymock" name="com.springsource.org.easymock" rev="2.3.0" conf="test->runtime,provided"/>
<dependency org="org.easymock" name="com.springsource.org.easymock" rev="2.4.0" conf="test->runtime,provided"/>
<dependency org="org.junit" name="com.springsource.org.junit" rev="4.4.0" conf="test->runtime,provided"/>
<dependency org="org.aspectj" name="com.springsource.org.aspectj.runtime" rev="1.5.4" conf="compile->runtime,provided"/>
<dependency org="org.aspectj" name="com.springsource.org.aspectj.weaver" rev="1.5.4" conf="compile->runtime,provided"/>

View File

@@ -21,7 +21,7 @@
<dependencies>
<dependency org="org.junit" name="com.springsource.org.junit" rev="4.4.0" conf="test->runtime,provided"/>
<dependency org="org.easymock" name="com.springsource.org.easymock" rev="2.3.0" conf="test->runtime,provided"/>
<dependency org="org.easymock" name="com.springsource.org.easymock" rev="2.4.0" conf="test->runtime,provided"/>
<dependency org="org.aspectj" name="com.springsource.org.aspectj.runtime" rev="1.5.4" conf="test->runtime,provided"/>
<dependency org="org.aspectj" name="com.springsource.org.aspectj.weaver" rev="1.5.4" conf="test->runtime,provided"/>
<dependency org="net.sourceforge.cglib" name="com.springsource.net.sf.cglib" rev="2.1.3" conf="optional->runtime,provided"/>

View File

@@ -43,6 +43,10 @@
<groupId>org.easymock</groupId>
<artifactId>easymock</artifactId>
</dependency>
<dependency>
<groupId>org.easymock</groupId>
<artifactId>easymockclassextension</artifactId>
</dependency>
<dependency>
<groupId>org.aspectj</groupId>
<artifactId>aspectjrt</artifactId>
@@ -195,6 +199,16 @@
<artifactId>woodstox-core-asl</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.amqp</groupId>
<artifactId>spring-amqp</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.amqp</groupId>
<artifactId>spring-rabbit</artifactId>
<optional>true</optional>
</dependency>
</dependencies>
<reporting>
<plugins>

View File

@@ -0,0 +1,62 @@
/*
* Copyright 2012 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.batch.item.amqp;
import org.springframework.amqp.core.AmqpTemplate;
import org.springframework.amqp.core.Message;
import org.springframework.batch.item.ItemReader;
import org.springframework.util.Assert;
/**
* <p>
* AMQP {@link ItemReader} implementation using an {@link AmqpTemplate} to
* receive and/or convert messages.
* </p>
*
* @author Chris Schaefer
*/
public class AmqpItemReader<T> implements ItemReader<T> {
private final AmqpTemplate amqpTemplate;
private Class<? extends T> itemType;
public AmqpItemReader(final AmqpTemplate amqpTemplate) {
Assert.notNull(amqpTemplate, "AmpqTemplate must not be null");
this.amqpTemplate = amqpTemplate;
}
@SuppressWarnings("unchecked")
public T read() {
if (itemType != null && itemType.isAssignableFrom(Message.class)) {
return (T) amqpTemplate.receive();
}
Object result = amqpTemplate.receiveAndConvert();
if (itemType != null && result != null) {
Assert.state(itemType.isAssignableFrom(result.getClass()),
"Received message payload of wrong type: expected [" + itemType + "]");
}
return (T) result;
}
public void setItemType(Class<? extends T> itemType) {
Assert.notNull(itemType, "Item type cannot be null");
this.itemType = itemType;
}
}

View File

@@ -0,0 +1,55 @@
/*
* Copyright 2012 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.batch.item.amqp;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.amqp.core.AmqpTemplate;
import org.springframework.batch.item.ItemWriter;
import org.springframework.util.Assert;
import java.util.List;
/**
* <p>
* AMQP {@link ItemWriter} implementation using an {@link AmqpTemplate} to
* send messages. Messages will be sent to the nameless exchange if not specified
* on the provided {@link AmqpTemplate}.
* </p>
*
* @author Chris Schaefer
*/
public class AmqpItemWriter<T> implements ItemWriter<T> {
private final AmqpTemplate amqpTemplate;
private final Log log = LogFactory.getLog(getClass());
public AmqpItemWriter(final AmqpTemplate amqpTemplate) {
Assert.notNull(amqpTemplate, "AmpqTemplate must not be null");
this.amqpTemplate = amqpTemplate;
}
public void write(final List<? extends T> items) throws Exception {
if (log.isDebugEnabled()) {
log.debug("Writing to AMQP with " + items.size() + " items.");
}
for (T item : items) {
amqpTemplate.convertAndSend(item);
}
}
}

View File

@@ -0,0 +1,109 @@
/*
* Copyright 2012 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.batch.item.amqp;
import org.easymock.classextension.EasyMock;
import org.junit.Test;
import org.springframework.amqp.core.AmqpTemplate;
import org.springframework.amqp.core.Message;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;
/**
* <p>
* Test cases around {@link AmqpItemReader}.
* </p>
*
* @author Chris Schaefer
*/
public class AmqpItemReaderTests {
@Test(expected = IllegalArgumentException.class)
public void testNullAmqpTemplate() {
new AmqpItemReader<String>(null);
}
@Test
public void testNoItemType() {
final AmqpTemplate amqpTemplate = EasyMock.createMock(AmqpTemplate.class);
EasyMock.expect(amqpTemplate.receiveAndConvert()).andReturn("foo");
EasyMock.replay(amqpTemplate);
final AmqpItemReader<String> amqpItemReader = new AmqpItemReader<String>(amqpTemplate);
assertEquals("foo", amqpItemReader.read());
EasyMock.verify(amqpTemplate);
}
@Test
public void testNonMessageItemType() {
final AmqpTemplate amqpTemplate = EasyMock.createMock(AmqpTemplate.class);
EasyMock.expect(amqpTemplate.receiveAndConvert()).andReturn("foo");
EasyMock.replay(amqpTemplate);
final AmqpItemReader<String> amqpItemReader = new AmqpItemReader<String>(amqpTemplate);
amqpItemReader.setItemType(String.class);
assertEquals("foo", amqpItemReader.read());
EasyMock.verify(amqpTemplate);
}
@Test
public void testMessageItemType() {
final AmqpTemplate amqpTemplate = EasyMock.createMock(AmqpTemplate.class);
final Message message = EasyMock.createMock(Message.class);
EasyMock.expect(amqpTemplate.receive()).andReturn(message);
EasyMock.replay(amqpTemplate, message);
final AmqpItemReader<Message> amqpItemReader = new AmqpItemReader<Message>(amqpTemplate);
amqpItemReader.setItemType(Message.class);
assertEquals(message, amqpItemReader.read());
EasyMock.verify(amqpTemplate);
}
@Test
public void testTypeMismatch() {
final AmqpTemplate amqpTemplate = EasyMock.createMock(AmqpTemplate.class);
EasyMock.expect(amqpTemplate.receiveAndConvert()).andReturn("foo");
EasyMock.replay(amqpTemplate);
final AmqpItemReader<Integer> amqpItemReader = new AmqpItemReader<Integer>(amqpTemplate);
amqpItemReader.setItemType(Integer.class);
try {
amqpItemReader.read();
fail("Expected IllegalStateException");
} catch (IllegalStateException e) {
assertTrue(e.getMessage().contains("wrong type"));
}
EasyMock.verify(amqpTemplate);
}
@Test(expected = IllegalArgumentException.class)
public void testNullItemType() {
final AmqpTemplate amqpTemplate = EasyMock.createMock(AmqpTemplate.class);
final AmqpItemReader<String> amqpItemReader = new AmqpItemReader<String>(amqpTemplate);
amqpItemReader.setItemType(null);
}
}

View File

@@ -0,0 +1,56 @@
/*
* Copyright 2012 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.batch.item.amqp;
import org.easymock.EasyMock;
import org.junit.Test;
import org.springframework.amqp.core.AmqpTemplate;
import java.util.Arrays;
/**
* <p>
* Test cases around {@link AmqpItemWriter}.
* </p>
*
* @author Chris Schaefer
*/
public class AmqpItemWriterTests {
@Test(expected = IllegalArgumentException.class)
public void testNullAmqpTemplate() {
new AmqpItemWriter<String>(null);
}
@Test
public void voidTestWrite() throws Exception {
AmqpTemplate amqpTemplate = EasyMock.createMock(AmqpTemplate.class);
amqpTemplate.convertAndSend("foo");
EasyMock.expectLastCall();
amqpTemplate.convertAndSend("bar");
EasyMock.expectLastCall();
EasyMock.replay(amqpTemplate);
AmqpItemWriter<String> amqpItemWriter = new AmqpItemWriter<String>(amqpTemplate);
amqpItemWriter.write(Arrays.asList("foo", "bar"));
EasyMock.verify(amqpTemplate);
}
}

View File

@@ -32,6 +32,7 @@
<spring.oxm.group>org.springframework.ws</spring.oxm.group>
<spring.oxm.artifact>spring-oxm-tiger</spring.oxm.artifact>
<spring.oxm.version>1.5.8</spring.oxm.version>
<spring.amqp.version>1.0.0.RELEASE</spring.amqp.version>
<junit.version>4.4</junit.version>
<bundlor.version>1.0.0.RELEASE</bundlor.version>
<dependency.locations.enabled>false</dependency.locations.enabled>
@@ -418,9 +419,15 @@
<dependency>
<groupId>org.easymock</groupId>
<artifactId>easymock</artifactId>
<version>2.3</version>
<version>2.4</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.easymock</groupId>
<artifactId>easymockclassextension</artifactId>
<version>2.4</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.apache.geronimo.specs</groupId>
<artifactId>geronimo-jms_1.1_spec</artifactId>
@@ -709,6 +716,16 @@
<version>1.1.2</version>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.amqp</groupId>
<artifactId>spring-amqp</artifactId>
<version>${spring.amqp.version}</version>
</dependency>
<dependency>
<groupId>org.springframework.amqp</groupId>
<artifactId>spring-rabbit</artifactId>
<version>${spring.amqp.version}</version>
</dependency>
</dependencies>
</dependencyManagement>
<distributionManagement>