RESOLVED - issue BATCH-76: ItemWriter for SQL batch updates

http://jira.springframework.org/browse/BATCH-76
This commit is contained in:
dsyer
2008-03-06 13:19:55 +00:00
parent 1184d30d54
commit 915997cd28
15 changed files with 823 additions and 4 deletions

View File

@@ -0,0 +1,227 @@
/*
* Copyright 2006-2007 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.io.support;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.util.HashSet;
import java.util.Iterator;
import java.util.Set;
import org.springframework.batch.item.ItemWriter;
import org.springframework.batch.item.exception.ClearFailedException;
import org.springframework.batch.item.exception.FlushFailedException;
import org.springframework.batch.repeat.RepeatContext;
import org.springframework.batch.repeat.synch.RepeatSynchronizationManager;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.dao.DataAccessException;
import org.springframework.jdbc.core.JdbcOperations;
import org.springframework.jdbc.core.PreparedStatementCallback;
import org.springframework.transaction.support.TransactionSynchronizationManager;
import org.springframework.util.Assert;
/**
* {@link ItemWriter} that uses the batching features from
* {@link PreparedStatement} if available and can take some rudimentary steps to
* locate a failure during a flush, and identify the items that failed. When one
* of those items is encountered again the batch is flushed aggressively so that
* the bad item is eventually identified and can be dealt with in isolation.<br/>
*
* The user must provide an SQL query and a special callback
* {@link ItemPreparedStatementSetter}, which is responsible for mapping the
* item to a PreparedStatement.
*
* @author Dave Syer
*
*/
public class BatchSqlUpdateItemWriter implements ItemWriter, InitializingBean {
/**
* Key for items processed in the current transaction {@link RepeatContext}.
*/
protected static final String ITEMS_PROCESSED = BatchSqlUpdateItemWriter.class.getName() + ".ITEMS_PROCESSED";
private Set failed = new HashSet();
private JdbcOperations jdbcTemplate;
private ItemPreparedStatementSetter preparedStatementSetter;
private String sql;
/**
* Public setter for the query string to execute on write. The parameters
* should correspond to those known to the
* {@link ItemPreparedStatementSetter}.
* @param sql the query to set
*/
public void setSql(String sql) {
this.sql = sql;
}
/**
* Public setter for the {@link ItemPreparedStatementSetter}.
* @param preparedStatementSetter the {@link ItemPreparedStatementSetter} to
* set
*/
public void setItemPreparedStatementSetter(ItemPreparedStatementSetter preparedStatementSetter) {
this.preparedStatementSetter = preparedStatementSetter;
}
/**
* Public setter for the {@link JdbcOperations}.
* @param jdbcTemplate the {@link JdbcOperations} to set
*/
public void setJdbcTemplate(JdbcOperations jdbcTemplate) {
this.jdbcTemplate = jdbcTemplate;
}
/**
* Check mandatory properties - there must be a delegate.
*
* @see org.springframework.dao.support.DaoSupport#initDao()
*/
public void afterPropertiesSet() throws Exception {
Assert.notNull(jdbcTemplate, "BatchSqlUpdateItemWriter requires an data source.");
Assert.notNull(preparedStatementSetter, "BatchSqlUpdateItemWriter requires a ItemPreparedStatementSetter");
}
/**
* Buffer the item in a transaction resource, but flush aggressively if the
* item was previously part of a failed chunk.
*
* @throws Exception
*
* @see org.springframework.batch.io.OutputSource#write(java.lang.Object)
*/
public void write(Object output) throws Exception {
bindTransactionResources();
getProcessed().add(output);
flushIfNecessary(output);
}
/**
* Accessor for the list of processed items in this transaction.
*
* @return the processed
*/
private Set getProcessed() {
Assert.state(TransactionSynchronizationManager.hasResource(ITEMS_PROCESSED),
"Processed items not bound to transaction.");
Set processed = (Set) TransactionSynchronizationManager.getResource(ITEMS_PROCESSED);
return processed;
}
/**
* Set up the {@link RepeatContext} as a transaction resource.
*
* @param context the context to set
*/
private void bindTransactionResources() {
if (TransactionSynchronizationManager.hasResource(ITEMS_PROCESSED)) {
return;
}
TransactionSynchronizationManager.bindResource(ITEMS_PROCESSED, new HashSet());
}
/**
* Remove the transaction resource associated with this context.
*/
private void unbindTransactionResources() {
if (!TransactionSynchronizationManager.hasResource(ITEMS_PROCESSED)) {
return;
}
TransactionSynchronizationManager.unbindResource(ITEMS_PROCESSED);
}
/**
* Accessor for the context property.
*
* @param output
*
* @return the context
*/
private void flushIfNecessary(Object output) throws Exception {
boolean flush;
synchronized (failed) {
flush = failed.contains(output);
}
if (flush) {
RepeatContext context = RepeatSynchronizationManager.getContext();
// Force early completion to commit aggressively if we encounter a
// failed item (from a failed chunk but we don't know which one was
// the problem).
context.setCompleteOnly();
// Flush now, so that if there is a failure this record can be
// skipped.
doFlush();
}
}
/**
* Flush the hibernate session from within a repeat context.
*/
private void doFlush() {
final Set processed = getProcessed();
try {
if (!processed.isEmpty()) {
jdbcTemplate.execute(sql, new PreparedStatementCallback() {
public Object doInPreparedStatement(PreparedStatement ps) throws SQLException, DataAccessException {
for (Iterator iterator = processed.iterator(); iterator.hasNext();) {
Object item = (Object) iterator.next();
preparedStatementSetter.setValues(item, ps);
ps.addBatch();
}
return ps.executeBatch();
}
});
}
}
catch (RuntimeException e) {
synchronized (failed) {
failed.addAll(processed);
}
throw e;
}
finally {
getProcessed().clear();
}
}
/**
* Unbind transaction resources, effectively clearing the item buffer.
*
* @see org.springframework.batch.item.ItemWriter#clear()
*/
public void clear() throws ClearFailedException {
unbindTransactionResources();
}
/**
* Flush the internal item buffer and record failures if there are any.
*
* @see org.springframework.batch.item.ItemWriter#flush()
*/
public void flush() throws FlushFailedException {
try {
doFlush();
}
finally {
unbindTransactionResources();
}
}
}

