fix for poll where no rows available
This commit is contained in:
@@ -62,7 +62,7 @@ public class JdbcPollingChannelAdapter implements MessageSource<Object>,
|
||||
|
||||
private volatile boolean updatePerRow = false;
|
||||
|
||||
private volatile String updateQuery;
|
||||
private volatile String updateSql;
|
||||
|
||||
private volatile SqlParamterSourceFactory sqlParameterSourceFactoryForUpdate = new DefaultSqlParamterSourceFactory();
|
||||
|
||||
@@ -91,8 +91,8 @@ public class JdbcPollingChannelAdapter implements MessageSource<Object>,
|
||||
this.rowMapper = rowMapper;
|
||||
}
|
||||
|
||||
public void setUpdateQuery(String updateQuery) {
|
||||
this.updateQuery = updateQuery;
|
||||
public void setUpdatesql(String updateSql) {
|
||||
this.updateSql = updateSql;
|
||||
}
|
||||
|
||||
public void setUpdatePerRow(boolean updatePerRow){
|
||||
@@ -127,6 +127,9 @@ public class JdbcPollingChannelAdapter implements MessageSource<Object>,
|
||||
} else {
|
||||
payload = pollAndUpdate();
|
||||
}
|
||||
if(payload == null){
|
||||
return null;
|
||||
}
|
||||
return MessageBuilder.withPayload(payload).build();
|
||||
}
|
||||
|
||||
@@ -138,7 +141,12 @@ public class JdbcPollingChannelAdapter implements MessageSource<Object>,
|
||||
payload = this.jdbcOperations.queryForList(this.selectQuery,
|
||||
this.sqlQueryParameterSource);
|
||||
}
|
||||
if (updateQuery != null) {
|
||||
|
||||
if(payload.size() < 1){
|
||||
payload = null;
|
||||
}
|
||||
|
||||
if (payload != null && updateSql != null) {
|
||||
if (this.updatePerRow) {
|
||||
for (Object row : payload) {
|
||||
executeUpdateQuery(row);
|
||||
@@ -157,9 +165,9 @@ public class JdbcPollingChannelAdapter implements MessageSource<Object>,
|
||||
|
||||
updateParamaterSource = this.sqlParameterSourceFactoryForUpdate
|
||||
.createParamterSource(obj);
|
||||
this.jdbcOperations.update(this.updateQuery, updateParamaterSource);
|
||||
this.jdbcOperations.update(this.updateSql, updateParamaterSource);
|
||||
} else {
|
||||
this.jdbcOperations.update(this.updateQuery);
|
||||
this.jdbcOperations.update(this.updateSql);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
package org.springframework.integration.jdbc;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
import static org.junit.Assert.*;
|
||||
|
||||
import java.sql.ResultSet;
|
||||
import java.sql.SQLException;
|
||||
@@ -23,119 +22,145 @@ import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseType;
|
||||
*/
|
||||
public class JdbcPollingChannelAdapterIntegrationTest {
|
||||
|
||||
EmbeddedDatabase embeddedDatabase;
|
||||
EmbeddedDatabase embeddedDatabase;
|
||||
|
||||
SimpleJdbcTemplate jdbcTemplate;
|
||||
SimpleJdbcTemplate jdbcTemplate;
|
||||
|
||||
@Before
|
||||
public void setUp(){
|
||||
EmbeddedDatabaseBuilder builder = new EmbeddedDatabaseBuilder();
|
||||
builder.setType(EmbeddedDatabaseType.DERBY).addScript("classpath:org/springframework/integration/jdbc/pollingChannelAdapterIntegrationTest.sql");
|
||||
this.embeddedDatabase = builder.build();
|
||||
this.jdbcTemplate = new SimpleJdbcTemplate(this.embeddedDatabase);
|
||||
}
|
||||
@Before
|
||||
public void setUp() {
|
||||
EmbeddedDatabaseBuilder builder = new EmbeddedDatabaseBuilder();
|
||||
builder
|
||||
.setType(EmbeddedDatabaseType.DERBY)
|
||||
.addScript(
|
||||
"classpath:org/springframework/integration/jdbc/pollingChannelAdapterIntegrationTest.sql");
|
||||
this.embeddedDatabase = builder.build();
|
||||
this.jdbcTemplate = new SimpleJdbcTemplate(this.embeddedDatabase);
|
||||
}
|
||||
|
||||
|
||||
@After
|
||||
public void tearDown(){
|
||||
this.embeddedDatabase.shutdown();
|
||||
}
|
||||
@After
|
||||
public void tearDown() {
|
||||
this.embeddedDatabase.shutdown();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSimplePollForListOfMapsNoUpdate(){
|
||||
JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter(this.embeddedDatabase, "select * from item");
|
||||
this.jdbcTemplate.update("insert into item values(1,2)") ;
|
||||
Message<Object> message = adapter.receive();
|
||||
Object payload = message.getPayload();
|
||||
assertTrue("Wrong payload type", payload instanceof List);
|
||||
List rows = (List)payload;
|
||||
assertEquals("Wrong number of elements" , 1, rows.size());
|
||||
assertTrue("Returned row not a map", rows.get(0) instanceof Map);
|
||||
Map<String,Object> row = (Map<String, Object>)rows.get(0);
|
||||
assertEquals("Wrong id", 1, row.get("id"));
|
||||
assertEquals("Wrong status", 2, row.get("status"));
|
||||
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
public void testSimplePollForListWithRowMapperNoUpdate(){
|
||||
JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter(this.embeddedDatabase, "select * from item");
|
||||
adapter.setRowMapper(new ItemRowMapper());
|
||||
this.jdbcTemplate.update("insert into item values(1,2)") ;
|
||||
Message<Object> message = adapter.receive();
|
||||
Object payload = message.getPayload();
|
||||
List rows = (List)payload;
|
||||
assertEquals("Wrong number of elements" , 1, rows.size());
|
||||
assertTrue("Wrong payload type", rows.get(0) instanceof Item);
|
||||
Item item = (Item)rows.get(0);
|
||||
assertEquals("Wrong id", 1, item.getId());
|
||||
assertEquals("Wrong status", 2, item.getStatus());
|
||||
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSimplePollForListOfMapsNoUpdate() {
|
||||
JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter(
|
||||
this.embeddedDatabase, "select * from item");
|
||||
this.jdbcTemplate.update("insert into item values(1,2)");
|
||||
Message<Object> message = adapter.receive();
|
||||
Object payload = message.getPayload();
|
||||
assertTrue("Wrong payload type", payload instanceof List);
|
||||
List rows = (List) payload;
|
||||
assertEquals("Wrong number of elements", 1, rows.size());
|
||||
assertTrue("Returned row not a map", rows.get(0) instanceof Map);
|
||||
Map<String, Object> row = (Map<String, Object>) rows.get(0);
|
||||
assertEquals("Wrong id", 1, row.get("id"));
|
||||
assertEquals("Wrong status", 2, row.get("status"));
|
||||
|
||||
@Test
|
||||
public void testSimplePollForListWithRowMapperAndOneUpdate(){
|
||||
JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter(this.embeddedDatabase, "select * from item where status=2");
|
||||
adapter.setUpdateQuery("update item set status = 10 where id in (:idList)");
|
||||
adapter.setRowMapper(new ItemRowMapper());
|
||||
|
||||
this.jdbcTemplate.update("insert into item values(1,2)") ;
|
||||
this.jdbcTemplate.update("insert into item values(2,2)") ;
|
||||
|
||||
Message<Object> message = adapter.receive();
|
||||
Object payload = message.getPayload();
|
||||
List rows = (List)payload;
|
||||
assertEquals("Wrong number of elements" , 2, rows.size());
|
||||
assertTrue("Wrong payload type", rows.get(0) instanceof Item);
|
||||
Item item = (Item)rows.get(0);
|
||||
assertEquals("Wrong id", 1, item.getId());
|
||||
assertEquals("Wrong status", 2, item.getStatus());
|
||||
|
||||
int countOfStatusTwo = this.jdbcTemplate.queryForInt("select count(*) from item where status = 2");
|
||||
assertEquals("Status not updated incorect number of rows with status 2", 0, countOfStatusTwo);
|
||||
|
||||
int countOfStatusTen = this.jdbcTemplate.queryForInt("select count(*) from item where status = 10");
|
||||
assertEquals("Status not updated incorect number of rows with status 10", 2, countOfStatusTen);
|
||||
|
||||
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSimplePollForListWithRowMapperAndUpdatePerRow(){
|
||||
JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter(this.embeddedDatabase, "select * from item where status=2");
|
||||
adapter.setUpdateQuery("update item set status = 10 where id = :id");
|
||||
adapter.setUpdatePerRow(true);
|
||||
adapter.setRowMapper(new ItemRowMapper());
|
||||
|
||||
this.jdbcTemplate.update("insert into item values(1,2)") ;
|
||||
this.jdbcTemplate.update("insert into item values(2,2)") ;
|
||||
|
||||
Message<Object> message = adapter.receive();
|
||||
Object payload = message.getPayload();
|
||||
List rows = (List)payload;
|
||||
assertEquals("Wrong number of elements" , 2, rows.size());
|
||||
assertTrue("Wrong payload type", rows.get(0) instanceof Item);
|
||||
Item item = (Item)rows.get(0);
|
||||
assertEquals("Wrong id", 1, item.getId());
|
||||
assertEquals("Wrong status", 2, item.getStatus());
|
||||
|
||||
int countOfStatusTwo = this.jdbcTemplate.queryForInt("select count(*) from item where status = 2");
|
||||
assertEquals("Status not updated incorect number of rows with status 2", 0, countOfStatusTwo);
|
||||
|
||||
int countOfStatusTen = this.jdbcTemplate.queryForInt("select count(*) from item where status = 10");
|
||||
assertEquals("Status not updated incorect number of rows with status 10", 2, countOfStatusTen);
|
||||
|
||||
|
||||
}
|
||||
|
||||
private static class Item{
|
||||
|
||||
|
||||
private int id;
|
||||
|
||||
private int status;
|
||||
}
|
||||
|
||||
|
||||
|
||||
@Test
|
||||
public void testSimplePollForListWithRowMapperNoUpdate() {
|
||||
JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter(
|
||||
this.embeddedDatabase, "select * from item");
|
||||
adapter.setRowMapper(new ItemRowMapper());
|
||||
this.jdbcTemplate.update("insert into item values(1,2)");
|
||||
Message<Object> message = adapter.receive();
|
||||
Object payload = message.getPayload();
|
||||
List rows = (List) payload;
|
||||
assertEquals("Wrong number of elements", 1, rows.size());
|
||||
assertTrue("Wrong payload type", rows.get(0) instanceof Item);
|
||||
Item item = (Item) rows.get(0);
|
||||
assertEquals("Wrong id", 1, item.getId());
|
||||
assertEquals("Wrong status", 2, item.getStatus());
|
||||
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSimplePollForListWithRowMapperAndOneUpdate() {
|
||||
JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter(
|
||||
this.embeddedDatabase, "select * from item where status=2");
|
||||
adapter
|
||||
.setUpdatesql("update item set status = 10 where id in (:idList)");
|
||||
adapter.setRowMapper(new ItemRowMapper());
|
||||
|
||||
this.jdbcTemplate.update("insert into item values(1,2)");
|
||||
this.jdbcTemplate.update("insert into item values(2,2)");
|
||||
|
||||
Message<Object> message = adapter.receive();
|
||||
Object payload = message.getPayload();
|
||||
List rows = (List) payload;
|
||||
assertEquals("Wrong number of elements", 2, rows.size());
|
||||
assertTrue("Wrong payload type", rows.get(0) instanceof Item);
|
||||
Item item = (Item) rows.get(0);
|
||||
assertEquals("Wrong id", 1, item.getId());
|
||||
assertEquals("Wrong status", 2, item.getStatus());
|
||||
|
||||
int countOfStatusTwo = this.jdbcTemplate
|
||||
.queryForInt("select count(*) from item where status = 2");
|
||||
assertEquals(
|
||||
"Status not updated incorect number of rows with status 2", 0,
|
||||
countOfStatusTwo);
|
||||
|
||||
int countOfStatusTen = this.jdbcTemplate
|
||||
.queryForInt("select count(*) from item where status = 10");
|
||||
assertEquals(
|
||||
"Status not updated incorect number of rows with status 10", 2,
|
||||
countOfStatusTen);
|
||||
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSimplePollForListWithRowMapperAndUpdatePerRow() {
|
||||
JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter(
|
||||
this.embeddedDatabase, "select * from item where status=2");
|
||||
adapter.setUpdatesql("update item set status = 10 where id = :id");
|
||||
adapter.setUpdatePerRow(true);
|
||||
adapter.setRowMapper(new ItemRowMapper());
|
||||
|
||||
this.jdbcTemplate.update("insert into item values(1,2)");
|
||||
this.jdbcTemplate.update("insert into item values(2,2)");
|
||||
|
||||
Message<Object> message = adapter.receive();
|
||||
Object payload = message.getPayload();
|
||||
List rows = (List) payload;
|
||||
assertEquals("Wrong number of elements", 2, rows.size());
|
||||
assertTrue("Wrong payload type", rows.get(0) instanceof Item);
|
||||
Item item = (Item) rows.get(0);
|
||||
assertEquals("Wrong id", 1, item.getId());
|
||||
assertEquals("Wrong status", 2, item.getStatus());
|
||||
|
||||
int countOfStatusTwo = this.jdbcTemplate
|
||||
.queryForInt("select count(*) from item where status = 2");
|
||||
assertEquals(
|
||||
"Status not updated incorect number of rows with status 2", 0,
|
||||
countOfStatusTwo);
|
||||
|
||||
int countOfStatusTen = this.jdbcTemplate
|
||||
.queryForInt("select count(*) from item where status = 10");
|
||||
assertEquals(
|
||||
"Status not updated incorect number of rows with status 10", 2,
|
||||
countOfStatusTen);
|
||||
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEmptyPoll() {
|
||||
JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter(
|
||||
this.embeddedDatabase, "select * from item");
|
||||
Message<Object> message = adapter.receive();
|
||||
assertNull("Message received when no rows in table", message);
|
||||
|
||||
|
||||
}
|
||||
|
||||
private static class Item {
|
||||
|
||||
private int id;
|
||||
|
||||
private int status;
|
||||
|
||||
public int getId() {
|
||||
return id;
|
||||
@@ -152,22 +177,17 @@ public class JdbcPollingChannelAdapterIntegrationTest {
|
||||
public void setStatus(int status) {
|
||||
this.status = status;
|
||||
}
|
||||
}
|
||||
|
||||
private static class ItemRowMapper implements RowMapper<Item>{
|
||||
}
|
||||
|
||||
private static class ItemRowMapper implements RowMapper<Item> {
|
||||
|
||||
|
||||
public Item mapRow(ResultSet rs, int rowNum) throws SQLException {
|
||||
Item item = new Item();
|
||||
item.setId(rs.getInt(1));
|
||||
item.setStatus(rs.getInt(2));
|
||||
return item;
|
||||
}
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
Reference in New Issue
Block a user