diff --git a/java-bigquery/google-cloud-bigquery/pom.xml b/java-bigquery/google-cloud-bigquery/pom.xml index 6716595c1995..12e2e97ef0c1 100644 --- a/java-bigquery/google-cloud-bigquery/pom.xml +++ b/java-bigquery/google-cloud-bigquery/pom.xml @@ -120,6 +120,15 @@ arrow-memory-netty + + com.google.api + gax-grpc + + + io.grpc + grpc-api + + com.google.errorprone error_prone_annotations diff --git a/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/BigQueryImpl.java b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/BigQueryImpl.java index 2ad09c33d7cb..07162ee6eb70 100644 --- a/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/BigQueryImpl.java +++ b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/BigQueryImpl.java @@ -22,7 +22,9 @@ import com.google.api.core.BetaApi; import com.google.api.core.InternalApi; +import com.google.api.gax.core.FixedCredentialsProvider; import com.google.api.gax.paging.Page; +import com.google.api.gax.rpc.ServerStream; import com.google.api.services.bigquery.model.ErrorProto; import com.google.api.services.bigquery.model.GetQueryResultsResponse; import com.google.api.services.bigquery.model.ProjectList; @@ -43,6 +45,10 @@ import com.google.cloud.bigquery.InsertAllRequest.RowToInsert; import com.google.cloud.bigquery.spi.v2.BigQueryRpc; import com.google.cloud.bigquery.spi.v2.HttpBigQueryRpc; +import com.google.cloud.bigquery.storage.v1.BigQueryReadClient; +import com.google.cloud.bigquery.storage.v1.BigQueryReadSettings; +import com.google.cloud.bigquery.storage.v1.ReadRowsRequest; +import com.google.cloud.bigquery.storage.v1.ReadRowsResponse; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Function; import com.google.common.base.Strings; @@ -52,14 +58,17 @@ import com.google.common.collect.Iterables; import com.google.common.collect.Lists; import com.google.common.collect.Maps; +import com.google.common.net.HostAndPort; import io.opentelemetry.api.common.Attributes; import io.opentelemetry.api.trace.Span; import io.opentelemetry.context.Scope; import java.io.IOException; +import java.net.URI; import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.concurrent.Callable; +import java.util.concurrent.locks.ReentrantLock; import java.util.regex.Matcher; import java.util.regex.Pattern; import org.checkerframework.checker.nullness.qual.NonNull; @@ -264,6 +273,219 @@ public Page getNextPage() { } } + /** + * Paging implementation that streams subsequent result pages in Apache Arrow format using the + * BigQuery Storage Read API. + */ + private static class ArrowQueryPageFetcher implements NextPageFetcher { + private static final long serialVersionUID = 1L; + + private static final long DEFAULT_PAGE_SIZE = 10000L; + + private final JobId jobId; + private final Schema schema; + private final String arrowSchemaJson; + private final BigQueryOptions serviceOptions; + private final long maxResults; + + private transient Object arrowSchemaPojo; + private final long totalRowsReturned; + + /** + * Constructs an {@link ArrowQueryPageFetcher} for streaming subsequent Arrow result pages. + * + * @param jobId the query job identifier + * @param schema the BigQuery table schema + * @param arrowSchemaPojo the Apache Arrow schema definition + * @param serviceOptions the BigQueryOptions configuration + * @param initialRowOffset the starting row offset within the stream + * @param maxResults the maximum total rows to return across all pages + */ + ArrowQueryPageFetcher( + JobId jobId, + Schema schema, + Object arrowSchemaPojo, + BigQueryOptions serviceOptions, + long initialRowOffset, + Long maxResults) { + this.jobId = jobId; + this.schema = schema; + this.arrowSchemaJson = ArrowDeserializer.arrowSchemaToJson(arrowSchemaPojo); + this.arrowSchemaPojo = arrowSchemaPojo; + this.serviceOptions = serviceOptions; + this.totalRowsReturned = initialRowOffset; + this.maxResults = maxResults != null ? maxResults : Long.MAX_VALUE; + } + + /** + * Fetches the next page of rows by opening a read stream to the query job's default stream and + * consuming up to {@code DEFAULT_PAGE_SIZE} rows. + * + * @return the next {@link Page} of {@link FieldValueList} rows, or null if no more rows are + * available + */ + @Override + public Page getNextPage() { + if (totalRowsReturned >= maxResults) { + return null; + } + + List rowBatch = + new ArrayList<>((int) Math.min(DEFAULT_PAGE_SIZE, maxResults - totalRowsReturned)); + boolean hasMore = false; + BigQueryReadClient ownedClient = null; + + try { + if (arrowSchemaPojo == null && arrowSchemaJson != null) { + arrowSchemaPojo = ArrowDeserializer.jsonToArrowSchema(arrowSchemaJson); + } + + String location = + jobId.getLocation() != null ? jobId.getLocation() : serviceOptions.getLocation(); + if (location == null) { + throw new BigQueryException( + 0, "Job location is required to read Arrow rows from storage stream"); + } + + BigQuery service = serviceOptions.getService(); + BigQueryReadClient client; + if (service instanceof BigQueryImpl) { + client = ((BigQueryImpl) service).getBigQueryReadClient(); + } else { + BigQueryReadSettings.Builder settingsBuilder = BigQueryReadSettings.newBuilder(); + configureReadSettings(settingsBuilder, serviceOptions); + ownedClient = BigQueryReadClient.create(settingsBuilder.build()); + client = ownedClient; + } + + String projectId = + jobId.getProject() != null ? jobId.getProject() : serviceOptions.getProjectId(); + if (projectId == null) { + throw new BigQueryException( + 0, "Project ID is required to read Arrow rows from storage stream"); + } + String streamName = + String.format( + "projects/%s/locations/%s/jobs/%s/streams/_default", + projectId, location, jobId.getJob()); + + ReadRowsRequest readRowsRequest = + ReadRowsRequest.newBuilder() + .setReadStream(streamName) + .setOffset(totalRowsReturned) + .build(); + + ServerStream stream = client.readRowsCallable().call(readRowsRequest); + try { + hasMore = + ArrowDeserializer.loadArrowRows( + stream.iterator(), + arrowSchemaPojo != null ? arrowSchemaPojo : arrowSchemaJson, + schema, + rowBatch, + DEFAULT_PAGE_SIZE, + totalRowsReturned, + maxResults); + } finally { + stream.cancel(); + } + } catch (BigQueryException e) { + throw e; + } catch (Exception e) { + throw new BigQueryException(0, "Failed to read Arrow rows from storage stream", e); + } finally { + if (ownedClient != null) { + try { + ownedClient.close(); + } catch (Exception e) { + // ignore + } + } + } + + if (rowBatch.isEmpty()) { + return null; + } + + long nextOffset = totalRowsReturned + rowBatch.size(); + String nextPageToken = hasMore ? String.valueOf(nextOffset) : null; + ArrowQueryPageFetcher nextPageFetcher = + new ArrowQueryPageFetcher( + jobId, schema, arrowSchemaPojo, serviceOptions, nextOffset, maxResults); + return new PageImpl<>(nextPageFetcher, nextPageToken, rowBatch); + } + } + + private final ReentrantLock readClientLock = new ReentrantLock(); + private transient BigQueryReadClient bqReadClient; + + /** + * Lazily creates or retrieves the shared {@link BigQueryReadClient} instance used for streaming + * Arrow query results, reusing credentials and channel configuration from this {@link + * BigQueryImpl}. + * + * @return the active BigQueryReadClient instance + * @throws IOException if initializing the storage read client fails + */ + BigQueryReadClient getBigQueryReadClient() throws IOException { + readClientLock.lock(); + try { + if (bqReadClient == null) { + BigQueryReadSettings.Builder settingsBuilder = BigQueryReadSettings.newBuilder(); + configureReadSettings(settingsBuilder, getOptions()); + bqReadClient = BigQueryReadClient.create(settingsBuilder.build()); + } + return bqReadClient; + } finally { + readClientLock.unlock(); + } + } + + /** + * Configures a {@link BigQueryReadSettings.Builder} with credentials, universe domain, custom + * endpoint, and transport settings mapped from the given {@link BigQueryOptions}. + * + * @param settingsBuilder the builder to configure + * @param options the source BigQueryOptions + */ + private static void configureReadSettings( + BigQueryReadSettings.Builder settingsBuilder, BigQueryOptions options) { + if (options.getCredentials() != null) { + settingsBuilder.setCredentialsProvider( + FixedCredentialsProvider.create(options.getCredentials())); + } + if (options.getUniverseDomain() != null) { + settingsBuilder.setUniverseDomain(options.getUniverseDomain()); + } + if (options.getHost() != null) { + String host = options.getHost(); + String target = host; + if (target.contains("://")) { + target = URI.create(target).getAuthority(); + } + HostAndPort hostAndPort = HostAndPort.fromString(target); + String endpointHost = hostAndPort.getHost(); + if (endpointHost.contains("bigquery.googleapis.com")) { + endpointHost = + endpointHost.replace("bigquery.googleapis.com", "bigquerystorage.googleapis.com"); + } else if (endpointHost.contains("bigquery.private.googleapis.com")) { + endpointHost = + endpointHost.replace( + "bigquery.private.googleapis.com", "bigquerystorage.private.googleapis.com"); + } else if (endpointHost.startsWith("bigquery.")) { + endpointHost = endpointHost.replaceFirst("^bigquery\\.", "bigquerystorage."); + } + int port = hostAndPort.getPortOrDefault(443); + settingsBuilder.setEndpoint(endpointHost + ":" + port); + if (endpointHost.contains("localhost") || endpointHost.contains("127.0.0.1")) { + settingsBuilder.setTransportChannelProvider( + BigQueryReadSettings.defaultGrpcTransportProviderBuilder() + .setChannelConfigurator(io.grpc.ManagedChannelBuilder::usePlaintext) + .build()); + } + } + } + private final HttpBigQueryRpc bigQueryRpc; private static final BigQueryRetryConfig EMPTY_RETRY_CONFIG = @@ -2077,8 +2299,26 @@ public com.google.api.services.bigquery.model.QueryResponse call() long numRows; Schema schema; - if (results.getJobComplete() && results.getSchema() != null) { - schema = Schema.fromPb(results.getSchema()); + boolean isArrow = false; + Object arrowSchemaPojo = null; + + if (results.getJobComplete()) { + if (results.getArrowSchema() != null) { + isArrow = true; + try { + arrowSchemaPojo = + ArrowDeserializer.deserializeSchema( + results.getArrowSchema().decodeSerializedSchema()); + schema = ArrowDeserializer.arrowSchemaToBigQuerySchema(arrowSchemaPojo); + } catch (IOException e) { + throw new BigQueryException(0, "Failed to deserialize Arrow schema from response", e); + } + } else if (results.getSchema() != null) { + schema = Schema.fromPb(results.getSchema()); + } else { + schema = null; + } + if (results.getNumDmlAffectedRows() == null && results.getTotalRows() == null) { numRows = 0L; } else if (results.getNumDmlAffectedRows() != null) { @@ -2095,45 +2335,75 @@ public com.google.api.services.bigquery.model.QueryResponse call() return job; } + List firstPageRows; + if (isArrow) { + if (results.getArrowRecordBatch() != null) { + try { + firstPageRows = + ArrowDeserializer.deserializeRecordBatch( + results.getArrowRecordBatch().decodeSerializedRecordBatch(), + schema, + arrowSchemaPojo); + } catch (IOException e) { + throw new BigQueryException(0, "Failed to deserialize Arrow record batch", e); + } + } else { + firstPageRows = ImmutableList.of(); + } + } else { + firstPageRows = + ImmutableList.copyOf( + transformTableData( + results.getRows(), + schema, + getOptions().getDataFormatOptions().useInt64Timestamp())); + } + if (results.getPageToken() != null) { JobId jobId = JobId.fromPb(results.getJobReference()); String cursor = results.getPageToken(); + + NextPageFetcher pageFetcher; + if (isArrow) { + long initialRowOffset = (long) firstPageRows.size(); + Map optionsMap = optionMap(options); + Number maxResultsOpt = (Number) optionsMap.get(BigQueryRpc.Option.MAX_RESULTS); + Long maxResults = maxResultsOpt != null ? maxResultsOpt.longValue() : null; + pageFetcher = + new ArrowQueryPageFetcher( + jobId, schema, arrowSchemaPojo, getOptions(), initialRowOffset, maxResults); + } else { + pageFetcher = new QueryPageFetcher(jobId, schema, getOptions(), cursor, optionMap(options)); + } + return TableResult.newBuilder() .setSchema(schema) .setTotalRows(numRows) - .setPageNoSchema( - new PageImpl<>( - // fetch next pages of results - new QueryPageFetcher(jobId, schema, getOptions(), cursor, optionMap(options)), - cursor, - transformTableData( - results.getRows(), - schema, - getOptions().getDataFormatOptions().useInt64Timestamp()))) + .setPageNoSchema(new PageImpl<>(pageFetcher, cursor, firstPageRows)) .setJobId(jobId) .setQueryId(results.getQueryId()) .setJobCreationReason(JobCreationReason.fromPb(results.getJobCreationReason())) - .setRowsInPage(results.getRows() != null ? (long) results.getRows().size() : 0L) + .setRowsInPage((long) firstPageRows.size()) .build(); } - // only 1 page of result + return TableResult.newBuilder() .setSchema(schema) .setTotalRows(numRows) .setPageNoSchema( new PageImpl<>( - new TableDataPageFetcher(null, schema, getOptions(), null, optionMap(options)), + isArrow + ? null + : new TableDataPageFetcher( + null, schema, getOptions(), null, optionMap(options)), null, - transformTableData( - results.getRows(), - schema, - getOptions().getDataFormatOptions().useInt64Timestamp()))) + firstPageRows)) // Return the JobID of the successful job .setJobId( results.getJobReference() != null ? JobId.fromPb(results.getJobReference()) : null) .setQueryId(results.getQueryId()) .setJobCreationReason(JobCreationReason.fromPb(results.getJobCreationReason())) - .setRowsInPage(results.getRows() != null ? (long) results.getRows().size() : 0L) + .setRowsInPage((long) firstPageRows.size()) .build(); } @@ -2207,6 +2477,10 @@ && getOptions().getOpenTelemetryTracer() != null) { return queryRpc(projectId, content, options); } + if (configuration.getQueryResultsFormat() == QueryResultsFormat.ARROW) { + throw new IllegalArgumentException( + "Arrow results format is only supported for fast query path execution (e.g. no destination table, no custom clustering, etc.)."); + } return create(JobInfo.of(jobId, configuration), options); } finally { if (querySpan != null) { diff --git a/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/QueryRequestInfo.java b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/QueryRequestInfo.java index c224bed5cc58..14d2c65fe78a 100644 --- a/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/QueryRequestInfo.java +++ b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/QueryRequestInfo.java @@ -46,6 +46,8 @@ final class QueryRequestInfo { private final DataFormatOptions formatOptions; private final String reservation; private final Long jobTimeoutMs; + private final QueryResultsFormat queryResultsFormat; + private final ArrowSerializationOptions arrowSerializationOptions; QueryRequestInfo( QueryJobConfiguration config, com.google.cloud.bigquery.DataFormatOptions dataFormatOptions) { @@ -63,9 +65,11 @@ final class QueryRequestInfo { this.useLegacySql = config.useLegacySql(); this.useQueryCache = config.useQueryCache(); this.jobCreationMode = config.getJobCreationMode(); - this.formatOptions = dataFormatOptions.toPb(); + this.formatOptions = dataFormatOptions != null ? dataFormatOptions.toPb() : null; this.reservation = config.getReservation(); this.jobTimeoutMs = config.getJobTimeoutMs(); + this.queryResultsFormat = config.getQueryResultsFormat(); + this.arrowSerializationOptions = config.getArrowSerializationOptions(); } /** @@ -142,6 +146,12 @@ QueryRequest toPb() { if (jobTimeoutMs != null) { request.setJobTimeoutMs(jobTimeoutMs); } + if (queryResultsFormat != null) { + request.setQueryResultsFormat(queryResultsFormat.toString()); + } + if (arrowSerializationOptions != null) { + request.setArrowSerializationOptions(arrowSerializationOptions.toPb()); + } return request; } @@ -161,7 +171,7 @@ public String toString() { .add("useQueryCache", useQueryCache) .add("useLegacySql", useLegacySql) .add("jobCreationMode", jobCreationMode) - .add("formatOptions", formatOptions.getUseInt64Timestamp()) + .add("formatOptions", formatOptions != null ? formatOptions.getUseInt64Timestamp() : null) .add("reservation", reservation) .add("jobTimeoutMs", jobTimeoutMs) .toString(); diff --git a/java-bigquery/google-cloud-bigquery/src/test/java/com/google/cloud/bigquery/it/ITBigQueryTest.java b/java-bigquery/google-cloud-bigquery/src/test/java/com/google/cloud/bigquery/it/ITBigQueryTest.java index 3c9613d20758..19c15eb9ee27 100644 --- a/java-bigquery/google-cloud-bigquery/src/test/java/com/google/cloud/bigquery/it/ITBigQueryTest.java +++ b/java-bigquery/google-cloud-bigquery/src/test/java/com/google/cloud/bigquery/it/ITBigQueryTest.java @@ -119,6 +119,7 @@ import com.google.cloud.bigquery.QueryJobConfiguration.JobCreationMode; import com.google.cloud.bigquery.QueryJobConfiguration.Priority; import com.google.cloud.bigquery.QueryParameterValue; +import com.google.cloud.bigquery.QueryResultsFormat; import com.google.cloud.bigquery.Range; import com.google.cloud.bigquery.RangePartitioning; import com.google.cloud.bigquery.Routine; @@ -160,6 +161,7 @@ import com.google.common.collect.ImmutableMap; import com.google.common.collect.ImmutableSet; import com.google.common.collect.Iterables; +import com.google.common.collect.Lists; import com.google.common.collect.Sets; import com.google.common.io.BaseEncoding; import com.google.common.util.concurrent.ListenableFuture; @@ -7561,6 +7563,45 @@ void testQueryWithTimeout() throws InterruptedException { assertTrue(millis < 1_000_000 * 2); } + @Test + void testQueryResultsFormatArrow() throws InterruptedException { + RemoteBigQueryHelper bigqueryHelper = RemoteBigQueryHelper.create(); + BigQuery bigQuery = bigqueryHelper.getOptions().getService(); + String query = "SELECT 1 as id, 'hello' as name, TIMESTAMP('2026-08-10T12:00:00Z') as ts"; + QueryJobConfiguration config = + QueryJobConfiguration.newBuilder(query) + .setQueryResultsFormat(QueryResultsFormat.ARROW) + .setJobCreationMode(JobCreationMode.JOB_CREATION_OPTIONAL) + .build(); + TableResult result = bigQuery.query(config); + assertNotNull(result); + List rows = Lists.newArrayList(result.iterateAll()); + assertEquals(1, rows.size()); + FieldValueList row = rows.get(0); + assertEquals(1L, row.get("id").getLongValue()); + assertEquals("hello", row.get("name").getStringValue()); + assertNotNull(row.get("ts").getValue()); + } + + @Test + void testQueryResultsFormatArrowMultiPage() throws InterruptedException { + RemoteBigQueryHelper bigqueryHelper = RemoteBigQueryHelper.create(); + BigQuery bigQuery = bigqueryHelper.getOptions().getService(); + String query = "SELECT x FROM UNNEST(GENERATE_ARRAY(1, 15000)) AS x"; + QueryJobConfiguration config = + QueryJobConfiguration.newBuilder(query) + .setQueryResultsFormat(QueryResultsFormat.ARROW) + .setJobCreationMode(JobCreationMode.JOB_CREATION_OPTIONAL) + .build(); + TableResult result = bigQuery.query(config); + assertNotNull(result); + List rows = Lists.newArrayList(result.iterateAll()); + assertEquals(15000, rows.size()); + for (int i = 0; i < 15000; i++) { + assertEquals((long) (i + 1), rows.get(i).get(0).getLongValue()); + } + } + @Test void testUniverseDomainWithInvalidUniverseDomain() { RemoteBigQueryHelper bigqueryHelper = RemoteBigQueryHelper.create();