You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@beam.apache.org by lc...@apache.org on 2016/10/03 15:33:55 UTC
[2/3] incubator-beam git commit: Add initial bucket stuff.
Add initial bucket stuff.
Project: http://git-wip-us.apache.org/repos/asf/incubator-beam/repo
Commit: http://git-wip-us.apache.org/repos/asf/incubator-beam/commit/ccd9bad1
Tree: http://git-wip-us.apache.org/repos/asf/incubator-beam/tree/ccd9bad1
Diff: http://git-wip-us.apache.org/repos/asf/incubator-beam/diff/ccd9bad1
Branch: refs/heads/master
Commit: ccd9bad18dcc0a1df7ebc82f43e89ddec838c037
Parents: 202acd1
Author: sammcveety <sa...@gmail.com>
Authored: Sat Sep 17 21:19:53 2016 -0400
Committer: Luke Cwik <lc...@google.com>
Committed: Mon Oct 3 08:20:12 2016 -0700
----------------------------------------------------------------------
.../java/org/apache/beam/sdk/util/GcsUtil.java | 67 +++++++++++++++++++-
1 file changed, 64 insertions(+), 3 deletions(-)
----------------------------------------------------------------------
http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/ccd9bad1/sdks/java/core/src/main/java/org/apache/beam/sdk/util/GcsUtil.java
----------------------------------------------------------------------
diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/GcsUtil.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/GcsUtil.java
index 41c372e..4befb1a 100644
--- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/GcsUtil.java
+++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/GcsUtil.java
@@ -349,16 +349,45 @@ public class GcsUtil {
}
/**
+ * Returns the project number of the project which owns this bucket.
+ * If the bucket exists, it must be accessible otherwise the permissions
+ * exception will be propagated.
+ */
+ public long bucketOwner(GcsPath path) throws IOException {
+ return getBucket(
+ path,
+ BACKOFF_FACTORY.backoff(),
+ Sleeper.DEFAULT).getProjectNumber();
+ }
+
+ /**
+ * Creates a bucket for the provided project or propagates an error.
+ */
+ public void createBucket(GcsPath path, long projectNumber) throws IOException {
+ return createBucket(
+ path,
+ projectNumber,
+ BACKOFF_FACTORY.backoff(),
+ Sleeper.DEFAULT);
+ }
+
+ /**
* Returns whether the GCS bucket exists. This will return false if the bucket
* is inaccessible due to permissions.
*/
@VisibleForTesting
boolean bucketExists(GcsPath path, BackOff backoff, Sleeper sleeper) throws IOException {
+ return getBucket(path, backoff, sleeper) != null;
+ }
+
+ @VisibleForTesting
+ @Nullable
+ Storage.Bucket getBucket(GcsPath path, BackOff backoff, Sleeper sleeper) throws IOException {
Storage.Buckets.Get getBucket =
storageClient.buckets().get(path.getBucket());
try {
- ResilientOperation.retry(
+ Storage.Bucket bucket = ResilientOperation.retry(
ResilientOperation.getGoogleRequestCallable(getBucket),
backoff,
new RetryDeterminer<IOException>() {
@@ -372,10 +401,11 @@ public class GcsUtil {
},
IOException.class,
sleeper);
- return true;
+
+ return bucket;
} catch (GoogleJsonResponseException e) {
if (errorExtractor.itemNotFound(e) || errorExtractor.accessDenied(e)) {
- return false;
+ return null;
}
throw e;
} catch (InterruptedException e) {
@@ -386,6 +416,37 @@ public class GcsUtil {
}
}
+ @VisibleForTesting
+ void createBucket(GcsPath path, long projectNumber, BackOff backoff, Sleeper sleeper)
+ throws IOException {
+ Storage.Buckets.Insert insertBucket =
+ storageClient.buckets().insert(path.getBucket());
+ insertBucket.setProject(String.valueOf(projectNumber));
+
+ try {
+ ResilientOperation.retry(
+ ResilientOperation.getGoogleRequestCallable(insertBucket),
+ backoff,
+ new RetryDeterminer<IOException>() {
+ @Override
+ public boolean shouldRetry(IOException e) {
+ if (errorExtractor.itemNotFound(e) || errorExtractor.accessDenied(e)) {
+ return false;
+ }
+ return RetryDeterminer.SOCKET_ERRORS.shouldRetry(e);
+ }
+ },
+ IOException.class,
+ sleeper);
+ return;
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new IOException(
+ String.format("Error while attempting to create bucket gs://%s for rproject %s",
+ path.getBucket(), projectNumber), e);
+ }
+ }
+
private static void executeBatches(List<BatchRequest> batches) throws IOException {
ListeningExecutorService executor = MoreExecutors.listeningDecorator(
MoreExecutors.getExitingExecutorService(