Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 6 additions & 2 deletions build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -36,13 +36,17 @@ allprojects {
implementation 'com.google.guava:guava:21.0'
implementation 'com.google.code.findbugs:jsr305:3.0.2'

// AWS Services
// AWS Services - SDK v1 (for non-S3 services)
implementation 'com.amazonaws:aws-java-sdk-core:latest.release'
implementation 'com.amazonaws:aws-java-sdk-s3:latest.release'
implementation 'com.amazonaws:aws-java-sdk-sns:latest.release'
implementation 'com.amazonaws:aws-java-sdk-ec2:latest.release'
implementation 'com.amazonaws:aws-java-sdk-autoscaling:latest.release'
implementation 'com.amazonaws:aws-java-sdk-sts:latest.release'

// AWS Services - SDK v2 (only for S3)
implementation platform('software.amazon.awssdk:bom:2.20.0')
implementation 'software.amazon.awssdk:s3'
implementation 'software.amazon.awssdk:sts' // Needed for S3 role assumption

implementation 'com.google.inject:guice:4.2.2'
implementation 'com.google.inject.extensions:guice-servlet:4.2.2'
Expand Down
20 changes: 10 additions & 10 deletions priam/src/main/java/com/netflix/priam/aws/IAMCredential.java
Original file line number Diff line number Diff line change
Expand Up @@ -16,31 +16,31 @@
*/
package com.netflix.priam.aws;

import com.amazonaws.auth.AWSCredentials;
import com.amazonaws.auth.AWSCredentialsProvider;
import com.amazonaws.auth.InstanceProfileCredentialsProvider;
import com.netflix.priam.cred.ICredential;
import software.amazon.awssdk.auth.credentials.AwsCredentials;
import software.amazon.awssdk.auth.credentials.AwsCredentialsProvider;
import software.amazon.awssdk.auth.credentials.InstanceProfileCredentialsProvider;

public class IAMCredential implements ICredential {
private final InstanceProfileCredentialsProvider iamCredProvider;

public IAMCredential() {
this.iamCredProvider = InstanceProfileCredentialsProvider.getInstance();
this.iamCredProvider = InstanceProfileCredentialsProvider.create();
}

public String getAccessKeyId() {
return iamCredProvider.getCredentials().getAWSAccessKeyId();
return iamCredProvider.resolveCredentials().accessKeyId();
}

public String getSecretAccessKey() {
return iamCredProvider.getCredentials().getAWSSecretKey();
return iamCredProvider.resolveCredentials().secretAccessKey();
}

public AWSCredentials getCredentials() {
return iamCredProvider.getCredentials();
public AwsCredentials getCredentials() {
return iamCredProvider.resolveCredentials();
}

public AWSCredentialsProvider getAwsCredentialProvider() {
public AwsCredentialsProvider getAwsCredentialProvider() {
return iamCredProvider;
}
}
}
113 changes: 72 additions & 41 deletions priam/src/main/java/com/netflix/priam/aws/S3FileSystem.java
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,6 @@
*/
package com.netflix.priam.aws;

import com.amazonaws.services.s3.AmazonS3Client;
import com.amazonaws.services.s3.S3ResponseMetadata;
import com.amazonaws.services.s3.model.*;
import com.google.common.base.Preconditions;
import com.netflix.priam.aws.auth.IS3Credential;
import com.netflix.priam.backup.AbstractBackupPath;
Expand Down Expand Up @@ -50,6 +47,10 @@
import org.apache.commons.io.IOUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import software.amazon.awssdk.core.sync.RequestBody;
import software.amazon.awssdk.regions.Region;
import software.amazon.awssdk.services.s3.S3Client;
import software.amazon.awssdk.services.s3.model.*;

/** Implementation of IBackupFileSystem for S3 */
@Singleton
Expand All @@ -70,9 +71,9 @@ public S3FileSystem(
DynamicRateLimiter dynamicRateLimiter) {
super(pathProvider, compress, config, backupMetrics, backupNotificationMgr);
s3Client =
AmazonS3Client.builder()
.withCredentials(cred.getAwsCredentialProvider())
.withRegion(instanceInfo.getRegion())
S3Client.builder()
.credentialsProvider(cred.getAwsCredentialProvider())
.region(Region.of(instanceInfo.getRegion()))
.build();
this.dynamicRateLimiter = dynamicRateLimiter;
}
Expand Down Expand Up @@ -104,19 +105,25 @@ protected void downloadFileImpl(AbstractBackupPath path, String suffix)
}
}

