diff --git a/pom.xml b/pom.xml
index 038c18e5a..7a93f6ccd 100644
--- a/pom.xml
+++ b/pom.xml
@@ -20,7 +20,7 @@
io.cdap.plugin
google-cloud
- 0.24.6
+ 0.24.7-SNAPSHOT
Google Cloud Plugins
jar
Plugins for Google Big Query
diff --git a/src/main/java/io/cdap/plugin/gcp/bigquery/source/BigQuerySourceUtils.java b/src/main/java/io/cdap/plugin/gcp/bigquery/source/BigQuerySourceUtils.java
index aa44e48f9..98904b74e 100644
--- a/src/main/java/io/cdap/plugin/gcp/bigquery/source/BigQuerySourceUtils.java
+++ b/src/main/java/io/cdap/plugin/gcp/bigquery/source/BigQuerySourceUtils.java
@@ -87,30 +87,43 @@ public static String getOrCreateBucket(Configuration configuration,
// By default, this option is false, meaning the job can not delete the bucket. So enable it only when bucket name
// is not provided.
configuration.setBoolean("fs.gs.bucket.delete.enable", true);
- GCPUtils.createBucket(storage, bucket, dataset.getLocation(), cmekKeyName);
+ createBucket(storage, bucket, dataset, cmekKeyName);
} else if (storage != null && storage.get(bucket) == null) {
- try {
- GCPUtils.createBucket(storage, bucket, dataset.getLocation(), cmekKeyName);
- } catch (StorageException e) {
- if (e.getCode() == 409) {
- // A conflict means the bucket already exists
- // This most likely means multiple stages in the same pipeline are trying to create the same bucket.
- // Ignore this and move on, since all that matters is that the bucket exists.
- return bucket;
- }
- String errorMessage = String.format("Unable to create Cloud Storage bucket '%s' in the same "
- + "location ('%s') as BigQuery dataset '%s'. " + "Please use a bucket "
- + "that is in the same location as the dataset. For more details, see %s",
- bucket, dataset.getLocation(), dataset.getDatasetId().getDataset(),
- GCPUtils.GCS_SUPPORTED_DOC_URL);
- throw ErrorUtils.getProgramFailureException(new ErrorCategory(ErrorCategory.ErrorCategoryEnum.PLUGIN),
- errorMessage, e.getMessage(), ErrorType.USER, true, ErrorCodeType.HTTP,
- String.valueOf(e.getCode()), GCPUtils.GCS_SUPPORTED_DOC_URL, e);
- }
+ createBucket(storage, bucket, dataset, cmekKeyName);
}
return bucket;
}
+ private static void createBucket(@Nullable Storage storage,
+ String bucket,
+ Dataset dataset,
+ @Nullable CryptoKeyName cmekKeyName) {
+ if (storage == null) {
+ return;
+ }
+ try {
+ GCPUtils.createBucket(storage, bucket, dataset.getLocation(), cmekKeyName);
+ } catch (StorageException e) {
+ if (e.getCode() == 409) {
+ // A conflict means the bucket already exists.
+ // This most likely means multiple stages in the same pipeline are trying to create the same bucket,
+ // or a retry occurred after a successful bucket creation.
+ // Ignore this and move on, since all that matters is that the bucket exists.
+ LOG.debug("Bucket '{}' already exists, ignoring 409 Conflict: {}", bucket, e.getMessage());
+ return;
+ }
+ String datasetName = dataset.getDatasetId() != null ? dataset.getDatasetId().getDataset() : "";
+ String errorMessage = String.format("Unable to create Cloud Storage bucket '%s' in the same "
+ + "location ('%s') as BigQuery dataset '%s'. " + "Please use a bucket "
+ + "that is in the same location as the dataset. For more details, see %s",
+ bucket, dataset.getLocation(), datasetName,
+ GCPUtils.GCS_SUPPORTED_DOC_URL);
+ throw ErrorUtils.getProgramFailureException(new ErrorCategory(ErrorCategory.ErrorCategoryEnum.PLUGIN),
+ errorMessage, e.getMessage(), ErrorType.USER, true, ErrorCodeType.HTTP,
+ String.valueOf(e.getCode()), GCPUtils.GCS_SUPPORTED_DOC_URL, e);
+ }
+ }
+
/**
* Sets up service account credentials into supplied Hadoop configuration.
*
diff --git a/src/test/java/io/cdap/plugin/gcp/bigquery/source/BigQuerySourceUtilsTest.java b/src/test/java/io/cdap/plugin/gcp/bigquery/source/BigQuerySourceUtilsTest.java
index ad9236061..6efaf5090 100644
--- a/src/test/java/io/cdap/plugin/gcp/bigquery/source/BigQuerySourceUtilsTest.java
+++ b/src/test/java/io/cdap/plugin/gcp/bigquery/source/BigQuerySourceUtilsTest.java
@@ -17,15 +17,18 @@
package io.cdap.plugin.gcp.bigquery.source;
import com.google.cloud.bigquery.Dataset;
+import com.google.cloud.bigquery.DatasetId;
import com.google.cloud.storage.Bucket;
+import com.google.cloud.storage.BucketInfo;
import com.google.cloud.storage.Storage;
+import com.google.cloud.storage.StorageException;
+import io.cdap.cdap.api.exception.ProgramFailureException;
import org.apache.hadoop.conf.Configuration;
import org.junit.Assert;
import org.junit.Test;
import org.mockito.Mockito;
import java.io.IOException;
-import java.lang.reflect.Field;
public class BigQuerySourceUtilsTest {
@@ -50,4 +53,39 @@ public void getOrCreateBucket() throws IllegalAccessException, NoSuchFieldExcept
"some-path", null);
Assert.assertEquals("a-bucket", bucket2);
}
+
+ @Test
+ public void getOrCreateBucket_nullBucketConflict409_returnsBucketAndEnablesDelete() throws IOException {
+ Configuration configuration = new Configuration();
+ BigQuerySourceConfig config = BigQuerySourceConfig.builder().build();
+ Dataset ds = Mockito.mock(Dataset.class);
+ Storage st = Mockito.mock(Storage.class);
+ Mockito.when(st.create(Mockito.any(BucketInfo.class)))
+ .thenThrow(new StorageException(409, "Your previous request to create the named bucket succeeded "
+ + "and you already own it."));
+
+ String bucket = BigQuerySourceUtils.getOrCreateBucket(configuration, st, config.getBucket(), ds,
+ "some-path", null);
+
+ Assert.assertEquals("bq-source-bucket-some-path", bucket);
+ Assert.assertTrue(configuration.getBoolean("fs.gs.bucket.delete.enable", false));
+ }
+
+ @Test
+ public void getOrCreateBucket_nullBucketNon409Error_throwsProgramFailureException() {
+ Configuration configuration = new Configuration();
+ BigQuerySourceConfig config = BigQuerySourceConfig.builder().build();
+ Dataset ds = Mockito.mock(Dataset.class);
+ DatasetId datasetId = DatasetId.of("project", "dataset");
+ Mockito.when(ds.getDatasetId()).thenReturn(datasetId);
+ Mockito.when(ds.getLocation()).thenReturn("US");
+ Storage st = Mockito.mock(Storage.class);
+ Mockito.when(st.create(Mockito.any(BucketInfo.class)))
+ .thenThrow(new StorageException(403, "Access denied"));
+
+ ProgramFailureException exception = Assert.assertThrows(ProgramFailureException.class, () ->
+ BigQuerySourceUtils.getOrCreateBucket(configuration, st, config.getBucket(), ds, "some-path", null));
+
+ Assert.assertTrue(exception.getMessage().contains("Access denied"));
+ }
}