From 2272e92b3255b415c8662ddf7ce1b05d57b13138 Mon Sep 17 00:00:00 2001 From: giwa Date: Fri, 5 Nov 2021 16:01:44 +0900 Subject: [PATCH 1/4] add destination_project --- README.md | 1 + .../output/bigquery_java/BigqueryClient.java | 55 +++++++++++-------- .../config/BigqueryConfigValidator.java | 11 ++++ .../config/BigqueryTaskBuilder.java | 23 ++++++++ .../bigquery_java/config/PluginTask.java | 23 +++++++- .../config/TestBigqueryConfigValidator.java | 29 ++++++++++ .../config/TestBigqueryTaskBuilder.java | 18 ++++++ .../embulk/output/bigquery_java/json_key.json | 3 + 8 files changed, 138 insertions(+), 25 deletions(-) create mode 100644 src/test/resources/java/org/embulk/output/bigquery_java/json_key.json diff --git a/README.md b/README.md index 7e1a142..343754b 100644 --- a/README.md +++ b/README.md @@ -39,6 +39,7 @@ Under construction | auth_method (service_account is supported) | string | optional | "application\_default" | See [Authentication](#authentication) | | json_keyfile | string | optional | | keyfile path or `content` | | project (x) | string | required unless service\_account's `json_keyfile` is given. | | project\_id | +| destination_project | string | optional | `project` value | A destination project to which the data will be loaded. Use this if you want to separate a billing project (the `project` value) and a destination project (the `destination_project` value). | | dataset | string | required | | dataset | | location | string | optional | nil | geographic location of dataset. See [Location](#location) | | table | string | required | | table name, or table name with a partition decorator such as `table_name$20160929`| diff --git a/src/main/java/org/embulk/output/bigquery_java/BigqueryClient.java b/src/main/java/org/embulk/output/bigquery_java/BigqueryClient.java index 66ef05e..a630eb0 100644 --- a/src/main/java/org/embulk/output/bigquery_java/BigqueryClient.java +++ b/src/main/java/org/embulk/output/bigquery_java/BigqueryClient.java @@ -43,6 +43,7 @@ public class BigqueryClient { private final Logger logger = LoggerFactory.getLogger(BigqueryClient.class); private BigQuery bigquery; + private String destinationProject; private String dataset; private String location; private String locationForLog; @@ -53,16 +54,17 @@ public class BigqueryClient { public BigqueryClient(PluginTask task, Schema schema) { this.task = task; this.schema = schema; + this.destinationProject = task.getDestinationProject().get(); this.dataset = task.getDataset(); - if (task.getLocation().isPresent()){ + if (task.getLocation().isPresent()) { this.location = task.getLocation().get(); this.locationForLog = task.getLocation().get(); - }else{ + } else { this.locationForLog = "us/eu"; } this.columnOptions = this.task.getColumnOptions().orElse(Collections.emptyList()); try { - this.bigquery = getClientWithJsonKey(this.task.getJsonKeyfile()); + this.bigquery = getClientWithJsonKey(this.task.getJsonKeyfile().get()); } catch (IOException e) { throw new RuntimeException(e); } @@ -75,15 +77,17 @@ private static BigQuery getClientWithJsonKey(String key) throws IOException { .getService(); } - public Dataset createDataset(String datasetId){ - DatasetInfo.Builder builder = DatasetInfo.newBuilder(datasetId); - if (this.location != null){ + public Dataset createDataset(String datasetId) { + logger.info(String.format("embulk-output-bigquery_java: Create dataset... %s:%s in %s", this.destinationProject, datasetId, this.locationForLog)); + DatasetInfo.Builder builder = DatasetInfo.newBuilder(this.destinationProject, datasetId); + if (this.location != null) { builder.setLocation(this.location); } return bigquery.create(builder.build()); } public Dataset getDataset(String datasetId) { + logger.info(String.format("embulk-output-bigquery_java: Get dataset... %s:%s", this.destinationProject, datasetId)); return bigquery.getDataset(datasetId); } @@ -92,14 +96,14 @@ public Job getJob(JobId jobId) { } public Table getTable(String name) { - return getTable(TableId.of(this.dataset, name)); + return getTable(TableId.of(this.destinationProject, this.dataset, name)); } public Table getTable(TableId tableId) { return this.bigquery.getTable(tableId); } - public void createTableIfNotExist(String table) { + public void createTableIfNotExist(String table) { createTableIfNotExist(table, dataset); } @@ -112,14 +116,15 @@ public void createTableIfNotExist(String table, String dataset) { } TableDefinition tableDefinition = tableDefinitionBuilder.build(); - try{ - bigquery.create(TableInfo.newBuilder(TableId.of(dataset, table), tableDefinition).build()); - }catch (BigQueryException e){ - if(e.getCode() == 409 && e.getMessage().contains("Already Exists:")){ + logger.info(String.format("embulk-output-bigquery: Create table... %s:%s.%s", this.destinationProject, dataset, table)); + try { + bigquery.create(TableInfo.newBuilder(TableId.of(this.destinationProject, dataset, table), tableDefinition).build()); + } catch (BigQueryException e) { + if (e.getCode() == 409 && e.getMessage().contains("Already Exists:")) { return; } - logger.error(String.format("embulk-out_bigquery: insert_table(%s, %s)", dataset, table)); - throw new BigqueryException(String.format("failed to create table %s.%s, response: %s", dataset, table, e)); + logger.error(String.format("embulk-out_bigquery: insert_table(%s, %s, %s)", this.destinationProject, dataset, table)); + throw new BigqueryException(String.format("failed to create table %s:%s.%s, response: %s", this.destinationProject, dataset, table, e)); } } @@ -147,6 +152,7 @@ public JobStatistics.LoadStatistics load(Path loadFile, String table, JobInfo.Wr int retries = this.task.getRetries(); PluginTask task = this.task; Schema schema = this.schema; + String destinationProject = this.task.getDestinationProject().get(); List columnOptions = this.columnOptions; try { @@ -171,7 +177,7 @@ public JobStatistics.LoadStatistics call() { return null; } - TableId tableId = TableId.of(dataset, table); + TableId tableId = TableId.of(destinationProject, dataset, table); WriteChannelConfiguration writeChannelConfiguration = WriteChannelConfiguration.newBuilder(tableId) .setFormatOptions(FormatOptions.json()) @@ -234,6 +240,7 @@ public JobStatistics.CopyStatistics copy(String sourceTable, JobInfo.WriteDisposition writeDestination) throws BigqueryException { String dataset = this.dataset; int retries = this.task.getRetries(); + String destinationProject = this.task.getDestinationProject().get(); try { return retryExecutor() @@ -245,8 +252,8 @@ public JobStatistics.CopyStatistics copy(String sourceTable, public JobStatistics.CopyStatistics call() { UUID uuid = UUID.randomUUID(); String jobId = String.format("embulk_load_job_%s", uuid.toString()); - TableId destTableId = TableId.of(destinationDataset, destinationTable); - TableId srcTableId = TableId.of(dataset, sourceTable); + TableId destTableId = TableId.of(destinationProject, destinationDataset, destinationTable); + TableId srcTableId = TableId.of(destinationProject, dataset, sourceTable); CopyJobConfiguration copyJobConfiguration = CopyJobConfiguration.newBuilder(destTableId, srcTableId) .setWriteDisposition(writeDestination) @@ -312,7 +319,7 @@ public JobStatistics.QueryStatistics call() { .setUseLegacySql(false) .build(); JobId.Builder jobIdBuilder = JobId.newBuilder().setJob(jobId); - if (location != null){ + if (location != null) { jobIdBuilder.setLocation(location); } @@ -362,24 +369,24 @@ public boolean deleteTable(String table) { } public boolean deleteTable(String table, String dataset) { - if (dataset == null){ + if (dataset == null) { dataset = this.dataset; } String chompedTable = BigqueryUtil.chompPartitionDecorator(table); return deleteTableOrPartition(chompedTable, dataset); } - public boolean deleteTableOrPartition(String table){ + public boolean deleteTableOrPartition(String table) { return deleteTableOrPartition(table, null); } // if `table` with a partition decorator is given, a partition is deleted. - public boolean deleteTableOrPartition(String table, String dataset){ - if (dataset == null){ + public boolean deleteTableOrPartition(String table, String dataset) { + if (dataset == null) { dataset = this.dataset; } - logger.info(String.format("embulk-output-bigquery: Delete table... %s.%s", dataset, table)); - return this.bigquery.delete(TableId.of(dataset, table)); + logger.info(String.format("embulk-output-bigquery: Delete table... %s:%s.%s", this.destinationProject, dataset, table)); + return this.bigquery.delete(TableId.of(this.destinationProject, dataset, table)); } private JobStatistics waitForLoad(Job job) throws BigqueryException { diff --git a/src/main/java/org/embulk/output/bigquery_java/config/BigqueryConfigValidator.java b/src/main/java/org/embulk/output/bigquery_java/config/BigqueryConfigValidator.java index 6a3fe01..2d68c85 100644 --- a/src/main/java/org/embulk/output/bigquery_java/config/BigqueryConfigValidator.java +++ b/src/main/java/org/embulk/output/bigquery_java/config/BigqueryConfigValidator.java @@ -8,6 +8,7 @@ public class BigqueryConfigValidator { public static void validate(PluginTask task) { validateMode(task); validateModeAndAutoCreteTable(task); + validateProject(task); } public static void validateMode(PluginTask task) throws ConfigException { @@ -33,4 +34,14 @@ public static void validateTimePartitioning(PluginTask task) throws ConfigExcept } } } + + public static void validateProject(PluginTask task) throws ConfigException { + if (!task.getProject().isPresent()){ + throw new ConfigException("project is empty"); + } + String project = task.getProject().get(); + if (project.isEmpty()){ + throw new ConfigException("project is empty string"); + } + } } diff --git a/src/main/java/org/embulk/output/bigquery_java/config/BigqueryTaskBuilder.java b/src/main/java/org/embulk/output/bigquery_java/config/BigqueryTaskBuilder.java index 99d9d2c..3d27e95 100644 --- a/src/main/java/org/embulk/output/bigquery_java/config/BigqueryTaskBuilder.java +++ b/src/main/java/org/embulk/output/bigquery_java/config/BigqueryTaskBuilder.java @@ -6,7 +6,11 @@ import java.util.UUID; import java.util.Optional; +import com.fasterxml.jackson.annotation.JacksonAnnotation; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.annotations.VisibleForTesting; +import org.embulk.config.ConfigException; public class BigqueryTaskBuilder { private static final String uniqueName = UUID.randomUUID().toString().replace("-", "_"); @@ -16,6 +20,7 @@ public static PluginTask build(PluginTask task) { setFileExt(task); setTempTable(task); setAbortOnError(task); + setProject(task); return task; } @@ -63,4 +68,22 @@ protected static void setAbortOnError(PluginTask task) { task.setAbortOnError(Optional.of(task.getMaxBadRecords() == 0)); } } + + @VisibleForTesting + protected static void setProject(PluginTask task){ + if (task.getJsonKeyfile().isPresent()){ + JsonNode root; + try { + ObjectMapper mapper = new ObjectMapper(); + root = mapper.readTree(task.getJsonKeyfile().get()); + }catch (IOException e){ + throw new ConfigException(String.format("Parsing 'json_keyfile' failed with error: %s %s", e.getClass(), e.getMessage())); + } + task.setProject(Optional.of(root.get("project_id").toString())); + } + + if (!task.getDestinationProject().isPresent()){ + task.setDestinationProject(Optional.of(task.getProject().get())); + } + } } diff --git a/src/main/java/org/embulk/output/bigquery_java/config/PluginTask.java b/src/main/java/org/embulk/output/bigquery_java/config/PluginTask.java index 81517c4..7de9f2d 100644 --- a/src/main/java/org/embulk/output/bigquery_java/config/PluginTask.java +++ b/src/main/java/org/embulk/output/bigquery_java/config/PluginTask.java @@ -25,7 +25,20 @@ public interface PluginTask String getAuthMethod(); @Config("json_keyfile") - String getJsonKeyfile(); + @ConfigDefault("null") + Optional getJsonKeyfile(); + + @Config("project") + @ConfigDefault("null") + Optional getProject(); + + void setProject(Optional project); + + @Config("destination_project") + @ConfigDefault("null") + Optional getDestinationProject(); + + void setDestinationProject(Optional destinationProject); @Config("dataset") String getDataset(); @@ -154,4 +167,12 @@ public interface PluginTask Optional getTimePartitioning(); void setTimePartitioning(Optional bigqueryTimePartitioning); + + @Config("gcs_bucket") + @ConfigDefault("null") + Optional getGcsBucket(); + + @Config("auto_create_gcs_bucket") + @ConfigDefault("false") + Optional getAutoCreateGcsBucket(); } diff --git a/src/test/java/org/embulk/output/bigquery_java/config/TestBigqueryConfigValidator.java b/src/test/java/org/embulk/output/bigquery_java/config/TestBigqueryConfigValidator.java index 8cd4e02..455c27a 100644 --- a/src/test/java/org/embulk/output/bigquery_java/config/TestBigqueryConfigValidator.java +++ b/src/test/java/org/embulk/output/bigquery_java/config/TestBigqueryConfigValidator.java @@ -1,5 +1,6 @@ package org.embulk.output.bigquery_java.config; +import org.embulk.config.Config; import org.embulk.config.ConfigException; import org.embulk.config.ConfigSource; import org.embulk.output.bigquery_java.BigqueryJavaOutputPlugin; @@ -8,6 +9,8 @@ import org.junit.Rule; import org.junit.Test; +import java.util.Optional; + import static org.junit.Assert.*; public class TestBigqueryConfigValidator { @@ -57,4 +60,30 @@ public void validateModeAndAutoCreteTable_autoCreateTable_False_configException( task.setAutoCreateTable(false); BigqueryConfigValidator.validateModeAndAutoCreteTable(task); } + + @Test + public void validateProject() { + config = loadYamlResource(embulk, "base.yml"); + PluginTask task = config.loadConfig(PluginTask.class); + task.setProject(Optional.of("project_id")); + BigqueryConfigValidator.validateProject(task); + + assertEquals("project_id", task.getProject().get()); + } + + @Test(expected = ConfigException.class) + public void validateProject_project_null_configException() { + config = loadYamlResource(embulk, "base.yml"); + PluginTask task = config.loadConfig(PluginTask.class); + BigqueryConfigValidator.validateProject(task); + } + + @Test(expected = ConfigException.class) + public void validateProject_project_empty_string_configException() { + config = loadYamlResource(embulk, "base.yml"); + PluginTask task = config.loadConfig(PluginTask.class); + task.setProject(Optional.empty()); + BigqueryConfigValidator.validateProject(task); + } + } diff --git a/src/test/java/org/embulk/output/bigquery_java/config/TestBigqueryTaskBuilder.java b/src/test/java/org/embulk/output/bigquery_java/config/TestBigqueryTaskBuilder.java index ac73481..532dee7 100644 --- a/src/test/java/org/embulk/output/bigquery_java/config/TestBigqueryTaskBuilder.java +++ b/src/test/java/org/embulk/output/bigquery_java/config/TestBigqueryTaskBuilder.java @@ -1,5 +1,6 @@ package org.embulk.output.bigquery_java.config; +import org.embulk.config.ConfigException; import org.embulk.config.ConfigSource; import org.embulk.output.bigquery_java.BigqueryJavaOutputPlugin; import org.embulk.spi.OutputPlugin; @@ -44,4 +45,21 @@ public void setFileExt_JSONL_GZIP_JSONL_GZ() { assertEquals("GZIP", task.getCompression()); assertEquals(".jsonl.gz", task.getFileExt().get()); } + + @Test + public void setProject() { + config = loadYamlResource(embulk, "base.yml"); + config.set("json_keyfile", BASIC_RESOURCE_PATH + "json_key.json"); + PluginTask task = config.loadConfig(PluginTask.class); + BigqueryTaskBuilder.setProject(task); + + assertEquals("project_id", task.getProject().get()); + } + + @Test(expected = ConfigException.class) + public void setProject_config_exception() { + config = loadYamlResource(embulk, "base.yml"); + PluginTask task = config.loadConfig(PluginTask.class); + BigqueryTaskBuilder.setProject(task); + } } diff --git a/src/test/resources/java/org/embulk/output/bigquery_java/json_key.json b/src/test/resources/java/org/embulk/output/bigquery_java/json_key.json new file mode 100644 index 0000000..b054a80 --- /dev/null +++ b/src/test/resources/java/org/embulk/output/bigquery_java/json_key.json @@ -0,0 +1,3 @@ +{ + "project_id": "project_id" +} \ No newline at end of file From 9e922ec15f345b9e5da6f967812862623f944ce1 Mon Sep 17 00:00:00 2001 From: giwa Date: Fri, 5 Nov 2021 16:22:22 +0900 Subject: [PATCH 2/4] pass file object not path --- .../output/bigquery_java/config/BigqueryTaskBuilder.java | 2 +- .../output/bigquery_java/config/TestBigqueryTaskBuilder.java | 3 ++- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/src/main/java/org/embulk/output/bigquery_java/config/BigqueryTaskBuilder.java b/src/main/java/org/embulk/output/bigquery_java/config/BigqueryTaskBuilder.java index 3d27e95..709b770 100644 --- a/src/main/java/org/embulk/output/bigquery_java/config/BigqueryTaskBuilder.java +++ b/src/main/java/org/embulk/output/bigquery_java/config/BigqueryTaskBuilder.java @@ -75,7 +75,7 @@ protected static void setProject(PluginTask task){ JsonNode root; try { ObjectMapper mapper = new ObjectMapper(); - root = mapper.readTree(task.getJsonKeyfile().get()); + root = mapper.readTree(new File(task.getJsonKeyfile().get())); }catch (IOException e){ throw new ConfigException(String.format("Parsing 'json_keyfile' failed with error: %s %s", e.getClass(), e.getMessage())); } diff --git a/src/test/java/org/embulk/output/bigquery_java/config/TestBigqueryTaskBuilder.java b/src/test/java/org/embulk/output/bigquery_java/config/TestBigqueryTaskBuilder.java index 532dee7..553dbbc 100644 --- a/src/test/java/org/embulk/output/bigquery_java/config/TestBigqueryTaskBuilder.java +++ b/src/test/java/org/embulk/output/bigquery_java/config/TestBigqueryTaskBuilder.java @@ -1,5 +1,6 @@ package org.embulk.output.bigquery_java.config; +import com.google.common.io.Resources; import org.embulk.config.ConfigException; import org.embulk.config.ConfigSource; import org.embulk.output.bigquery_java.BigqueryJavaOutputPlugin; @@ -49,7 +50,7 @@ public void setFileExt_JSONL_GZIP_JSONL_GZ() { @Test public void setProject() { config = loadYamlResource(embulk, "base.yml"); - config.set("json_keyfile", BASIC_RESOURCE_PATH + "json_key.json"); + config.set("json_keyfile", Resources.getResource(BASIC_RESOURCE_PATH+"json_key.json").getPath()); PluginTask task = config.loadConfig(PluginTask.class); BigqueryTaskBuilder.setProject(task); From 4133f37797ee1122581c532003dd327bfd68eb07 Mon Sep 17 00:00:00 2001 From: giwa Date: Fri, 5 Nov 2021 16:48:30 +0900 Subject: [PATCH 3/4] pass value as text --- .../embulk/output/bigquery_java/config/BigqueryTaskBuilder.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/main/java/org/embulk/output/bigquery_java/config/BigqueryTaskBuilder.java b/src/main/java/org/embulk/output/bigquery_java/config/BigqueryTaskBuilder.java index 709b770..505a22a 100644 --- a/src/main/java/org/embulk/output/bigquery_java/config/BigqueryTaskBuilder.java +++ b/src/main/java/org/embulk/output/bigquery_java/config/BigqueryTaskBuilder.java @@ -79,7 +79,7 @@ protected static void setProject(PluginTask task){ }catch (IOException e){ throw new ConfigException(String.format("Parsing 'json_keyfile' failed with error: %s %s", e.getClass(), e.getMessage())); } - task.setProject(Optional.of(root.get("project_id").toString())); + task.setProject(Optional.of(root.get("project_id").asText())); } if (!task.getDestinationProject().isPresent()){ From 1d165f1e7743d6eee1e118a7e1f672a142cce670 Mon Sep 17 00:00:00 2001 From: giwa Date: Wed, 10 Nov 2021 18:33:50 +0900 Subject: [PATCH 4/4] validate time partitioning --- .../output/bigquery_java/config/BigqueryConfigValidator.java | 1 + 1 file changed, 1 insertion(+) diff --git a/src/main/java/org/embulk/output/bigquery_java/config/BigqueryConfigValidator.java b/src/main/java/org/embulk/output/bigquery_java/config/BigqueryConfigValidator.java index 2d68c85..24b3a8c 100644 --- a/src/main/java/org/embulk/output/bigquery_java/config/BigqueryConfigValidator.java +++ b/src/main/java/org/embulk/output/bigquery_java/config/BigqueryConfigValidator.java @@ -8,6 +8,7 @@ public class BigqueryConfigValidator { public static void validate(PluginTask task) { validateMode(task); validateModeAndAutoCreteTable(task); + validateTimePartitioning(task); validateProject(task); }