View File

@@ -0,0 +1,39 @@
/*
* Copyright 2006-2007 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.io.support;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import org.springframework.jdbc.core.RowMapper;
/**
* A convenient strategy for SQL updates, acting effectively as the inverse of
* {@link RowMapper}.
* @author Dave Syer
*
*/
public interface ItemPreparedStatementSetter {
/**
* Set parameter values on the given PreparedStatement as determined from
* the provided item.
* @param ps the PreparedStatement to invoke setter methods on
* @throws SQLException if a SQLException is encountered (i.e. there is no
* need to catch SQLException)
*/
void setValues(Object item, PreparedStatement ps) throws SQLException;
}

View File

@@ -0,0 +1,208 @@
/*
* Copyright 2006-2007 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.io.support;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import junit.framework.TestCase;
import org.easymock.MockControl;
import org.springframework.batch.repeat.RepeatContext;
import org.springframework.batch.repeat.context.RepeatContextSupport;
import org.springframework.batch.repeat.synch.RepeatSynchronizationManager;
import org.springframework.dao.DataAccessException;
import org.springframework.jdbc.UncategorizedSQLException;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.core.PreparedStatementCallback;
import org.springframework.transaction.support.TransactionSynchronizationManager;
/**
* @author Dave Syer
*
*/
public class BatchSqlUpdateItemWriterTests extends TestCase {
private BatchSqlUpdateItemWriter writer = new BatchSqlUpdateItemWriter();
private JdbcTemplate jdbcTemplate;
protected List list = new ArrayList();
private RepeatContext context = new RepeatContextSupport(null);
private PreparedStatement ps;
private MockControl control = MockControl.createControl(PreparedStatement.class);
/*
* (non-Javadoc)
* @see junit.framework.TestCase#setUp()
*/
protected void setUp() throws Exception {
ps = (PreparedStatement) control.getMock();
jdbcTemplate = new JdbcTemplate() {
public Object execute(String sql, PreparedStatementCallback action) throws DataAccessException {
list.add(sql);
try {
return action.doInPreparedStatement(ps);
}
catch (SQLException e) {
throw new UncategorizedSQLException("doInPreparedStatement", sql, e);
}
}
};
writer.setSql("SQL");
writer.setJdbcTemplate(jdbcTemplate);
writer.setItemPreparedStatementSetter(new ItemPreparedStatementSetter() {
public void setValues(Object item, PreparedStatement ps) throws SQLException {
list.add(item);
}
});
TransactionSynchronizationManager.bindResource(BatchSqlUpdateItemWriter.ITEMS_PROCESSED, new HashSet(
Collections.singleton("spam")));
RepeatSynchronizationManager.register(context);
}
/*
* (non-Javadoc)
* @see junit.framework.TestCase#tearDown()
*/
protected void tearDown() throws Exception {
Map map = TransactionSynchronizationManager.getResourceMap();
for (Iterator iterator = map.keySet().iterator(); iterator.hasNext();) {
String key = (String) iterator.next();
TransactionSynchronizationManager.unbindResource(key);
}
RepeatSynchronizationManager.clear();
}
/**
* Test method for
* {@link org.springframework.batch.io.support.BatchSqlUpdateItemWriter#afterPropertiesSet()}.
* @throws Exception
*/
public void testAfterPropertiesSet() throws Exception {
try {
writer.afterPropertiesSet();
}
catch (IllegalArgumentException e) {
// expected
String message = e.getMessage().toLowerCase();
assertTrue("Message does not contain 'query'.", message.indexOf("query") >= 0);
}
}
/**
* Test method for
* {@link org.springframework.batch.io.support.BatchSqlUpdateItemWriter#write(java.lang.Object)}.
* @throws Exception
*/
public void testWrite() throws Exception {
writer.setSql("foo");
writer.write("bar");
// Nothing happens till we flush
assertEquals(0, list.size());
}
/**
* Test method for
* {@link org.springframework.batch.io.support.BatchSqlUpdateItemWriter#clear()}.
*/
public void testClear() {
assertTrue(TransactionSynchronizationManager.hasResource(BatchSqlUpdateItemWriter.ITEMS_PROCESSED));
writer.clear();
assertFalse(TransactionSynchronizationManager.hasResource(BatchSqlUpdateItemWriter.ITEMS_PROCESSED));
}
/**
* Test method for
* {@link org.springframework.batch.io.support.BatchSqlUpdateItemWriter#flush()}.
*/
public void testFlush() {
assertTrue(TransactionSynchronizationManager.hasResource(BatchSqlUpdateItemWriter.ITEMS_PROCESSED));
writer.flush();
assertFalse(TransactionSynchronizationManager.hasResource(BatchSqlUpdateItemWriter.ITEMS_PROCESSED));
assertEquals(2, list.size());
assertTrue(list.contains("SQL"));
}
/**
* Test method for
* {@link org.springframework.batch.io.support.BatchSqlUpdateItemWriter#flush()}.
* @throws Exception
*/
public void testWriteAndFlush() throws Exception {
assertTrue(TransactionSynchronizationManager.hasResource(BatchSqlUpdateItemWriter.ITEMS_PROCESSED));
writer.write("bar");
writer.flush();
assertFalse(TransactionSynchronizationManager.hasResource(BatchSqlUpdateItemWriter.ITEMS_PROCESSED));
assertEquals(3, list.size());
assertTrue(list.contains("SQL"));
}
public void testFlushWithFailure() throws Exception{
tearDown();
try {
writer.flush();
fail("Expected IllegalStateException");
} catch (IllegalStateException e) {
assertTrue("Message should contain hint about transaction: "+e.getMessage(), e.getMessage().indexOf("transaction")>=0);
}
}
public void testWriteAndFlushWithFailure() throws Exception {
final RuntimeException ex = new RuntimeException("bar");
writer.setItemPreparedStatementSetter(new ItemPreparedStatementSetter() {
public void setValues(Object item, PreparedStatement ps) throws SQLException {
list.add(item);
throw ex;
}
});
ps.addBatch();
control.setVoidCallable();
control.expectAndReturn(ps.executeBatch(), new int[] {123});
control.replay();
writer.write("foo");
try {
writer.flush();
fail("Expected RuntimeException");
} catch (RuntimeException e) {
assertEquals("bar", e.getMessage());
}
assertFalse(TransactionSynchronizationManager.hasResource(BatchSqlUpdateItemWriter.ITEMS_PROCESSED));
assertEquals(2, list.size());
writer.setItemPreparedStatementSetter(new ItemPreparedStatementSetter() {
public void setValues(Object item, PreparedStatement ps) throws SQLException {
list.add(item);
}
});
writer.write("foo");
writer.flush();
control.verify();
assertEquals(4, list.size());
assertTrue(list.contains("SQL"));
assertTrue(list.contains("foo"));
assertTrue(context.isCompleteOnly());
}
}

View File

@@ -88,6 +88,7 @@ public class HibernateAwareItemWriterTests extends TestCase {
String key = (String) iterator.next();
TransactionSynchronizationManager.unbindResource(key);
}
RepeatSynchronizationManager.clear();
}
/**