Skip to content
Merged
Show file tree
Hide file tree
Changes from 4 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
Original file line number Diff line number Diff line change
Expand Up @@ -52,11 +52,17 @@
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.ByteArrayOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.nio.charset.StandardCharsets;
import java.sql.Connection;
import java.sql.DatabaseMetaData;
import java.sql.PreparedStatement;
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;
Expand Down Expand Up @@ -94,6 +100,8 @@ class BigQueryDatabaseMetaData implements DatabaseMetaData {
private static final String PROCEDURE_TERM = "Procedure";
private static final int DEFAULT_PAGE_SIZE = 500;
private static final int DEFAULT_QUEUE_CAPACITY = 5000;
private static final String GET_EXPORTED_KEYS_SQL = "DatabaseMetaData_GetExportedKeys.sql";
Comment thread
keshavdandeva marked this conversation as resolved.
private static String exportedKeysSqlContent;
// Declared package-private for testing.
static final String GOOGLE_SQL_QUOTED_IDENTIFIER = "`";
// Does not include SQL:2003 Keywords as per JDBC spec.
Expand Down Expand Up @@ -2559,40 +2567,63 @@ public ResultSet getExportedKeys(String catalog, String schema, String table)
final Schema resultSchema = defineForeignKeyResultSetSchema();
final FieldList resultSchemaFields = resultSchema.getFields();

final List<FieldValueList> collectedResults = Collections.synchronizedList(new ArrayList<>());
List<DatasetId> targetDatasets = getTargetDatasets(catalog, null);
// Early return for PCNT catalog schemas (containing '.') as they do not support table
// constraints.
if (schema != null && schema.contains(".")) {
final BlockingQueue<BigQueryFieldValueListWrapper> queue = new LinkedBlockingQueue<>(1);
signalEndOfData(queue, resultSchemaFields);
return BigQueryJsonResultSet.of(resultSchema, 0, queue, 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;
// Fallback Path: If catalog or schema is null, fall back to REST API metadata scan.
if (catalog == null || schema == null) {
final List<FieldValueList> collectedResults = Collections.synchronizedList(new ArrayList<>());
List<DatasetId> targetDatasets = getTargetDatasets(catalog, schema);

boolean ignoreAccessErrors = (catalog == null);
processTargetTablesConcurrently(
targetDatasets,
null,
collectedResults,
resultSchemaFields,
ignoreAccessErrors,
(bqTable, results, fields) -> {
TableConstraints constraints = bqTable.getTableConstraints();
if (constraints == null || constraints.getForeignKeys() == null) {
return;
}
processForeignKey(fk, pkTableId, bqTable.getTableId(), results, fields);
}
});
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<FieldValueList> comparator = defineFkTableSortComparator(resultSchemaFields);
sortResults(collectedResults, comparator, "getExportedKeys", LOG);
Comparator<FieldValueList> comparator = defineFkTableSortComparator(resultSchemaFields);
sortResults(collectedResults, comparator, "getExportedKeys", LOG);

final BlockingQueue<BigQueryFieldValueListWrapper> queue =
new LinkedBlockingQueue<>(DEFAULT_QUEUE_CAPACITY);
Future<?> fetcherFuture = populateQueueAsync(collectedResults, queue, resultSchemaFields);
return BigQueryJsonResultSet.of(resultSchema, -1, queue, null, fetcherFuture);
final BlockingQueue<BigQueryFieldValueListWrapper> queue =
new LinkedBlockingQueue<>(DEFAULT_QUEUE_CAPACITY);
Future<?> fetcherFuture = populateQueueAsync(collectedResults, queue, resultSchemaFields);
return BigQueryJsonResultSet.of(resultSchema, -1, queue, null, fetcherFuture);
}

String sql = getExportedKeysSqlContent();
String formattedSql = replaceSqlParameters(sql, catalog, schema);
PreparedStatement stmt = this.connection.prepareStatement(formattedSql);
try {
stmt.closeOnCompletion();
stmt.setString(1, table);
return stmt.executeQuery();
} catch (SQLException e) {
closeStatementIgnoreException(stmt);
throw new BigQueryJdbcException("Error executing getExportedKeys", e);
}
}

@Override
Expand Down Expand Up @@ -5353,4 +5384,43 @@ private Comparator<FieldValueList> defineFkTableSortComparator(FieldList resultS
(FieldValueList fvl) -> getLongValueOrNull(fvl, KEY_SEQ_IDX),
Comparator.nullsFirst(Long::compareTo));
}

private static synchronized String getExportedKeysSqlContent() {
if (exportedKeysSqlContent == null) {
exportedKeysSqlContent = readSqlFromFile(GET_EXPORTED_KEYS_SQL);
}
return exportedKeysSqlContent;
}
Comment thread
keshavdandeva marked this conversation as resolved.