private ObjectMetadata getObjectMetadata(File file) {
ObjectMetadata ret = new ObjectMetadata();
long lastModified = file.lastModified();

if (lastModified != 0) {
ret.addUserMetadata("local-modification-time", Long.toString(lastModified));
}
private PutObjectRequest.Builder getObjectMetadataBuilder(File file, String bucket, String key) {
PutObjectRequest.Builder builder = PutObjectRequest.builder()
.bucket(bucket)
.key(key);

long lastModified = file.lastModified();
long fileSize = file.length();
if (fileSize != 0) {
ret.addUserMetadata("local-size", Long.toString(fileSize));

if (lastModified != 0 || fileSize != 0) {
java.util.Map<String, String> metadata = new java.util.HashMap<>();
if (lastModified != 0) {
metadata.put("local-modification-time", Long.toString(lastModified));
}
if (fileSize != 0) {
metadata.put("local-size", Long.toString(fileSize));
}
builder.metadata(metadata);
}
return ret;
return builder;
}

private long uploadMultipart(AbstractBackupPath path, Instant target)
Expand All @@ -128,12 +135,29 @@ private long uploadMultipart(AbstractBackupPath path, Instant target)
if (logger.isDebugEnabled())
logger.debug("Uploading to {}/{} with chunk size {}", prefix, remotePath, chunkSize);
File localFile = localPath.toFile();
InitiateMultipartUploadRequest initRequest =
new InitiateMultipartUploadRequest(prefix, remotePath)
.withObjectMetadata(getObjectMetadata(localFile));
String uploadId = s3Client.initiateMultipartUpload(initRequest).getUploadId();

CreateMultipartUploadRequest.Builder initRequestBuilder = CreateMultipartUploadRequest.builder()
.bucket(prefix)
.key(remotePath);

// Add metadata
long lastModified = localFile.lastModified();
long fileSize = localFile.length();
java.util.Map<String, String> metadata = new java.util.HashMap<>();
if (lastModified != 0) {
metadata.put("local-modification-time", Long.toString(lastModified));
}
if (fileSize != 0) {
metadata.put("local-size", Long.toString(fileSize));
}
if (!metadata.isEmpty()) {
initRequestBuilder.metadata(metadata);
}

CreateMultipartUploadRequest initRequest = initRequestBuilder.build();
String uploadId = s3Client.createMultipartUpload(initRequest).uploadId();
DataPart part = new DataPart(prefix, remotePath, uploadId);
List<PartETag> partETags = Collections.synchronizedList(new ArrayList<>());
List<CompletedPart> partETags = Collections.synchronizedList(new ArrayList<>());

try (InputStream in = new FileInputStream(localFile)) {
Iterator<byte[]> chunks = new ChunkedStream(in, chunkSize, path.getCompression());
Expand All @@ -155,15 +179,10 @@ private long uploadMultipart(AbstractBackupPath path, Instant target)
executor.sleepTillEmpty();
logger.info("{} done. part count: {} expected: {}", localFile, partsPut.get(), partNum);
Preconditions.checkState(partNum == partETags.size(), "part count mismatch");
CompleteMultipartUploadResult resultS3MultiPartUploadComplete =
CompleteMultipartUploadResponse resultS3MultiPartUploadComplete =
new S3PartUploader(s3Client, part, partETags).completeUpload();
checkSuccessfulUpload(resultS3MultiPartUploadComplete, localPath);

if (logger.isDebugEnabled()) {
final S3ResponseMetadata info = s3Client.getCachedResponseMetadata(initRequest);
logger.debug("Request Id: {}, Host Id: {}", info.getRequestId(), info.getHostId());
}

return compressedFileSize;
} catch (Exception e) {
new S3PartUploader(s3Client, part, partETags).abortUpload();
Expand All @@ -182,10 +201,10 @@ protected long uploadFileImpl(AbstractBackupPath path, Instant target)
dynamicRateLimiter.acquire(path, target, chunk.length);
}
try {
new BoundedExponentialRetryCallable<PutObjectResult>(1000, 10000, 5) {
new BoundedExponentialRetryCallable<PutObjectResponse>(1000, 10000, 5) {
@Override
public PutObjectResult retriableCall() {
return s3Client.putObject(generatePut(path, chunk));
public PutObjectResponse retriableCall() {
return s3Client.putObject(generatePut(path, chunk), RequestBody.fromBytes(chunk));
}
}.call();
} catch (Exception e) {
Expand All @@ -196,18 +215,30 @@ public PutObjectResult retriableCall() {

private PutObjectRequest generatePut(AbstractBackupPath path, byte[] chunk) {
File localFile = Paths.get(path.getBackupFile().getAbsolutePath()).toFile();
ObjectMetadata metadata = getObjectMetadata(localFile);
metadata.setContentLength(chunk.length);
PutObjectRequest put =
new PutObjectRequest(
config.getBackupPrefix(),
path.getRemotePath(),
new ByteArrayInputStream(chunk),
metadata);

PutObjectRequest.Builder builder = PutObjectRequest.builder()
.bucket(config.getBackupPrefix())
.key(path.getRemotePath())
.contentLength((long) chunk.length);

// Add metadata
long lastModified = localFile.lastModified();
long fileSize = localFile.length();
java.util.Map<String, String> metadata = new java.util.HashMap<>();
if (lastModified != 0) {
metadata.put("local-modification-time", Long.toString(lastModified));
}
if (fileSize != 0) {
metadata.put("local-size", Long.toString(fileSize));
}
if (!metadata.isEmpty()) {
builder.metadata(metadata);
}

if (config.addMD5ToBackupUploads()) {
put.getMetadata().setContentMD5(SystemUtils.toBase64(SystemUtils.md5(chunk)));
builder.contentMD5(SystemUtils.toBase64(SystemUtils.md5(chunk)));
}
return put;
return builder.build();
}

private byte[] getFileContents(AbstractBackupPath path) throws BackupRestoreException {
Expand All @@ -224,4 +255,4 @@ private byte[] getFileContents(AbstractBackupPath path) throws BackupRestoreExce
throw new BackupRestoreException("Error reading file: " + localFile.getName(), e);
}
}
}
}
62 changes: 40 additions & 22 deletions priam/src/main/java/com/netflix/priam/aws/S3FileSystemBase.java
Original file line number Diff line number Diff line change
Expand Up @@ -13,13 +13,6 @@
*/
package com.netflix.priam.aws;

import com.amazonaws.AmazonClientException;
import com.amazonaws.services.s3.AmazonS3;
import com.amazonaws.services.s3.model.BucketLifecycleConfiguration;
import com.amazonaws.services.s3.model.BucketLifecycleConfiguration.Rule;
import com.amazonaws.services.s3.model.CompleteMultipartUploadResult;
import com.amazonaws.services.s3.model.DeleteObjectsRequest;
import com.amazonaws.services.s3.model.lifecycle.*;
import com.google.common.collect.Lists;
import com.google.common.util.concurrent.RateLimiter;
import com.netflix.priam.backup.AbstractBackupPath;
Expand All @@ -39,11 +32,21 @@
import javax.inject.Provider;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import software.amazon.awssdk.core.exception.SdkException;
import software.amazon.awssdk.services.s3.S3Client;
import software.amazon.awssdk.services.s3.model.BucketLifecycleConfiguration;
import software.amazon.awssdk.services.s3.model.CompleteMultipartUploadResponse;
import software.amazon.awssdk.services.s3.model.Delete;
import software.amazon.awssdk.services.s3.model.DeleteObjectsRequest;
import software.amazon.awssdk.services.s3.model.HeadObjectRequest;
import software.amazon.awssdk.services.s3.model.HeadObjectResponse;
import software.amazon.awssdk.services.s3.model.LifecycleRule;
import software.amazon.awssdk.services.s3.model.ObjectIdentifier;

public abstract class S3FileSystemBase extends AbstractFileSystem {
private static final int MAX_CHUNKS = 9995; // 10K is AWS limit, minus a small buffer
private static final Logger logger = LoggerFactory.getLogger(S3FileSystemBase.class);
AmazonS3 s3Client;
S3Client s3Client;
final IConfiguration config;
final ICompression compress;
final BlockingSubmitThreadPoolExecutor executor;
Expand Down Expand Up @@ -88,45 +91,55 @@ public void configChangeListener() {
objectExistLimiter.getRate());
}

private AmazonS3 getS3Client() {
private S3Client getS3Client() {
return s3Client;
}

/*
* A means to change the default handle to the S3 client.
*/
public void setS3Client(AmazonS3 client) {
public void setS3Client(S3Client client) {
s3Client = client;
}

void checkSuccessfulUpload(
CompleteMultipartUploadResult resultS3MultiPartUploadComplete, Path localPath)
CompleteMultipartUploadResponse resultS3MultiPartUploadComplete, Path localPath)
throws BackupRestoreException {
if (null != resultS3MultiPartUploadComplete
&& null != resultS3MultiPartUploadComplete.getETag()) {
&& null != resultS3MultiPartUploadComplete.eTag()) {
logger.info(
"Uploaded file: {}, object eTag: {}",
localPath,
resultS3MultiPartUploadComplete.getETag());
resultS3MultiPartUploadComplete.eTag());
} else {
throw new BackupRestoreException(
"Error uploading file as ETag or CompleteMultipartUploadResult is NULL -"
"Error uploading file as ETag or CompleteMultipartUploadResponse is NULL -"
+ localPath);
}
}

@Override
public long getFileSize(String remotePath) throws BackupRestoreException {
return s3Client.getObjectMetadata(getShard(), remotePath).getContentLength();
HeadObjectRequest request = HeadObjectRequest.builder()
.bucket(getShard())
.key(remotePath)
.build();
HeadObjectResponse response = s3Client.headObject(request);
return response.contentLength();
}

@Override
protected boolean doesRemoteFileExist(Path remotePath) {
objectExistLimiter.acquire();
boolean exists = false;
try {
exists = s3Client.doesObjectExist(getShard(), remotePath.toString());
} catch (AmazonClientException ex) {
HeadObjectRequest request = HeadObjectRequest.builder()
.bucket(getShard())
.key(remotePath.toString())
.build();
s3Client.headObject(request);
exists = true;
} catch (SdkException ex) {
// No point throwing this exception up.
logger.error(
"Exception while checking existence of object: {}. Error: {}",
Expand All @@ -152,16 +165,21 @@ public void deleteFiles(List<Path> remotePaths) throws BackupRestoreException {
if (remotePaths.isEmpty()) return;

try {
List<DeleteObjectsRequest.KeyVersion> keys =
List<ObjectIdentifier> keys =
remotePaths
.stream()
.map(
remotePath ->
new DeleteObjectsRequest.KeyVersion(
remotePath.toString()))
ObjectIdentifier.builder()
.key(remotePath.toString())
.build())
.collect(Collectors.toList());
s3Client.deleteObjects(
new DeleteObjectsRequest(getShard()).withKeys(keys).withQuiet(true));
Delete delete = Delete.builder().objects(keys).quiet(true).build();
DeleteObjectsRequest request = DeleteObjectsRequest.builder()
.bucket(getShard())
.delete(delete)
.build();
s3Client.deleteObjects(request);
logger.info("Deleted {} objects from S3", remotePaths.size());
} catch (Exception e) {
logger.error(
Expand Down
Loading
Loading