From b44bf47d622f6158672c1ebc3a20e360e0f79ad1 Mon Sep 17 00:00:00 2001 From: Christos Kapasakalidis Date: Thu, 6 Nov 2014 17:31:13 +0200 Subject: [PATCH] INTEXT-80 - FIX S3 ConnectionPoolTimeoutException JIRA: https://jira.spring.io/browse/INTEXT-80 When synchronizing to local directory, getObject was called without ever closing the InputStream that the SDK opens causing an "org.apache.http.conn.ConnectionPoolTimeoutException: Timeout waiting for connection" exception. To fix that, after synchronization, the inputStream of the s3Object is closed, causing the s3 client to release the connection. --- .../s3/InboundFileSynchronizationImpl.java | 179 ++++++++++-------- 1 file changed, 104 insertions(+), 75 deletions(-) diff --git a/spring-integration-aws/src/main/java/org/springframework/integration/aws/s3/InboundFileSynchronizationImpl.java b/spring-integration-aws/src/main/java/org/springframework/integration/aws/s3/InboundFileSynchronizationImpl.java index 50e6aa3..106ade7 100644 --- a/spring-integration-aws/src/main/java/org/springframework/integration/aws/s3/InboundFileSynchronizationImpl.java +++ b/spring-integration-aws/src/main/java/org/springframework/integration/aws/s3/InboundFileSynchronizationImpl.java @@ -29,6 +29,7 @@ import java.util.concurrent.locks.ReentrantLock; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; + import org.springframework.beans.factory.InitializingBean; import org.springframework.integration.aws.s3.core.AmazonS3Object; import org.springframework.integration.aws.s3.core.AmazonS3Operations; @@ -45,22 +46,31 @@ import org.springframework.util.StringUtils; * be checked against the * * @author Amol Nayak + * @author Christos Kapasakalidis * * @since 0.5 * */ -public class InboundFileSynchronizationImpl implements InboundFileSynchronizer,InitializingBean { +public class InboundFileSynchronizationImpl implements InboundFileSynchronizer, InitializingBean { private final Log logger = LogFactory.getLog(getClass()); public static final String CONTENT_MD5 = "Content-MD5"; + private final AmazonS3Operations client; - private volatile int maxObjectsPerBatch = 100; //default + + private volatile int maxObjectsPerBatch = 100; //default + private final InboundLocalFileOperations fileOperations; + private volatile FileNameFilter filter; + private volatile String fileWildcard; + private volatile String fileNameRegex; + private final Lock lock = new ReentrantLock(); + private volatile boolean acceptSubFolders; /** @@ -68,31 +78,30 @@ public class InboundFileSynchronizationImpl implements InboundFileSynchronizer,I * @param client */ public InboundFileSynchronizationImpl(AmazonS3Operations client, - InboundLocalFileOperations fileOperations) { - Assert.notNull(client,"AmazonS3Client should be non null"); - Assert.notNull(fileOperations,"fileOperations should be non null"); + InboundLocalFileOperations fileOperations) { + Assert.notNull(client, "AmazonS3Client should be non null"); + Assert.notNull(fileOperations, "fileOperations should be non null"); this.client = client; this.fileOperations = fileOperations; } - public void afterPropertiesSet() throws Exception { Assert.isTrue(!(StringUtils.hasText(fileWildcard) && StringUtils.hasText(fileNameRegex)), - "Only one of the file name wildcard string or file name regex can be specified"); + "Only one of the file name wildcard string or file name regex can be specified"); - if(StringUtils.hasText(fileWildcard)) { + if (StringUtils.hasText(fileWildcard)) { filter = new WildcardFileNameFilter(fileWildcard); } - else if(StringUtils.hasText(fileNameRegex)) { + else if (StringUtils.hasText(fileNameRegex)) { filter = new RegexFileNameFilter(fileNameRegex); } else { - filter = new AlwaysTrueFileNamefilter(); //Match all + filter = new AlwaysTrueFileNamefilter(); //Match all } - if(acceptSubFolders) { - ((AbstractFileNameFilter)filter).setAcceptSubFolders(true); + if (acceptSubFolders) { + ((AbstractFileNameFilter) filter).setAcceptSubFolders(true); fileOperations.setCreateDirectoriesIfRequired(true); } } @@ -103,83 +112,96 @@ public class InboundFileSynchronizationImpl implements InboundFileSynchronizer,I */ public void synchronizeToLocalDirectory(File localDirectory, String bucketName, String remoteFolder) { - if(!lock.tryLock()) { - if(logger.isInfoEnabled()) { - logger.info("Sync already in progess"); - } - //Prevent concurrent synchronization requests - return; + if (!lock.tryLock()) { + if (logger.isInfoEnabled()) { + logger.info("Sync already in progess"); + } + //Prevent concurrent synchronization requests + return; + } + + if (logger.isInfoEnabled()) { + logger.info("Starting sync with local directory"); + } + //Below sync can take long, above lock ensures only one thread is synchronizing + try { + if (remoteFolder != null && "/".equals(remoteFolder)) { + remoteFolder = null; } - if(logger.isInfoEnabled()) { - logger.info("Starting sync with local directory"); + //Set the remote folder for the filter + if (filter instanceof AbstractFileNameFilter) { + ((AbstractFileNameFilter) filter).setFolderName(remoteFolder); } - //Below sync can take long, above lock ensures only one thread is synchronizing - try { - if(remoteFolder != null && "/".equals(remoteFolder)) { - remoteFolder = null; - } - //Set the remote folder for the filter - if(filter instanceof AbstractFileNameFilter) { - ((AbstractFileNameFilter)filter).setFolderName(remoteFolder); - } - - String nextMarker = null; - do { - PaginatedObjectsView paginatedView = client.listObjects(bucketName, remoteFolder,nextMarker,maxObjectsPerBatch); - if(paginatedView == null) - break; //No files to sync - nextMarker = paginatedView.getNextMarker(); - List summaries = paginatedView.getObjectSummary(); - for(S3ObjectSummary summary:summaries) { - String key = summary.getKey(); - if(key.endsWith("/")) { - continue; - } - if(!filter.accept(key)) - continue; - //The folder is the root as the key is relative to bucket - AmazonS3Object s3Object = client.getObject(bucketName, "/", key); - synchronizeObjectWithFile(localDirectory,summary,s3Object); + String nextMarker = null; + do { + PaginatedObjectsView paginatedView = client.listObjects(bucketName, remoteFolder, nextMarker, maxObjectsPerBatch); + if (paginatedView == null) + break; //No files to sync + nextMarker = paginatedView.getNextMarker(); + List summaries = paginatedView.getObjectSummary(); + for (S3ObjectSummary summary : summaries) { + String key = summary.getKey(); + if (key.endsWith("/")) { + continue; + } + if (!filter.accept(key)) + continue; + //The folder is the root as the key is relative to bucket + AmazonS3Object s3Object = null; + try { + s3Object = client.getObject(bucketName, "/", key); + synchronizeObjectWithFile(localDirectory, summary, s3Object); + } + finally { + if (s3Object != null && s3Object.getInputStream() != null) { + s3Object.getInputStream().close(); + } } - } while(nextMarker != null); - - } finally { - lock.unlock(); - if(logger.isInfoEnabled()) { - logger.info("Sync completed"); } + } while (nextMarker != null); + + } + catch (IOException e) { + logger.error("Caught Exception while trying to close s3object InputStream", e); + } + finally { + lock.unlock(); + if (logger.isInfoEnabled()) { + logger.info("Sync completed"); } + } } + /** * Synchronizes the Object with the File on the local file system * @param localDirectory * @param summary */ - private void synchronizeObjectWithFile(File localDirectory,S3ObjectSummary summary, + private void synchronizeObjectWithFile(File localDirectory, S3ObjectSummary summary, AmazonS3Object s3Object) { //Get the complete object data String key = summary.getKey(); - if(key.endsWith("/")) { + if (key.endsWith("/")) { return; } int lastIndex = key.lastIndexOf("/"); String fileName = key.substring(lastIndex + 1); String filePath = localDirectory.getAbsolutePath(); - if(!filePath.endsWith(File.separator)) { + if (!filePath.endsWith(File.separator)) { filePath += File.separator; } File baseDirectory; - if(lastIndex > 0) { + if (lastIndex > 0) { //there could very well be previous '/' and thus nested sub folders String prefixKey = key.substring(0, lastIndex); String[] folders = prefixKey.split("/"); - if(folders.length > 0) { - for(String folder:folders) { + if (folders.length > 0) { + for (String folder : folders) { filePath = filePath + folder + File.separator; } //create the directory structure @@ -195,38 +217,41 @@ public class InboundFileSynchronizationImpl implements InboundFileSynchronizer,I } File file = new File(filePath + fileName); - if(!file.exists()) { + if (!file.exists()) { //File doesnt exist, write the contents to it try { - fileOperations.writeToFile(baseDirectory, fileName,s3Object.getInputStream()); - } catch (IOException e) { + fileOperations.writeToFile(baseDirectory, fileName, s3Object.getInputStream()); + } + catch (IOException e) { logger.error("Caught Exception while writing to file " + file.getAbsolutePath()); //continue with next file. } } else { //Synchronize a file that exists - if(!file.isFile()) { - if(logger.isWarnEnabled()) { + if (!file.isFile()) { + if (logger.isWarnEnabled()) { logger.warn("The file " + file.getAbsolutePath() + " is not a regular file, probably a directory, "); } return; } String eTag = summary.getETag(); String md5Hex = null; - if(isEtagMD5Hash(eTag)) { + if (isEtagMD5Hash(eTag)) { //Single thread upload try { md5Hex = encodeHex(getContentsMD5AsBytes(file)); - } catch (UnsupportedEncodingException e) { + } + catch (UnsupportedEncodingException e) { logger.error("Exception encountered while generating the MD5 hash for the file " + file.getAbsolutePath(), e); } - if(!eTag.equals(md5Hex)) { + if (!eTag.equals(md5Hex)) { //The local file is different than the one on S3, could be latest but we will still //sync this with the copy on S3 try { fileOperations.writeToFile(baseDirectory, fileName, s3Object.getInputStream()); - } catch (IOException e) { + } + catch (IOException e) { logger.error("Caught Exception while writing to file " + file.getAbsolutePath()); } } @@ -236,27 +261,30 @@ public class InboundFileSynchronizationImpl implements InboundFileSynchronizer,I //Get the MD5 hash from the headers Map userMetaData = s3Object.getUserMetaData(); String b64MD5 = userMetaData.get(CONTENT_MD5); - if(b64MD5 != null) { + if (b64MD5 != null) { //Need to convert to Hex from Base64 try { md5Hex = encodeHex(getContentsMD5AsBytes(file)); - } catch (UnsupportedEncodingException e) { + } + catch (UnsupportedEncodingException e) { logger.error("Exception encountered while generating the MD5 hash for the file " + file.getAbsolutePath(), e); } try { String remoteHexMD5 = new String( encodeHex( decodeBase64(b64MD5.getBytes("UTF-8")))); - if(!md5Hex.equals(remoteHexMD5)) { + if (!md5Hex.equals(remoteHexMD5)) { //Update only if the local file is not same as remote file try { fileOperations.writeToFile(baseDirectory, fileName, s3Object.getInputStream()); - } catch (IOException e) { + } + catch (IOException e) { logger.error("Caught Exception while writing to file " + file.getAbsolutePath()); } } - } catch (UnsupportedEncodingException e) { + } + catch (UnsupportedEncodingException e) { //Should never get this, suppress } } @@ -264,7 +292,8 @@ public class InboundFileSynchronizationImpl implements InboundFileSynchronizer,I //Forcefully update the file try { fileOperations.writeToFile(baseDirectory, fileName, s3Object.getInputStream()); - } catch (IOException e) { + } + catch (IOException e) { logger.error("Caught Exception while writing to file " + file.getAbsolutePath()); } } @@ -293,7 +322,7 @@ public class InboundFileSynchronizationImpl implements InboundFileSynchronizer,I */ public void setSynchronizingBatchSize(int batchSize) { - if(batchSize > 0) + if (batchSize > 0) this.maxObjectsPerBatch = batchSize; }