static String readSqlFromFile(String filename) {
try (InputStream in = BigQueryDatabaseMetaData.class.getResourceAsStream(filename)) {
if (in == null) {
throw new IllegalArgumentException("SQL file not found: " + filename);
}
ByteArrayOutputStream result = new ByteArrayOutputStream();
byte[] buffer = new byte[1024];
int length;
while ((length = in.read(buffer)) != -1) {
result.write(buffer, 0, length);
}
return result.toString(StandardCharsets.UTF_8.name());
} catch (IOException e) {
throw new RuntimeException("Failed to read SQL file: " + filename, e);
}
}

String replaceSqlParameters(String sql, String... params) throws SQLException {
return String.format(sql, (Object[]) params);
}
Comment thread
keshavdandeva marked this conversation as resolved.

private void closeStatementIgnoreException(Statement stmt) {
if (stmt == null) {
return;
}
try {
stmt.close();
} catch (SQLException e) {
// ignore
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
/*
* Copyright 2026 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 table_name = ?) 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
Comment thread
keshavdandeva marked this conversation as resolved.
Outdated
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,8 @@ public class ITDatabaseMetadataTest extends ITBase {
private static final String CONSTRAINTS_TABLE_NAME = "JDBC_CONSTRAINTS_TEST_TABLE";
private static final String CONSTRAINTS_TABLE_NAME2 = "JDBC_CONSTRAINTS_TEST_TABLE2";
private static final String CONSTRAINTS_TABLE_NAME3 = "JDBC_CONSTRAINTS_TEST_TABLE3";
private static final String PCNT_SCHEMA = "bq-drivers-test-warehouse.jdbc_pcnt_test_namespace";
private static final String PCNT_TABLE_NAME = "PCNT_TEST_TABLE";
private static final Pattern VERSION_PATTERN =
Pattern.compile("^(\\d+)\\.(\\d+)(?:\\.\\d+)+\\s*.*");
private static final String DEFAULT_CATALOG = ServiceOptions.getDefaultProjectId();
Expand Down Expand Up @@ -387,6 +389,62 @@ public void testGetExportedKeys_noKeys() throws SQLException {
}
}

@Test
public void testGetPrimaryKeys_pcntTable() throws SQLException {
try (Connection connection = DriverManager.getConnection(ITBase.connectionUrl);
ResultSet primaryKeys =
connection.getMetaData().getPrimaryKeys(PROJECT_ID, PCNT_SCHEMA, PCNT_TABLE_NAME)) {
Assertions.assertNotNull(primaryKeys);
ResultSetMetaData pkMetaData = primaryKeys.getMetaData();
Assertions.assertEquals(6, pkMetaData.getColumnCount());
Assertions.assertFalse(primaryKeys.next());
}
}

@Test
public void testGetImportedKeys_pcntTable() throws SQLException {
try (Connection connection = DriverManager.getConnection(ITBase.connectionUrl);
ResultSet importedKeys =
connection.getMetaData().getImportedKeys(PROJECT_ID, PCNT_SCHEMA, PCNT_TABLE_NAME)) {
Assertions.assertNotNull(importedKeys);
ResultSetMetaData ikMetaData = importedKeys.getMetaData();
Assertions.assertEquals(14, ikMetaData.getColumnCount());
Assertions.assertFalse(importedKeys.next());
}
}

@Test
public void testGetExportedKeys_pcntTable() throws SQLException {
try (Connection connection = DriverManager.getConnection(ITBase.connectionUrl);
ResultSet exportedKeys =
connection.getMetaData().getExportedKeys(PROJECT_ID, PCNT_SCHEMA, PCNT_TABLE_NAME)) {
Assertions.assertNotNull(exportedKeys);
ResultSetMetaData ekMetaData = exportedKeys.getMetaData();
Assertions.assertEquals(14, ekMetaData.getColumnCount());
Assertions.assertFalse(exportedKeys.next());
}
}

@Test
public void testGetCrossReference_pcntTable() throws SQLException {
try (Connection connection = DriverManager.getConnection(ITBase.connectionUrl);
ResultSet crossReference =
connection
.getMetaData()
.getCrossReference(
PROJECT_ID,
PCNT_SCHEMA,
PCNT_TABLE_NAME,
PROJECT_ID,
PCNT_SCHEMA,
PCNT_TABLE_NAME)) {
Assertions.assertNotNull(crossReference);
ResultSetMetaData crMetaData = crossReference.getMetaData();
Assertions.assertEquals(14, crMetaData.getColumnCount());
Assertions.assertFalse(crossReference.next());
}
}

@Test
public void testMetadataResultSetsDoNotInterfere() throws SQLException {
try (Connection connection = DriverManager.getConnection(ITBase.connectionUrl)) {
Expand Down
Loading