diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryDatabaseMetaData.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryDatabaseMetaData.java index 75414771600c..4d5b2978d02f 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryDatabaseMetaData.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryDatabaseMetaData.java @@ -52,15 +52,11 @@ import com.google.cloud.bigquery.exception.BigQueryJdbcException; import com.google.cloud.bigquery.jdbc.BigQueryJdbcTypeMappings.ColumnTypeInfo; import com.google.cloud.bigquery.jdbc.utils.BigQueryJdbcVersionUtility; -import java.io.BufferedReader; -import java.io.InputStream; -import java.io.InputStreamReader; import java.sql.Connection; import java.sql.DatabaseMetaData; import java.sql.ResultSet; import java.sql.RowIdLifetime; import java.sql.SQLException; -import java.sql.Statement; import java.sql.Types; import java.util.ArrayList; import java.util.Arrays; @@ -68,7 +64,6 @@ import java.util.Comparator; import java.util.HashSet; import java.util.List; -import java.util.Scanner; import java.util.Set; import java.util.concurrent.BlockingQueue; import java.util.concurrent.Callable; @@ -89,17 +84,14 @@ * * @see BigQueryStatement */ -// TODO(neenu): test and verify after post MVP implementation. class BigQueryDatabaseMetaData implements DatabaseMetaData { final BigQueryJdbcCustomLogger LOG = new BigQueryJdbcCustomLogger(this.toString()); private static final String DATABASE_PRODUCT_NAME = "Google BigQuery"; private static final String DATABASE_PRODUCT_VERSION = "2.0"; private static final String DRIVER_NAME = "GoogleJDBCDriverForGoogleBigQuery"; - private static final String SCHEMA_TERM = "Dataset"; private static final String CATALOG_TERM = "Project"; private static final String PROCEDURE_TERM = "Procedure"; - private static final String GET_EXPORTED_KEYS_SQL = "DatabaseMetaData_GetExportedKeys.sql"; private static final int DEFAULT_PAGE_SIZE = 500; private static final int DEFAULT_QUEUE_CAPACITY = 5000; // Declared package-private for testing. @@ -1820,11 +1812,9 @@ public ResultSet getCatalogs() throws SQLException { final BlockingQueue queue = new LinkedBlockingQueue<>(catalogRows.isEmpty() ? 1 : catalogRows.size() + 1); - populateQueue(catalogRows, queue, schemaFields); - signalEndOfData(queue, schemaFields); + Future fetcherFuture = populateQueueAsync(catalogRows, queue, schemaFields); - return BigQueryJsonResultSet.of( - catalogsSchema, catalogRows.size(), queue, null, new Future[0]); + return BigQueryJsonResultSet.of(catalogsSchema, catalogRows.size(), queue, null, fetcherFuture); } Schema defineGetCatalogsSchema() { @@ -1852,11 +1842,11 @@ public ResultSet getTableTypes() { BlockingQueue queue = new LinkedBlockingQueue<>(tableTypeRows.size() + 1); - populateQueue(tableTypeRows, queue, tableTypesSchema.getFields()); - signalEndOfData(queue, tableTypesSchema.getFields()); + Future fetcherFuture = + populateQueueAsync(tableTypeRows, queue, tableTypesSchema.getFields()); return BigQueryJsonResultSet.of( - tableTypesSchema, tableTypeRows.size(), queue, null, new Future[0]); + tableTypesSchema, tableTypeRows.size(), queue, null, fetcherFuture); } static Schema defineGetTableTypesSchema() { @@ -2413,17 +2403,6 @@ Schema defineGetVersionColumnsSchema() { return Schema.of(fields); } - private void closeStatementIgnoreException(Statement statement) { - if (statement == null) { - return; - } - try { - statement.close(); - } catch (SQLException e) { - // pass - } - } - @Override public ResultSet getPrimaryKeys(String catalog, String schema, String table) throws SQLException { if ((catalog != null && catalog.isEmpty()) @@ -2458,9 +2437,8 @@ public ResultSet getPrimaryKeys(String catalog, String schema, String table) thr final BlockingQueue queue = new LinkedBlockingQueue<>(DEFAULT_QUEUE_CAPACITY); - populateQueue(collectedResults, queue, resultSchemaFields); - signalEndOfData(queue, resultSchemaFields); - return BigQueryJsonResultSet.of(resultSchema, -1, queue, null); + Future fetcherFuture = populateQueueAsync(collectedResults, queue, resultSchemaFields); + return BigQueryJsonResultSet.of(resultSchema, -1, queue, null, fetcherFuture); } private Schema defineGetPrimaryKeysSchema() { @@ -2561,24 +2539,59 @@ public ResultSet getImportedKeys(String catalog, String schema, String table) final BlockingQueue queue = new LinkedBlockingQueue<>(DEFAULT_QUEUE_CAPACITY); - populateQueue(collectedResults, queue, resultSchemaFields); - signalEndOfData(queue, resultSchemaFields); - return BigQueryJsonResultSet.of(resultSchema, -1, queue, null); + Future fetcherFuture = populateQueueAsync(collectedResults, queue, resultSchemaFields); + return BigQueryJsonResultSet.of(resultSchema, -1, queue, null, fetcherFuture); } @Override public ResultSet getExportedKeys(String catalog, String schema, String table) throws SQLException { - String sql = readSqlFromFile(GET_EXPORTED_KEYS_SQL); - Statement stmt = this.connection.createStatement(); - try { - stmt.closeOnCompletion(); - String formattedSql = replaceSqlParameters(sql, catalog, schema, table); - return stmt.executeQuery(formattedSql); - } catch (SQLException e) { - closeStatementIgnoreException(stmt); - throw new BigQueryJdbcException("Error executing getExportedKeys", e); + if ((catalog != null && catalog.isEmpty()) + || (schema != null && schema.isEmpty()) + || table == null + || table.isEmpty()) { + LOG.warning( + "Returning empty ResultSet as required parameters are null/empty, or catalog/schema parameters are empty."); + return new BigQueryJsonResultSet(); } + + final Schema resultSchema = defineForeignKeyResultSetSchema(); + final FieldList resultSchemaFields = resultSchema.getFields(); + + final List collectedResults = Collections.synchronizedList(new ArrayList<>()); + List targetDatasets = getTargetDatasets(catalog, null); + + boolean ignoreAccessErrors = (catalog == null); + processTargetTablesConcurrently( + targetDatasets, + null, + collectedResults, + resultSchemaFields, + ignoreAccessErrors, + (bqTable, results, fields) -> { + TableConstraints constraints = bqTable.getTableConstraints(); + if (constraints == null || constraints.getForeignKeys() == null) { + return; + } + for (ForeignKey fk : constraints.getForeignKeys()) { + TableId pkTableId = fk.getReferencedTable(); + if (pkTableId == null + || !equalsOrNullMatchesAll(catalog, pkTableId.getProject()) + || !equalsOrNullMatchesAll(schema, pkTableId.getDataset()) + || !table.equals(pkTableId.getTable())) { + continue; + } + processForeignKey(fk, pkTableId, bqTable.getTableId(), results, fields); + } + }); + + Comparator comparator = defineFkTableSortComparator(resultSchemaFields); + sortResults(collectedResults, comparator, "getExportedKeys", LOG); + + final BlockingQueue queue = + new LinkedBlockingQueue<>(DEFAULT_QUEUE_CAPACITY); + Future fetcherFuture = populateQueueAsync(collectedResults, queue, resultSchemaFields); + return BigQueryJsonResultSet.of(resultSchema, -1, queue, null, fetcherFuture); } @Override @@ -2635,9 +2648,8 @@ && equalsOrNullMatchesAll(parentSchema, pkTableId.getDataset()) final BlockingQueue queue = new LinkedBlockingQueue<>(DEFAULT_QUEUE_CAPACITY); - populateQueue(collectedResults, queue, resultSchemaFields); - signalEndOfData(queue, resultSchemaFields); - return BigQueryJsonResultSet.of(resultSchema, -1, queue, null); + Future fetcherFuture = populateQueueAsync(collectedResults, queue, resultSchemaFields); + return BigQueryJsonResultSet.of(resultSchema, -1, queue, null, fetcherFuture); } @Override @@ -2653,10 +2665,9 @@ public ResultSet getTypeInfo() { final BlockingQueue queue = new LinkedBlockingQueue<>(typeInfoRows.size() + 1); - populateQueue(typeInfoRows, queue, schemaFields); - signalEndOfData(queue, schemaFields); + Future fetcherFuture = populateQueueAsync(typeInfoRows, queue, schemaFields); return BigQueryJsonResultSet.of( - typeInfoSchema, typeInfoRows.size(), queue, null, new Future[0]); + typeInfoSchema, typeInfoRows.size(), queue, null, fetcherFuture); } Schema defineGetTypeInfoSchema() { @@ -3577,9 +3588,8 @@ public ResultSet getSchemas(String catalog, String schemaPattern) throws SQLExce } Comparator comparator = defineGetSchemasComparator(resultSchemaFields); sortResults(collectedResults, comparator, "getSchemas", LOG); - populateQueue(collectedResults, queue, resultSchemaFields); - signalEndOfData(queue, resultSchemaFields); - return BigQueryJsonResultSet.of(resultSchema, -1, queue, null); + Future fetcherFuture = populateQueueAsync(collectedResults, queue, resultSchemaFields); + return BigQueryJsonResultSet.of(resultSchema, -1, queue, null, fetcherFuture); } // Multi-Catalog Path: fan out using connection-scoped metadataExecutor @@ -4891,6 +4901,19 @@ private void waitForTasksCompletion(List> taskFutures) throws Executio LOG.info("Finished waiting for tasks."); } + private Future populateQueueAsync( + List collectedResults, + BlockingQueue queue, + FieldList resultSchemaFields) { + return connection + .getMetadataExecutor() + .submit( + () -> { + populateQueue(collectedResults, queue, resultSchemaFields); + signalEndOfData(queue, resultSchemaFields); + }); + } + private void populateQueue( List collectedResults, BlockingQueue queue, @@ -5063,24 +5086,6 @@ private boolean equalsOrNullMatchesAll(String expected, String actual) { return expected == null || expected.equals(actual); } - static String readSqlFromFile(String filename) { - InputStream in; - in = BigQueryDatabaseMetaData.class.getResourceAsStream(filename); - BufferedReader reader = new BufferedReader(new InputStreamReader(in)); - StringBuilder builder = new StringBuilder(); - try (Scanner scanner = new Scanner(reader)) { - while (scanner.hasNextLine()) { - String line = scanner.nextLine(); - builder.append(line).append("\n"); - } - } - return builder.toString(); - } - - String replaceSqlParameters(String sql, String... params) throws SQLException { - return String.format(sql, (Object[]) params); - } - private void writeErrorToQueue(BlockingQueue queue, Throwable t) { Exception ex = (t instanceof Exception) ? (Exception) t : new Exception(t); BigQueryFieldValueListWrapper element = BigQueryFieldValueListWrapper.ofError(ex); @@ -5168,7 +5173,7 @@ private void processTargetTablesConcurrently( boolean ignoreAccessErrors, TableProcessor processor) throws SQLException { - if (targetDatasets.size() == 1) { + if (targetDatasets.size() == 1 && tableName != null) { processSingleTable( targetDatasets.get(0), tableName, @@ -5183,19 +5188,59 @@ private void processTargetTablesConcurrently( List> taskFutures = new ArrayList<>(); try { + List> tasks = new ArrayList<>(); for (DatasetId datasetId : targetDatasets) { - taskFutures.add( - executor.submit( + if (tableName != null) { + tasks.add( + () -> { + processSingleTable( + datasetId, + tableName, + collectedResults, + resultSchemaFields, + ignoreAccessErrors, + processor); + return null; + }); + continue; + } + + try { + Page tablesPage = + bigquery.listTables(datasetId, TableListOption.pageSize(DEFAULT_PAGE_SIZE)); + if (tablesPage == null) { + continue; + } + for (Table table : tablesPage.iterateAll()) { + if (table.getDefinition() == null + || table.getDefinition().getType() != TableDefinition.Type.TABLE) { + continue; + } + tasks.add( () -> { processSingleTable( datasetId, - tableName, + table.getTableId().getTable(), collectedResults, resultSchemaFields, ignoreAccessErrors, processor); return null; - })); + }); + } + } catch (BigQueryException e) { + if (ignoreAccessErrors && (e.getCode() == 404 || e.getCode() == 403)) { + LOG.info( + "Dataset '%s' not found/accessible in project '%s' (API error %d). Skipping.", + datasetId.getDataset(), datasetId.getProject(), e.getCode()); + continue; + } + throw new SQLException("Error while listing tables: " + e.getMessage(), e); + } + } + + for (Callable task : tasks) { + taskFutures.add(executor.submit(task)); } waitForTasksCompletion(taskFutures); if (Thread.currentThread().isInterrupted()) { diff --git a/java-bigquery-jdbc/src/main/resources/com/google/cloud/bigquery/jdbc/DatabaseMetaData_GetExportedKeys.sql b/java-bigquery-jdbc/src/main/resources/com/google/cloud/bigquery/jdbc/DatabaseMetaData_GetExportedKeys.sql deleted file mode 100644 index 4058f6bff60a..000000000000 --- a/java-bigquery-jdbc/src/main/resources/com/google/cloud/bigquery/jdbc/DatabaseMetaData_GetExportedKeys.sql +++ /dev/null @@ -1,71 +0,0 @@ -/* - * Copyright 2024 Google LLC - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -SELECT PKTABLE_CAT, - PKTABLE_SCHEM, - PKTABLE_NAME, - PRIMARY.column_name AS PKCOLUMN_NAME, - FOREIGN.constraint_catalog AS FKTABLE_CAT, - FOREIGN.constraint_schema AS FKTABLE_SCHEM, - FOREIGN.table_name AS FKTABLE_NAME, - FOREIGN.column_name AS FKCOLUMN_NAME, - FOREIGN.ordinal_position AS KEY_SEQ, - NULL AS UPDATE_RULE, - NULL AS DELETE_RULE, - FOREIGN.constraint_name AS FK_NAME, - PRIMARY.constraint_name AS PK_NAME, - NULL AS DEFERRABILITY -FROM (SELECT DISTINCT CCU.table_catalog AS PKTABLE_CAT, - CCU.table_schema AS PKTABLE_SCHEM, - CCU.table_name AS PKTABLE_NAME, - TC.constraint_catalog, - TC.constraint_schema, - TC.constraint_name, - TC.table_catalog, - TC.table_schema, - TC.table_name, - TC.constraint_type, - KCU.column_name, - KCU.ordinal_position, - KCU.position_in_unique_constraint - FROM `%1$s.%2$s.INFORMATION_SCHEMA.TABLE_CONSTRAINTS` TC - INNER JOIN - `%1$s.%2$s.INFORMATION_SCHEMA.KEY_COLUMN_USAGE` KCU - USING - (constraint_catalog, - constraint_schema, - constraint_name, - table_catalog, - table_schema, - table_name) - INNER JOIN - `%1$s.%2$s.INFORMATION_SCHEMA.CONSTRAINT_COLUMN_USAGE` CCU - USING - (constraint_catalog, - constraint_schema, - constraint_name) - WHERE constraint_type = 'FOREIGN KEY') FOREIGN - INNER JOIN (SELECT * - FROM `%1$s.%2$s.INFORMATION_SCHEMA.KEY_COLUMN_USAGE` - WHERE position_in_unique_constraint IS NULL - AND RTRIM(table_name) = '%3$s') PRIMARY -ON - FOREIGN.PKTABLE_CAT = PRIMARY.table_catalog - AND FOREIGN.PKTABLE_SCHEM = PRIMARY.table_schema - AND FOREIGN.PKTABLE_NAME = PRIMARY.table_name - AND FOREIGN.position_in_unique_constraint = - PRIMARY.ordinal_position -ORDER BY FKTABLE_CAT, FKTABLE_SCHEM, FKTABLE_NAME, KEY_SEQ \ No newline at end of file diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITDatabaseMetadataTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITDatabaseMetadataTest.java index e40ba8bccd8c..046f49991dbe 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITDatabaseMetadataTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITDatabaseMetadataTest.java @@ -329,6 +329,64 @@ public void testTableConstraints() throws SQLException { connection.close(); } + @Test + public void testGetExportedKeys_multipleKeys() throws SQLException { + try (Connection connection = DriverManager.getConnection(ITBase.connectionUrl)) { + DatabaseMetaData metaData = connection.getMetaData(); + + // Table 2 exports keys to Table + ResultSet exportedKeys2 = + metaData.getExportedKeys(PROJECT_ID, CONSTRAINTS_DATASET, CONSTRAINTS_TABLE_NAME2); + Assertions.assertTrue(exportedKeys2.next()); + Assertions.assertEquals(CONSTRAINTS_TABLE_NAME2, exportedKeys2.getString(3)); // PKTABLE_NAME + Assertions.assertEquals("first_name", exportedKeys2.getString(4)); // PKCOLUMN_NAME + Assertions.assertEquals(CONSTRAINTS_TABLE_NAME, exportedKeys2.getString(7)); // FKTABLE_NAME + Assertions.assertEquals("name", exportedKeys2.getString(8)); // FKCOLUMN_NAME + Assertions.assertEquals(1, exportedKeys2.getInt(9)); // KEY_SEQ + Assertions.assertEquals("my_fk", exportedKeys2.getString(12)); // FK_NAME + + Assertions.assertTrue(exportedKeys2.next()); + Assertions.assertEquals(CONSTRAINTS_TABLE_NAME2, exportedKeys2.getString(3)); + Assertions.assertEquals("last_name", exportedKeys2.getString(4)); + Assertions.assertEquals(CONSTRAINTS_TABLE_NAME, exportedKeys2.getString(7)); + Assertions.assertEquals("second_name", exportedKeys2.getString(8)); + Assertions.assertEquals(2, exportedKeys2.getInt(9)); + Assertions.assertEquals("my_fk", exportedKeys2.getString(12)); + Assertions.assertFalse(exportedKeys2.next()); + } + } + + @Test + public void testGetExportedKeys_singleKey() throws SQLException { + try (Connection connection = DriverManager.getConnection(ITBase.connectionUrl)) { + DatabaseMetaData metaData = connection.getMetaData(); + + // Table 3 exports keys to Table + ResultSet exportedKeys3 = + metaData.getExportedKeys(PROJECT_ID, CONSTRAINTS_DATASET, CONSTRAINTS_TABLE_NAME3); + Assertions.assertTrue(exportedKeys3.next()); + Assertions.assertEquals(CONSTRAINTS_TABLE_NAME3, exportedKeys3.getString(3)); + Assertions.assertEquals("address", exportedKeys3.getString(4)); + Assertions.assertEquals(CONSTRAINTS_TABLE_NAME, exportedKeys3.getString(7)); + Assertions.assertEquals("address", exportedKeys3.getString(8)); + Assertions.assertEquals(1, exportedKeys3.getInt(9)); + Assertions.assertEquals("my_fk2", exportedKeys3.getString(12)); + Assertions.assertFalse(exportedKeys3.next()); + } + } + + @Test + public void testGetExportedKeys_noKeys() throws SQLException { + try (Connection connection = DriverManager.getConnection(ITBase.connectionUrl)) { + DatabaseMetaData metaData = connection.getMetaData(); + + // Table does not export keys to anything (it only imports them) + ResultSet exportedKeys1 = + metaData.getExportedKeys(PROJECT_ID, CONSTRAINTS_DATASET, CONSTRAINTS_TABLE_NAME); + Assertions.assertFalse(exportedKeys1.next()); + } + } + @Test public void testMetadataResultSetsDoNotInterfere() throws SQLException { try (Connection connection = DriverManager.getConnection(ITBase.connectionUrl)) {