diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcExecutionContextDao.java b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcExecutionContextDao.java index d6fd74cde..e55b791e4 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcExecutionContextDao.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcExecutionContextDao.java @@ -30,6 +30,8 @@ import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Map.Entry; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import org.springframework.batch.core.JobExecution; import org.springframework.batch.core.StepExecution; @@ -112,6 +114,8 @@ public class JdbcExecutionContextDao extends AbstractJdbcBatchMetadataDao implem private ExecutionContextSerializer serializer = new DefaultExecutionContextSerializer(); + private final Lock lock = new ReentrantLock(); + /** * Setter for {@link Serializer} implementation * @param serializer {@link ExecutionContextSerializer} instance to use. @@ -191,7 +195,8 @@ public class JdbcExecutionContextDao extends AbstractJdbcBatchMetadataDao implem public void updateExecutionContext(final StepExecution stepExecution) { // Attempt to prevent concurrent modification errors by blocking here if // someone is already trying to do it. - synchronized (stepExecution) { + this.lock.lock(); + try { Long executionId = stepExecution.getId(); ExecutionContext executionContext = stepExecution.getExecutionContext(); Assert.notNull(executionId, "ExecutionId must not be null."); @@ -201,6 +206,9 @@ public class JdbcExecutionContextDao extends AbstractJdbcBatchMetadataDao implem persistSerializedContext(executionId, serializedContext, UPDATE_STEP_EXECUTION_CONTEXT); } + finally { + this.lock.unlock(); + } } @Override diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcJobExecutionDao.java b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcJobExecutionDao.java index 3c260dd43..17ace20ad 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcJobExecutionDao.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcJobExecutionDao.java @@ -26,6 +26,8 @@ import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -158,6 +160,8 @@ public class JdbcJobExecutionDao extends AbstractJdbcBatchMetadataDao implements private ConfigurableConversionService conversionService; + private final Lock lock = new ReentrantLock(); + public JdbcJobExecutionDao() { DefaultConversionService conversionService = new DefaultConversionService(); conversionService.addConverter(new DateToStringConverter()); @@ -278,7 +282,8 @@ public class JdbcJobExecutionDao extends AbstractJdbcBatchMetadataDao implements Assert.notNull(jobExecution.getVersion(), "JobExecution version cannot be null. JobExecution must be saved before it can be updated"); - synchronized (jobExecution) { + this.lock.lock(); + try { Integer version = jobExecution.getVersion() + 1; String exitDescription = jobExecution.getExitStatus().getExitDescription(); @@ -323,6 +328,9 @@ public class JdbcJobExecutionDao extends AbstractJdbcBatchMetadataDao implements jobExecution.incrementVersion(); } + finally { + this.lock.unlock(); + } } @Nullable diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcStepExecutionDao.java b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcStepExecutionDao.java index c486a067b..5c53a6f4a 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcStepExecutionDao.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcStepExecutionDao.java @@ -26,6 +26,8 @@ import java.util.Arrays; import java.util.Collection; import java.util.Iterator; import java.util.List; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -119,6 +121,8 @@ public class JdbcStepExecutionDao extends AbstractJdbcBatchMetadataDao implement private DataFieldMaxValueIncrementer stepExecutionIncrementer; + private final Lock lock = new ReentrantLock(); + /** * Public setter for the exit message length in database. Do not set this if you * haven't modified the schema. @@ -256,7 +260,8 @@ public class JdbcStepExecutionDao extends AbstractJdbcBatchMetadataDao implement // Attempt to prevent concurrent modification errors by blocking here if // someone is already trying to do it. - synchronized (stepExecution) { + this.lock.lock(); + try { Integer version = stepExecution.getVersion() + 1; Timestamp startTime = stepExecution.getStartTime() == null ? null @@ -289,6 +294,9 @@ public class JdbcStepExecutionDao extends AbstractJdbcBatchMetadataDao implement stepExecution.incrementVersion(); } + finally { + this.lock.unlock(); + } } /** diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/data/AbstractPaginatedDataItemReader.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/data/AbstractPaginatedDataItemReader.java index 05580de0e..c7982e506 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/data/AbstractPaginatedDataItemReader.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/data/AbstractPaginatedDataItemReader.java @@ -22,6 +22,8 @@ import org.springframework.lang.Nullable; import org.springframework.util.Assert; import java.util.Iterator; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; /** * A base class that handles basic reading logic based on the paginated semantics of @@ -44,7 +46,7 @@ public abstract class AbstractPaginatedDataItemReader extends AbstractItemCou protected Iterator results; - private final Object lock = new Object(); + private final Lock lock = new ReentrantLock(); /** * The number of items to be read with each page. @@ -59,7 +61,8 @@ public abstract class AbstractPaginatedDataItemReader extends AbstractItemCou @Override protected T doRead() throws Exception { - synchronized (lock) { + this.lock.lock(); + try { if (results == null || !results.hasNext()) { results = doPageRead(); @@ -78,6 +81,9 @@ public abstract class AbstractPaginatedDataItemReader extends AbstractItemCou return null; } } + finally { + this.lock.unlock(); + } } /** @@ -101,7 +107,8 @@ public abstract class AbstractPaginatedDataItemReader extends AbstractItemCou @Override protected void jumpToItem(int itemLastIndex) throws Exception { - synchronized (lock) { + this.lock.lock(); + try { page = itemLastIndex / pageSize; int current = itemLastIndex % pageSize; @@ -111,6 +118,9 @@ public abstract class AbstractPaginatedDataItemReader extends AbstractItemCou initialPage.next(); } } + finally { + this.lock.unlock(); + } } } diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/data/RepositoryItemReader.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/data/RepositoryItemReader.java index cab3b1ee7..4494b2d8c 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/data/RepositoryItemReader.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/data/RepositoryItemReader.java @@ -19,6 +19,8 @@ import java.lang.reflect.InvocationTargetException; import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -98,7 +100,7 @@ public class RepositoryItemReader extends AbstractItemCountingItemStreamItemR private volatile List results; - private final Object lock = new Object(); + private final Lock lock = new ReentrantLock(); private String methodName; @@ -162,7 +164,8 @@ public class RepositoryItemReader extends AbstractItemCountingItemStreamItemR @Override protected T doRead() throws Exception { - synchronized (lock) { + this.lock.lock(); + try { boolean nextPageNeeded = (results != null && current >= results.size()); if (results == null || nextPageNeeded) { @@ -192,14 +195,21 @@ public class RepositoryItemReader extends AbstractItemCountingItemStreamItemR return null; } } + finally { + this.lock.unlock(); + } } @Override protected void jumpToItem(int itemLastIndex) throws Exception { - synchronized (lock) { + this.lock.lock(); + try { page = itemLastIndex / pageSize; current = itemLastIndex % pageSize; } + finally { + this.lock.unlock(); + } } /** @@ -236,11 +246,15 @@ public class RepositoryItemReader extends AbstractItemCountingItemStreamItemR @Override protected void doClose() throws Exception { - synchronized (lock) { + this.lock.lock(); + try { current = 0; page = 0; results = null; } + finally { + this.lock.unlock(); + } } private Sort convertToSort(Map sorts) { diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/database/AbstractPagingItemReader.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/database/AbstractPagingItemReader.java index 45504e7af..e66dc2bc2 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/database/AbstractPagingItemReader.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/database/AbstractPagingItemReader.java @@ -16,6 +16,8 @@ package org.springframework.batch.item.database; import java.util.List; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -58,7 +60,7 @@ public abstract class AbstractPagingItemReader extends AbstractItemCountingIt protected volatile List results; - private final Object lock = new Object(); + private final Lock lock = new ReentrantLock(); public AbstractPagingItemReader() { setName(ClassUtils.getShortName(AbstractPagingItemReader.class)); @@ -101,7 +103,8 @@ public abstract class AbstractPagingItemReader extends AbstractItemCountingIt @Override protected T doRead() throws Exception { - synchronized (lock) { + this.lock.lock(); + try { if (results == null || current >= pageSize) { @@ -126,6 +129,9 @@ public abstract class AbstractPagingItemReader extends AbstractItemCountingIt } } + finally { + this.lock.unlock(); + } } @@ -142,22 +148,30 @@ public abstract class AbstractPagingItemReader extends AbstractItemCountingIt @Override protected void doClose() throws Exception { - synchronized (lock) { + this.lock.lock(); + try { initialized = false; current = 0; page = 0; results = null; } + finally { + this.lock.unlock(); + } } @Override protected void jumpToItem(int itemIndex) throws Exception { - synchronized (lock) { + this.lock.lock(); + try { page = itemIndex / pageSize; current = itemIndex % pageSize; } + finally { + this.lock.unlock(); + } if (logger.isDebugEnabled()) { logger.debug("Jumping to page " + getPage() + " and index " + current); diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/database/ExtendedConnectionDataSourceProxy.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/database/ExtendedConnectionDataSourceProxy.java index 9fe20ac92..18f8a968a 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/database/ExtendedConnectionDataSourceProxy.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/database/ExtendedConnectionDataSourceProxy.java @@ -24,6 +24,8 @@ import java.lang.reflect.Proxy; import java.sql.Connection; import java.sql.SQLException; import java.sql.SQLFeatureNotSupportedException; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import java.util.logging.Logger; import javax.sql.DataSource; @@ -92,7 +94,7 @@ public class ExtendedConnectionDataSourceProxy implements SmartDataSource, Initi private boolean borrowedConnection = false; /** Synchronization monitor for the shared Connection */ - private final Object connectionMonitor = new Object(); + private final Lock connectionMonitor = new ReentrantLock(); /** * No arg constructor for use when configured using JavaBean style. @@ -143,12 +145,16 @@ public class ExtendedConnectionDataSourceProxy implements SmartDataSource, Initi * @param connection the {@link Connection} that close suppression is requested for */ public void startCloseSuppression(Connection connection) { - synchronized (this.connectionMonitor) { + this.connectionMonitor.lock(); + try { closeSuppressedConnection = connection; if (TransactionSynchronizationManager.isActualTransactionActive()) { borrowedConnection = true; } } + finally { + this.connectionMonitor.unlock(); + } } /** @@ -156,24 +162,36 @@ public class ExtendedConnectionDataSourceProxy implements SmartDataSource, Initi * off for */ public void stopCloseSuppression(Connection connection) { - synchronized (this.connectionMonitor) { + this.connectionMonitor.lock(); + try { closeSuppressedConnection = null; borrowedConnection = false; } + finally { + this.connectionMonitor.unlock(); + } } @Override public Connection getConnection() throws SQLException { - synchronized (this.connectionMonitor) { + this.connectionMonitor.lock(); + try { return initConnection(null, null); } + finally { + this.connectionMonitor.unlock(); + } } @Override public Connection getConnection(String username, String password) throws SQLException { - synchronized (this.connectionMonitor) { + this.connectionMonitor.lock(); + try { return initConnection(username, password); } + finally { + this.connectionMonitor.unlock(); + } } private boolean completeCloseCall(Connection connection) { diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/file/SimpleBinaryBufferedReaderFactory.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/file/SimpleBinaryBufferedReaderFactory.java index 10020f87b..dba352eff 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/file/SimpleBinaryBufferedReaderFactory.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/file/SimpleBinaryBufferedReaderFactory.java @@ -20,6 +20,8 @@ import java.io.IOException; import java.io.InputStreamReader; import java.io.Reader; import java.io.UnsupportedEncodingException; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import org.springframework.core.io.Resource; @@ -60,12 +62,15 @@ public class SimpleBinaryBufferedReaderFactory implements BufferedReaderFactory * usual plain text conventions. * * @author Dave Syer + * @author Mahmoud Ben Hassine * */ private static final class BinaryBufferedReader extends BufferedReader { private final String ending; + private final Lock lock = new ReentrantLock(); + private BinaryBufferedReader(Reader in, String ending) { super(in); this.ending = ending; @@ -76,7 +81,8 @@ public class SimpleBinaryBufferedReaderFactory implements BufferedReaderFactory StringBuilder buffer; - synchronized (lock) { + this.lock.lock(); + try { int next = read(); if (next == -1) { @@ -92,6 +98,9 @@ public class SimpleBinaryBufferedReaderFactory implements BufferedReaderFactory buffer.append(candidateEnding); } + finally { + this.lock.unlock(); + } if (buffer != null && buffer.length() > 0) { return buffer.toString(); diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/support/SynchronizedItemStreamReader.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/support/SynchronizedItemStreamReader.java index 1b33e9d3c..8cc1322d8 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/support/SynchronizedItemStreamReader.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/support/SynchronizedItemStreamReader.java @@ -15,6 +15,9 @@ */ package org.springframework.batch.item.support; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; + import org.springframework.batch.item.ExecutionContext; import org.springframework.batch.item.ItemStreamReader; import org.springframework.batch.item.NonTransientResourceException; @@ -35,6 +38,7 @@ import org.springframework.util.Assert; * Here is the motivation behind this class: https://stackoverflow.com/a/20002493/2910265 * * @author Matthew Ouyang + * @author Mahmoud Ben Hassine * @since 3.0.4 * @param type of object being read */ @@ -42,6 +46,8 @@ public class SynchronizedItemStreamReader implements ItemStreamReader, Ini private ItemStreamReader delegate; + private final Lock lock = new ReentrantLock(); + public void setDelegate(ItemStreamReader delegate) { this.delegate = delegate; } @@ -50,8 +56,14 @@ public class SynchronizedItemStreamReader implements ItemStreamReader, Ini * This delegates to the read method of the delegate */ @Nullable - public synchronized T read() throws Exception { - return this.delegate.read(); + public T read() throws Exception { + this.lock.lock(); + try { + return this.delegate.read(); + } + finally { + this.lock.unlock(); + } } public void close() { diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/support/SynchronizedItemStreamWriter.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/support/SynchronizedItemStreamWriter.java index 56fc6710b..ad69b89a5 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/support/SynchronizedItemStreamWriter.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/support/SynchronizedItemStreamWriter.java @@ -15,6 +15,9 @@ */ package org.springframework.batch.item.support; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; + import org.springframework.batch.item.Chunk; import org.springframework.batch.item.ExecutionContext; import org.springframework.batch.item.ItemStreamException; @@ -47,6 +50,8 @@ public class SynchronizedItemStreamWriter implements ItemStreamWriter, Ini private ItemStreamWriter delegate; + private final Lock lock = new ReentrantLock(); + /** * Set the delegate {@link ItemStreamWriter}. * @param delegate the delegate to set @@ -59,8 +64,14 @@ public class SynchronizedItemStreamWriter implements ItemStreamWriter, Ini * This method delegates to the {@code write} method of the {@code delegate}. */ @Override - public synchronized void write(Chunk items) throws Exception { - this.delegate.write(items); + public void write(Chunk items) throws Exception { + this.lock.lock(); + try { + this.delegate.write(items); + } + finally { + this.lock.unlock(); + } } @Override diff --git a/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/StagingItemReader.java b/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/StagingItemReader.java index 3472f00d2..d2d5f1324 100644 --- a/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/StagingItemReader.java +++ b/spring-batch-samples/src/main/java/org/springframework/batch/sample/common/StagingItemReader.java @@ -20,6 +20,8 @@ import java.io.InputStream; import java.io.ObjectInputStream; import java.util.Iterator; import java.util.List; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import javax.sql.DataSource; @@ -50,7 +52,7 @@ public class StagingItemReader private StepExecution stepExecution; - private final Object lock = new Object(); + private final Lock lock = new ReentrantLock(); private volatile boolean initialized = false; @@ -75,7 +77,8 @@ public class StagingItemReader private List retrieveKeys() { - synchronized (lock) { + this.lock.lock(); + try { return jdbcTemplate.query( @@ -86,6 +89,9 @@ public class StagingItemReader stepExecution.getJobExecution().getJobId(), StagingItemWriter.NEW); } + finally { + this.lock.unlock(); + } }