From 7b1352f0a464ba25f6e299704706083300648985 Mon Sep 17 00:00:00 2001 From: Jin Seop Kim Date: Wed, 29 Jul 2026 10:33:48 -0400 Subject: [PATCH 1/7] feat(bigquery): integrate Arrow query response processing and stream pagination --- .../google/cloud/bigquery/BigQueryImpl.java | 253 ++++++++++++++++-- .../cloud/bigquery/QueryRequestInfo.java | 10 + 2 files changed, 245 insertions(+), 18 deletions(-) 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..2bd731f0296f 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; @@ -57,6 +63,7 @@ import io.opentelemetry.context.Scope; import java.io.IOException; import java.util.ArrayList; +import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.concurrent.Callable; @@ -264,6 +271,143 @@ public Page getNextPage() { } } + private static class ArrowQueryPageFetcher implements NextPageFetcher { + private static final long serialVersionUID = 1L; + + private final JobId jobId; + private final Schema schema; + private final org.apache.arrow.vector.types.pojo.Schema arrowSchemaPojo; + private final BigQueryOptions serviceOptions; + private final long maxResults; + + private transient BigQueryReadClient bqReadClient; + private transient ServerStream stream; + private transient Iterator streamIterator; + private long totalRowsReturned = 0L; + private boolean streamClosed = false; + + ArrowQueryPageFetcher( + JobId jobId, + Schema schema, + org.apache.arrow.vector.types.pojo.Schema arrowSchemaPojo, + BigQueryOptions serviceOptions, + long initialRowOffset, + Long maxResults) { + this.jobId = jobId; + this.schema = schema; + this.arrowSchemaPojo = arrowSchemaPojo; + this.serviceOptions = serviceOptions; + this.totalRowsReturned = initialRowOffset; + this.maxResults = maxResults != null ? maxResults : Long.MAX_VALUE; + } + + @Override + public Page getNextPage() { + if (streamClosed || totalRowsReturned >= maxResults) { + closeClient(); + return null; + } + + long pageSize = 100000L; + List rowBatch = new ArrayList<>(); + + try { + if (bqReadClient == null) { + BigQueryReadSettings settings = + BigQueryReadSettings.newBuilder() + .setCredentialsProvider( + FixedCredentialsProvider.create(serviceOptions.getCredentials())) + .build(); + bqReadClient = BigQueryReadClient.create(settings); + } + + if (streamIterator == null) { + String streamName = + String.format( + "projects/%s/locations/%s/jobs/%s/streams/_default", + jobId.getProject() != null ? jobId.getProject() : serviceOptions.getProjectId(), + jobId.getLocation() != null ? jobId.getLocation() : serviceOptions.getLocation(), + jobId.getJob()); + + ReadRowsRequest readRowsRequest = + ReadRowsRequest.newBuilder() + .setReadStream(streamName) + .setOffset(totalRowsReturned) + .build(); + + stream = bqReadClient.readRowsCallable().call(readRowsRequest); + streamIterator = stream.iterator(); + } + + try (org.apache.arrow.memory.BufferAllocator allocator = + new org.apache.arrow.memory.RootAllocator(Long.MAX_VALUE)) { + List vectors = new ArrayList<>(); + for (org.apache.arrow.vector.types.pojo.Field field : arrowSchemaPojo.getFields()) { + vectors.add((org.apache.arrow.vector.FieldVector) field.createVector(allocator)); + } + try (org.apache.arrow.vector.VectorSchemaRoot root = + new org.apache.arrow.vector.VectorSchemaRoot(vectors)) { + org.apache.arrow.vector.VectorLoader loader = + new org.apache.arrow.vector.VectorLoader(root); + + while (rowBatch.size() < pageSize && streamIterator.hasNext()) { + ReadRowsResponse response = streamIterator.next(); + if (response.hasArrowRecordBatch()) { + com.google.cloud.bigquery.storage.v1.ArrowRecordBatch batch = + response.getArrowRecordBatch(); + org.apache.arrow.vector.ipc.message.ArrowRecordBatch deserializedBatch = + org.apache.arrow.vector.ipc.message.MessageSerializer.deserializeRecordBatch( + new org.apache.arrow.vector.ipc.ReadChannel( + new org.apache.arrow.vector.util.ByteArrayReadableSeekableByteChannel( + batch.getSerializedRecordBatch().toByteArray())), + allocator); + loader.load(deserializedBatch); + deserializedBatch.close(); + int batchRowCount = root.getRowCount(); + for (int i = 0; i < batchRowCount; i++) { + rowBatch.add(ArrowDeserializer.arrowRootToFieldValueList(root, i, schema)); + } + root.clear(); + } + } + } + } + + if (rowBatch.isEmpty()) { + streamClosed = true; + closeClient(); + return null; + } + + totalRowsReturned += rowBatch.size(); + + String nextPageToken = null; + if (streamIterator.hasNext() && totalRowsReturned < maxResults) { + nextPageToken = String.valueOf(totalRowsReturned); + } else { + streamClosed = true; + closeClient(); + } + + return new PageImpl<>(this, nextPageToken, rowBatch); + + } catch (Exception e) { + streamClosed = true; + closeClient(); + throw new BigQueryException(0, "Failed to read Arrow rows from storage stream", e); + } + } + + private void closeClient() { + if (bqReadClient != null) { + bqReadClient.close(); + bqReadClient = null; + } + streamIterator = null; + stream = null; + } + } + private final HttpBigQueryRpc bigQueryRpc; private static final BigQueryRetryConfig EMPTY_RETRY_CONFIG = @@ -2077,8 +2221,28 @@ 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; + org.apache.arrow.vector.types.pojo.Schema arrowSchemaPojo = null; + + if (results.getJobComplete()) { + if (results.getSchema() != null) { + schema = Schema.fromPb(results.getSchema()); + } else if (results.getArrowSchema() != null) { + isArrow = true; + try { + arrowSchemaPojo = + org.apache.arrow.vector.ipc.message.MessageSerializer.deserializeSchema( + new org.apache.arrow.vector.ipc.ReadChannel( + new org.apache.arrow.vector.util.ByteArrayReadableSeekableByteChannel( + results.getArrowSchema().decodeSerializedSchema()))); + schema = ArrowDeserializer.arrowSchemaToBigQuerySchema(arrowSchemaPojo); + } catch (IOException e) { + throw new BigQueryException(0, "Failed to deserialize Arrow schema from response", e); + } + } else { + schema = null; + } + if (results.getNumDmlAffectedRows() == null && results.getTotalRows() == null) { numRows = 0L; } else if (results.getNumDmlAffectedRows() != null) { @@ -2098,42 +2262,91 @@ public com.google.api.services.bigquery.model.QueryResponse call() if (results.getPageToken() != null) { JobId jobId = JobId.fromPb(results.getJobReference()); String cursor = results.getPageToken(); + + Iterable firstPageRows; + NextPageFetcher pageFetcher; + + 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(); + } + long initialRowOffset = + firstPageRows instanceof List ? ((List) firstPageRows).size() : 0L; + pageFetcher = + new ArrowQueryPageFetcher( + jobId, + schema, + arrowSchemaPojo, + getOptions(), + initialRowOffset, + null); // Or use maxResults from configuration if available + } else { + firstPageRows = + transformTableData( + results.getRows(), schema, getOptions().getDataFormatOptions().useInt64Timestamp()); + 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( + firstPageRows instanceof List ? (long) ((List) firstPageRows).size() : 0L) .build(); } // only 1 page of result + Iterable 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 = + transformTableData( + results.getRows(), schema, getOptions().getDataFormatOptions().useInt64Timestamp()); + } + 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(firstPageRows instanceof List ? (long) ((List) firstPageRows).size() : 0L) .build(); } @@ -2207,6 +2420,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..ec09d140a28d 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) { @@ -66,6 +68,8 @@ final class QueryRequestInfo { this.formatOptions = dataFormatOptions.toPb(); 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; } From 687badb5cad5798bcd0cab38c1557159f2092c8b Mon Sep 17 00:00:00 2001 From: Jin Seop Kim Date: Fri, 7 Aug 2026 15:51:16 -0400 Subject: [PATCH 2/7] refactor(bigquery): import Arrow IPC classes and simplify FQCNs --- .../java/com/google/cloud/bigquery/BigQueryImpl.java | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) 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 2bd731f0296f..546fc8c66301 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 @@ -69,6 +69,9 @@ import java.util.concurrent.Callable; import java.util.regex.Matcher; import java.util.regex.Pattern; +import org.apache.arrow.vector.ipc.ReadChannel; +import org.apache.arrow.vector.ipc.message.MessageSerializer; +import org.apache.arrow.vector.util.ByteArrayReadableSeekableByteChannel; import org.checkerframework.checker.nullness.qual.NonNull; final class BigQueryImpl extends BaseService implements BigQuery { @@ -2231,9 +2234,9 @@ public com.google.api.services.bigquery.model.QueryResponse call() isArrow = true; try { arrowSchemaPojo = - org.apache.arrow.vector.ipc.message.MessageSerializer.deserializeSchema( - new org.apache.arrow.vector.ipc.ReadChannel( - new org.apache.arrow.vector.util.ByteArrayReadableSeekableByteChannel( + MessageSerializer.deserializeSchema( + new ReadChannel( + new ByteArrayReadableSeekableByteChannel( results.getArrowSchema().decodeSerializedSchema()))); schema = ArrowDeserializer.arrowSchemaToBigQuerySchema(arrowSchemaPojo); } catch (IOException e) { From 6b30416582de56e7eee582b74db7709900aec7bc Mon Sep 17 00:00:00 2001 From: Jin Seop Kim Date: Fri, 7 Aug 2026 15:53:49 -0400 Subject: [PATCH 3/7] refactor(bigquery): clean up DataFormatOptions constructor parameter type in QueryRequestInfo --- .../java/com/google/cloud/bigquery/QueryRequestInfo.java | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) 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 ec09d140a28d..6255f330c00a 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 @@ -16,7 +16,6 @@ package com.google.cloud.bigquery; -import com.google.api.services.bigquery.model.DataFormatOptions; import com.google.api.services.bigquery.model.QueryParameter; import com.google.api.services.bigquery.model.QueryRequest; import com.google.cloud.bigquery.QueryJobConfiguration.JobCreationMode; @@ -43,14 +42,13 @@ final class QueryRequestInfo { private final Boolean useQueryCache; private final Boolean useLegacySql; private final JobCreationMode jobCreationMode; - private final DataFormatOptions formatOptions; + private final com.google.api.services.bigquery.model.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) { + QueryRequestInfo(QueryJobConfiguration config, DataFormatOptions dataFormatOptions) { this.config = config; this.connectionProperties = config.getConnectionProperties(); this.defaultDataset = config.getDefaultDataset(); From aeb011caeb389fe447b71cc241133908bd63d8fd Mon Sep 17 00:00:00 2001 From: Jin Seop Kim Date: Fri, 7 Aug 2026 15:56:22 -0400 Subject: [PATCH 4/7] refactor(bigquery): import Arrow IPC classes and simplify FQCNs --- .../bigquery/ArrowSerializationOptions.java | 3 +-- .../google/cloud/bigquery/BigQueryImpl.java | 27 ++++++++++--------- 2 files changed, 16 insertions(+), 14 deletions(-) diff --git a/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowSerializationOptions.java b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowSerializationOptions.java index 1ad88f8de84b..1c1c6295360f 100644 --- a/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowSerializationOptions.java +++ b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowSerializationOptions.java @@ -77,8 +77,7 @@ com.google.api.services.bigquery.model.ArrowSerializationOptions toPb() { return ArrowSerializationOptionsConverter.toPb(this); } - static ArrowSerializationOptions fromPb( - com.google.api.services.bigquery.model.ArrowSerializationOptions optionsPb) { + static ArrowSerializationOptions fromPb(Object optionsPb) { return ArrowSerializationOptionsConverter.fromPb(optionsPb); } 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 546fc8c66301..d359d5fc6f2e 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 @@ -69,7 +69,13 @@ import java.util.concurrent.Callable; import java.util.regex.Matcher; import java.util.regex.Pattern; +import org.apache.arrow.memory.BufferAllocator; +import org.apache.arrow.memory.RootAllocator; +import org.apache.arrow.vector.FieldVector; +import org.apache.arrow.vector.VectorLoader; +import org.apache.arrow.vector.VectorSchemaRoot; import org.apache.arrow.vector.ipc.ReadChannel; +import org.apache.arrow.vector.ipc.message.ArrowRecordBatch; import org.apache.arrow.vector.ipc.message.MessageSerializer; import org.apache.arrow.vector.util.ByteArrayReadableSeekableByteChannel; import org.checkerframework.checker.nullness.qual.NonNull; @@ -342,26 +348,23 @@ public Page getNextPage() { streamIterator = stream.iterator(); } - try (org.apache.arrow.memory.BufferAllocator allocator = - new org.apache.arrow.memory.RootAllocator(Long.MAX_VALUE)) { - List vectors = new ArrayList<>(); + try (BufferAllocator allocator = new RootAllocator(Long.MAX_VALUE)) { + List vectors = new ArrayList<>(); for (org.apache.arrow.vector.types.pojo.Field field : arrowSchemaPojo.getFields()) { - vectors.add((org.apache.arrow.vector.FieldVector) field.createVector(allocator)); + vectors.add((FieldVector) field.createVector(allocator)); } - try (org.apache.arrow.vector.VectorSchemaRoot root = - new org.apache.arrow.vector.VectorSchemaRoot(vectors)) { - org.apache.arrow.vector.VectorLoader loader = - new org.apache.arrow.vector.VectorLoader(root); + try (VectorSchemaRoot root = new VectorSchemaRoot(vectors)) { + VectorLoader loader = new VectorLoader(root); while (rowBatch.size() < pageSize && streamIterator.hasNext()) { ReadRowsResponse response = streamIterator.next(); if (response.hasArrowRecordBatch()) { com.google.cloud.bigquery.storage.v1.ArrowRecordBatch batch = response.getArrowRecordBatch(); - org.apache.arrow.vector.ipc.message.ArrowRecordBatch deserializedBatch = - org.apache.arrow.vector.ipc.message.MessageSerializer.deserializeRecordBatch( - new org.apache.arrow.vector.ipc.ReadChannel( - new org.apache.arrow.vector.util.ByteArrayReadableSeekableByteChannel( + ArrowRecordBatch deserializedBatch = + MessageSerializer.deserializeRecordBatch( + new ReadChannel( + new ByteArrayReadableSeekableByteChannel( batch.getSerializedRecordBatch().toByteArray())), allocator); loader.load(deserializedBatch); From c2af8738e30cf91434f37fd4c6a49085da2bd621 Mon Sep 17 00:00:00 2001 From: Jin Seop Kim Date: Fri, 7 Aug 2026 15:59:59 -0400 Subject: [PATCH 5/7] refactor(bigquery): delegate Arrow Pojo Schema/Field conversions to ArrowPojoUtils --- .../cloud/bigquery/ArrowDeserializer.java | 136 +++--------------- .../google/cloud/bigquery/ArrowPojoUtils.java | 116 +++++++++++++++ .../google/cloud/bigquery/BigQueryImpl.java | 11 +- 3 files changed, 139 insertions(+), 124 deletions(-) create mode 100644 java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowPojoUtils.java diff --git a/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowDeserializer.java b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowDeserializer.java index 7c5829734822..5d2433330460 100644 --- a/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowDeserializer.java +++ b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowDeserializer.java @@ -45,105 +45,13 @@ final class ArrowDeserializer { private ArrowDeserializer() {} /** - * Converts an Apache Arrow {@link org.apache.arrow.vector.types.pojo.Schema} to a BigQuery Veneer - * {@link Schema}. + * Converts an Apache Arrow Schema to a BigQuery Veneer {@link Schema}. * * @param arrowSchema the Apache Arrow schema to convert * @return the corresponding BigQuery Veneer Schema */ - static Schema arrowSchemaToBigQuerySchema(org.apache.arrow.vector.types.pojo.Schema arrowSchema) { - List fields = new ArrayList<>(); - for (org.apache.arrow.vector.types.pojo.Field arrowField : arrowSchema.getFields()) { - fields.add(arrowFieldToBigQueryField(arrowField)); - } - return Schema.of(fields); - } - - /** - * Recursively converts an Apache Arrow {@link org.apache.arrow.vector.types.pojo.Field} to a - * BigQuery Veneer {@link Field}. - * - * @param arrowField the Arrow field to convert - * @return the corresponding BigQuery Veneer Field - */ - private static Field arrowFieldToBigQueryField( - org.apache.arrow.vector.types.pojo.Field arrowField) { - String name = arrowField.getName(); - ArrowType type = arrowField.getType(); - Field.Builder builder; - - if (type instanceof ArrowType.List) { - if (arrowField.getChildren().isEmpty()) { - throw new IllegalArgumentException( - "Arrow List field must have at least one child field: " + name); - } - org.apache.arrow.vector.types.pojo.Field innerField = arrowField.getChildren().get(0); - LegacySQLTypeName innerType = arrowTypeToLegacySQLTypeName(innerField.getType()); - builder = Field.newBuilder(name, innerType); - builder.setMode(Field.Mode.REPEATED); - if (!innerField.getChildren().isEmpty()) { - List subFields = new ArrayList<>(); - for (org.apache.arrow.vector.types.pojo.Field childField : innerField.getChildren()) { - subFields.add(arrowFieldToBigQueryField(childField)); - } - builder.setType(LegacySQLTypeName.RECORD, FieldList.of(subFields)); - } - } else { - LegacySQLTypeName bqType = arrowTypeToLegacySQLTypeName(type); - builder = Field.newBuilder(name, bqType); - if (arrowField.isNullable()) { - builder.setMode(Field.Mode.NULLABLE); - } else { - builder.setMode(Field.Mode.REQUIRED); - } - if (!arrowField.getChildren().isEmpty()) { - List subFields = new ArrayList<>(); - for (org.apache.arrow.vector.types.pojo.Field childField : innerFieldChildren(arrowField)) { - subFields.add(arrowFieldToBigQueryField(childField)); - } - builder.setType(LegacySQLTypeName.RECORD, FieldList.of(subFields)); - } - } - return builder.build(); - } - - private static List innerFieldChildren( - org.apache.arrow.vector.types.pojo.Field arrowField) { - return arrowField.getChildren(); - } - - /** - * Maps an Apache Arrow data type {@link ArrowType} to a BigQuery {@link LegacySQLTypeName}. - * - * @param type the Arrow data type to map - * @return the corresponding BigQuery LegacySQLTypeName - * @throws IllegalArgumentException if the Arrow type is unsupported - */ - private static LegacySQLTypeName arrowTypeToLegacySQLTypeName(ArrowType type) { - switch (type.getTypeID()) { - case Int: - return LegacySQLTypeName.INTEGER; - case FloatingPoint: - return LegacySQLTypeName.FLOAT; - case Utf8: - return LegacySQLTypeName.STRING; - case Bool: - return LegacySQLTypeName.BOOLEAN; - case Binary: - return LegacySQLTypeName.BYTES; - case Decimal: - return LegacySQLTypeName.NUMERIC; - case Timestamp: - return LegacySQLTypeName.TIMESTAMP; - case Date: - return LegacySQLTypeName.DATE; - case Time: - return LegacySQLTypeName.TIME; - case Struct: - return LegacySQLTypeName.RECORD; - default: - throw new IllegalArgumentException("Unsupported Arrow type: " + type.getTypeID()); - } + static Schema arrowSchemaToBigQuerySchema(Object arrowSchema) { + return ArrowPojoUtils.arrowSchemaToBigQuerySchema(arrowSchema); } /** @@ -160,13 +68,25 @@ private static LegacySQLTypeName arrowTypeToLegacySQLTypeName(ArrowType type) { * @throws IOException if deserialization of the Arrow record batch fails */ static List deserializeRecordBatch( - byte[] recordBatchBytes, Schema schema, org.apache.arrow.vector.types.pojo.Schema arrowSchema) + byte[] recordBatchBytes, Schema schema, Object arrowSchema) throws IOException { try (BufferAllocator allocator = new RootAllocator(Long.MAX_VALUE)) { - List vectors = new ArrayList<>(); + List vectors = ArrowPojoUtils.createVectors(arrowSchema, allocator); try { - for (org.apache.arrow.vector.types.pojo.Field field : arrowSchema.getFields()) { - vectors.add(field.createVector(allocator)); + try (VectorSchemaRoot root = new VectorSchemaRoot(vectors)) { + VectorLoader loader = new VectorLoader(root); + try (ArrowRecordBatch deserializedBatch = + MessageSerializer.deserializeRecordBatch( + new ReadChannel(new ByteArrayReadableSeekableByteChannel(recordBatchBytes)), + allocator)) { + loader.load(deserializedBatch); + int rowCount = root.getRowCount(); + List rows = new ArrayList<>(rowCount); + for (int i = 0; i < rowCount; i++) { + rows.add(arrowRootToFieldValueList(root, i, schema)); + } + return ImmutableList.copyOf(rows); + } } } catch (Throwable t) { for (int i = vectors.size() - 1; i >= 0; i--) { @@ -178,21 +98,6 @@ static List deserializeRecordBatch( } throw t; } - try (VectorSchemaRoot root = new VectorSchemaRoot(vectors)) { - VectorLoader loader = new VectorLoader(root); - try (ArrowRecordBatch deserializedBatch = - MessageSerializer.deserializeRecordBatch( - new ReadChannel(new ByteArrayReadableSeekableByteChannel(recordBatchBytes)), - allocator)) { - loader.load(deserializedBatch); - int rowCount = root.getRowCount(); - List rows = new ArrayList<>(rowCount); - for (int i = 0; i < rowCount; i++) { - rows.add(arrowRootToFieldValueList(root, i, schema)); - } - return ImmutableList.copyOf(rows); - } - } } } @@ -281,9 +186,6 @@ private static FieldValue arrowVectorToFieldValue( // Handle primitive types String stringVal; if (bqField.getType() == LegacySQLTypeName.TIMESTAMP) { - // Arrow timestamps are long values representing epoch seconds/millis/micros/nanos. - // Standard BigQuery JSON returns timestamps as string of epoch seconds with micro precision - // (e.g. "1408452095.220000"). TimeStampVector tsVector = (TimeStampVector) vector; long rawVal = tsVector.get(rowIndex); ArrowType.Timestamp tsType = (ArrowType.Timestamp) vector.getField().getType(); diff --git a/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowPojoUtils.java b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowPojoUtils.java new file mode 100644 index 000000000000..828beceb23d9 --- /dev/null +++ b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowPojoUtils.java @@ -0,0 +1,116 @@ +/* + * 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 + * + * http://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. + */ + +package com.google.cloud.bigquery; + +import java.util.ArrayList; +import java.util.List; +import org.apache.arrow.memory.BufferAllocator; +import org.apache.arrow.vector.FieldVector; +import org.apache.arrow.vector.types.pojo.ArrowType; +import org.apache.arrow.vector.types.pojo.Field; +import org.apache.arrow.vector.types.pojo.Schema; + +/** Internal helper for Apache Arrow Schema/Field conversions. */ +final class ArrowPojoUtils { + + private ArrowPojoUtils() {} + + static com.google.cloud.bigquery.Schema arrowSchemaToBigQuerySchema(Object arrowSchemaObj) { + Schema arrowSchema = (Schema) arrowSchemaObj; + List fields = new ArrayList<>(); + for (Field arrowField : arrowSchema.getFields()) { + fields.add(arrowFieldToBigQueryField(arrowField)); + } + return com.google.cloud.bigquery.Schema.of(fields); + } + + static com.google.cloud.bigquery.Field arrowFieldToBigQueryField(Field arrowField) { + String name = arrowField.getName(); + ArrowType type = arrowField.getType(); + com.google.cloud.bigquery.Field.Builder builder; + + if (type instanceof ArrowType.List) { + if (arrowField.getChildren().isEmpty()) { + throw new IllegalArgumentException( + "Arrow List field must have at least one child field: " + name); + } + Field innerField = arrowField.getChildren().get(0); + LegacySQLTypeName innerType = arrowTypeToLegacySQLTypeName(innerField.getType()); + builder = com.google.cloud.bigquery.Field.newBuilder(name, innerType); + builder.setMode(com.google.cloud.bigquery.Field.Mode.REPEATED); + if (!innerField.getChildren().isEmpty()) { + List subFields = new ArrayList<>(); + for (Field childField : innerField.getChildren()) { + subFields.add(arrowFieldToBigQueryField(childField)); + } + builder.setType(LegacySQLTypeName.RECORD, FieldList.of(subFields)); + } + } else { + LegacySQLTypeName bqType = arrowTypeToLegacySQLTypeName(type); + builder = com.google.cloud.bigquery.Field.newBuilder(name, bqType); + if (arrowField.isNullable()) { + builder.setMode(com.google.cloud.bigquery.Field.Mode.NULLABLE); + } else { + builder.setMode(com.google.cloud.bigquery.Field.Mode.REQUIRED); + } + if (!arrowField.getChildren().isEmpty()) { + List subFields = new ArrayList<>(); + for (Field childField : arrowField.getChildren()) { + subFields.add(arrowFieldToBigQueryField(childField)); + } + builder.setType(LegacySQLTypeName.RECORD, FieldList.of(subFields)); + } + } + return builder.build(); + } + + static List createVectors(Object arrowSchemaObj, BufferAllocator allocator) { + Schema arrowSchema = (Schema) arrowSchemaObj; + List vectors = new ArrayList<>(); + for (Field field : arrowSchema.getFields()) { + vectors.add(field.createVector(allocator)); + } + return vectors; + } + + private static LegacySQLTypeName arrowTypeToLegacySQLTypeName(ArrowType type) { + switch (type.getTypeID()) { + case Int: + return LegacySQLTypeName.INTEGER; + case FloatingPoint: + return LegacySQLTypeName.FLOAT; + case Utf8: + return LegacySQLTypeName.STRING; + case Bool: + return LegacySQLTypeName.BOOLEAN; + case Binary: + return LegacySQLTypeName.BYTES; + case Decimal: + return LegacySQLTypeName.NUMERIC; + case Timestamp: + return LegacySQLTypeName.TIMESTAMP; + case Date: + return LegacySQLTypeName.DATE; + case Time: + return LegacySQLTypeName.TIME; + case Struct: + return LegacySQLTypeName.RECORD; + default: + throw new IllegalArgumentException("Unsupported Arrow type: " + type.getTypeID()); + } + } +} 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 d359d5fc6f2e..d40724363f6c 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 @@ -285,7 +285,7 @@ private static class ArrowQueryPageFetcher implements NextPageFetcher getNextPage() { } try (BufferAllocator allocator = new RootAllocator(Long.MAX_VALUE)) { - List vectors = new ArrayList<>(); - for (org.apache.arrow.vector.types.pojo.Field field : arrowSchemaPojo.getFields()) { - vectors.add((FieldVector) field.createVector(allocator)); - } + List vectors = ArrowPojoUtils.createVectors(arrowSchemaPojo, allocator); try (VectorSchemaRoot root = new VectorSchemaRoot(vectors)) { VectorLoader loader = new VectorLoader(root); @@ -2228,7 +2225,7 @@ public com.google.api.services.bigquery.model.QueryResponse call() long numRows; Schema schema; boolean isArrow = false; - org.apache.arrow.vector.types.pojo.Schema arrowSchemaPojo = null; + Object arrowSchemaPojo = null; if (results.getJobComplete()) { if (results.getSchema() != null) { From 9be8cdd5591d3e5f6f1025848659f87e1e4e7561 Mon Sep 17 00:00:00 2001 From: Jin Seop Kim Date: Fri, 7 Aug 2026 16:12:56 -0400 Subject: [PATCH 6/7] refactor(bigquery): eliminate FQCNs in ArrowPojoUtils, QueryRequestInfo, and ArrowDeserializerTest --- .../cloud/bigquery/ArrowDeserializer.java | 3 +- .../google/cloud/bigquery/ArrowPojoUtils.java | 44 ++++++++-------- .../google/cloud/bigquery/BigQueryImpl.java | 4 +- .../cloud/bigquery/QueryRequestInfo.java | 22 +++----- .../cloud/bigquery/ArrowDeserializerTest.java | 50 +++++++++---------- 5 files changed, 56 insertions(+), 67 deletions(-) diff --git a/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowDeserializer.java b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowDeserializer.java index 5d2433330460..e131c157d6d1 100644 --- a/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowDeserializer.java +++ b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowDeserializer.java @@ -68,8 +68,7 @@ static Schema arrowSchemaToBigQuerySchema(Object arrowSchema) { * @throws IOException if deserialization of the Arrow record batch fails */ static List deserializeRecordBatch( - byte[] recordBatchBytes, Schema schema, Object arrowSchema) - throws IOException { + byte[] recordBatchBytes, Schema schema, Object arrowSchema) throws IOException { try (BufferAllocator allocator = new RootAllocator(Long.MAX_VALUE)) { List vectors = ArrowPojoUtils.createVectors(arrowSchema, allocator); try { diff --git a/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowPojoUtils.java b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowPojoUtils.java index 828beceb23d9..f1bfa3e501a0 100644 --- a/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowPojoUtils.java +++ b/java-bigquery/google-cloud-bigquery/src/main/java/com/google/cloud/bigquery/ArrowPojoUtils.java @@ -21,55 +21,56 @@ import org.apache.arrow.memory.BufferAllocator; import org.apache.arrow.vector.FieldVector; import org.apache.arrow.vector.types.pojo.ArrowType; -import org.apache.arrow.vector.types.pojo.Field; -import org.apache.arrow.vector.types.pojo.Schema; /** Internal helper for Apache Arrow Schema/Field conversions. */ final class ArrowPojoUtils { private ArrowPojoUtils() {} - static com.google.cloud.bigquery.Schema arrowSchemaToBigQuerySchema(Object arrowSchemaObj) { - Schema arrowSchema = (Schema) arrowSchemaObj; - List fields = new ArrayList<>(); - for (Field arrowField : arrowSchema.getFields()) { + static Schema arrowSchemaToBigQuerySchema(Object arrowSchemaObj) { + org.apache.arrow.vector.types.pojo.Schema arrowSchema = + (org.apache.arrow.vector.types.pojo.Schema) arrowSchemaObj; + List fields = new ArrayList<>(); + for (org.apache.arrow.vector.types.pojo.Field arrowField : arrowSchema.getFields()) { fields.add(arrowFieldToBigQueryField(arrowField)); } - return com.google.cloud.bigquery.Schema.of(fields); + return Schema.of(fields); } - static com.google.cloud.bigquery.Field arrowFieldToBigQueryField(Field arrowField) { + static Field arrowFieldToBigQueryField(Object arrowFieldObj) { + org.apache.arrow.vector.types.pojo.Field arrowField = + (org.apache.arrow.vector.types.pojo.Field) arrowFieldObj; String name = arrowField.getName(); ArrowType type = arrowField.getType(); - com.google.cloud.bigquery.Field.Builder builder; + Field.Builder builder; if (type instanceof ArrowType.List) { if (arrowField.getChildren().isEmpty()) { throw new IllegalArgumentException( "Arrow List field must have at least one child field: " + name); } - Field innerField = arrowField.getChildren().get(0); + org.apache.arrow.vector.types.pojo.Field innerField = arrowField.getChildren().get(0); LegacySQLTypeName innerType = arrowTypeToLegacySQLTypeName(innerField.getType()); - builder = com.google.cloud.bigquery.Field.newBuilder(name, innerType); - builder.setMode(com.google.cloud.bigquery.Field.Mode.REPEATED); + builder = Field.newBuilder(name, innerType); + builder.setMode(Field.Mode.REPEATED); if (!innerField.getChildren().isEmpty()) { - List subFields = new ArrayList<>(); - for (Field childField : innerField.getChildren()) { + List subFields = new ArrayList<>(); + for (org.apache.arrow.vector.types.pojo.Field childField : innerField.getChildren()) { subFields.add(arrowFieldToBigQueryField(childField)); } builder.setType(LegacySQLTypeName.RECORD, FieldList.of(subFields)); } } else { LegacySQLTypeName bqType = arrowTypeToLegacySQLTypeName(type); - builder = com.google.cloud.bigquery.Field.newBuilder(name, bqType); + builder = Field.newBuilder(name, bqType); if (arrowField.isNullable()) { - builder.setMode(com.google.cloud.bigquery.Field.Mode.NULLABLE); + builder.setMode(Field.Mode.NULLABLE); } else { - builder.setMode(com.google.cloud.bigquery.Field.Mode.REQUIRED); + builder.setMode(Field.Mode.REQUIRED); } if (!arrowField.getChildren().isEmpty()) { - List subFields = new ArrayList<>(); - for (Field childField : arrowField.getChildren()) { + List subFields = new ArrayList<>(); + for (org.apache.arrow.vector.types.pojo.Field childField : arrowField.getChildren()) { subFields.add(arrowFieldToBigQueryField(childField)); } builder.setType(LegacySQLTypeName.RECORD, FieldList.of(subFields)); @@ -79,9 +80,10 @@ static com.google.cloud.bigquery.Field arrowFieldToBigQueryField(Field arrowFiel } static List createVectors(Object arrowSchemaObj, BufferAllocator allocator) { - Schema arrowSchema = (Schema) arrowSchemaObj; + org.apache.arrow.vector.types.pojo.Schema arrowSchema = + (org.apache.arrow.vector.types.pojo.Schema) arrowSchemaObj; List vectors = new ArrayList<>(); - for (Field field : arrowSchema.getFields()) { + for (org.apache.arrow.vector.types.pojo.Field field : arrowSchema.getFields()) { vectors.add(field.createVector(allocator)); } return vectors; 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 d40724363f6c..cc4ca5412e9f 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 @@ -45,6 +45,7 @@ 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.ArrowRecordBatch; import com.google.cloud.bigquery.storage.v1.BigQueryReadClient; import com.google.cloud.bigquery.storage.v1.BigQueryReadSettings; import com.google.cloud.bigquery.storage.v1.ReadRowsRequest; @@ -356,8 +357,7 @@ public Page getNextPage() { while (rowBatch.size() < pageSize && streamIterator.hasNext()) { ReadRowsResponse response = streamIterator.next(); if (response.hasArrowRecordBatch()) { - com.google.cloud.bigquery.storage.v1.ArrowRecordBatch batch = - response.getArrowRecordBatch(); + ArrowRecordBatch batch = response.getArrowRecordBatch(); ArrowRecordBatch deserializedBatch = MessageSerializer.deserializeRecordBatch( new ReadChannel( 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 6255f330c00a..ea5df0ff2cf2 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 @@ -16,6 +16,7 @@ package com.google.cloud.bigquery; +import com.google.api.services.bigquery.model.DataFormatOptions; import com.google.api.services.bigquery.model.QueryParameter; import com.google.api.services.bigquery.model.QueryRequest; import com.google.cloud.bigquery.QueryJobConfiguration.JobCreationMode; @@ -42,13 +43,14 @@ final class QueryRequestInfo { private final Boolean useQueryCache; private final Boolean useLegacySql; private final JobCreationMode jobCreationMode; - private final com.google.api.services.bigquery.model.DataFormatOptions formatOptions; + private final DataFormatOptions formatOptions; private final String reservation; private final Long jobTimeoutMs; private final QueryResultsFormat queryResultsFormat; private final ArrowSerializationOptions arrowSerializationOptions; - QueryRequestInfo(QueryJobConfiguration config, DataFormatOptions dataFormatOptions) { + QueryRequestInfo( + QueryJobConfiguration config, com.google.cloud.bigquery.DataFormatOptions dataFormatOptions) { this.config = config; this.connectionProperties = config.getConnectionProperties(); this.defaultDataset = config.getDefaultDataset(); @@ -63,25 +65,13 @@ 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(); } - /** - * Determines if the query can be executed via the "fast query" path (jobs.query API) instead of - * the "slow path" (jobs.insert API followed by jobs.getQueryResults). - * - *

The fast query path is preferred because it completes in a single RPC, significantly - * reducing end-to-end latency for small queries. - * - *

However, the jobs.query API does not support all configuration options available in - * jobs.insert (e.g., destination table, clustering, time partitioning). This method checks the - * QueryJobConfiguration for any unsupported options. If any are present, we must fall back to the - * jobs.insert path. - */ boolean isFastQuerySupported() { return config.getClustering() == null && config.getCreateDisposition() == null @@ -169,7 +159,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/ArrowDeserializerTest.java b/java-bigquery/google-cloud-bigquery/src/test/java/com/google/cloud/bigquery/ArrowDeserializerTest.java index f23dbe4b01a5..b1524d49d603 100644 --- a/java-bigquery/google-cloud-bigquery/src/test/java/com/google/cloud/bigquery/ArrowDeserializerTest.java +++ b/java-bigquery/google-cloud-bigquery/src/test/java/com/google/cloud/bigquery/ArrowDeserializerTest.java @@ -43,39 +43,42 @@ import org.apache.arrow.vector.ipc.message.MessageSerializer; import org.apache.arrow.vector.types.TimeUnit; import org.apache.arrow.vector.types.pojo.ArrowType; -import org.apache.arrow.vector.types.pojo.Field; import org.apache.arrow.vector.types.pojo.FieldType; -import org.apache.arrow.vector.types.pojo.Schema; import org.junit.jupiter.api.Test; public class ArrowDeserializerTest { @Test public void testArrowSchemaToBigQuerySchema() { - Field intField = new Field("int_col", FieldType.nullable(new ArrowType.Int(32, true)), null); - Field strField = new Field("str_col", FieldType.notNullable(new ArrowType.Utf8()), null); - Field boolField = new Field("bool_col", FieldType.nullable(new ArrowType.Bool()), null); - Field tsField = - new Field( + org.apache.arrow.vector.types.pojo.Field intField = + new org.apache.arrow.vector.types.pojo.Field( + "int_col", FieldType.nullable(new ArrowType.Int(32, true)), null); + org.apache.arrow.vector.types.pojo.Field strField = + new org.apache.arrow.vector.types.pojo.Field( + "str_col", FieldType.notNullable(new ArrowType.Utf8()), null); + org.apache.arrow.vector.types.pojo.Field boolField = + new org.apache.arrow.vector.types.pojo.Field( + "bool_col", FieldType.nullable(new ArrowType.Bool()), null); + org.apache.arrow.vector.types.pojo.Field tsField = + new org.apache.arrow.vector.types.pojo.Field( "ts_col", FieldType.nullable(new ArrowType.Timestamp(TimeUnit.MICROSECOND, "UTC")), null); - Schema arrowSchema = new Schema(ImmutableList.of(intField, strField, boolField, tsField)); + org.apache.arrow.vector.types.pojo.Schema arrowSchema = + new org.apache.arrow.vector.types.pojo.Schema( + ImmutableList.of(intField, strField, boolField, tsField)); - com.google.cloud.bigquery.Schema bqSchema = - ArrowDeserializer.arrowSchemaToBigQuerySchema(arrowSchema); + Schema bqSchema = ArrowDeserializer.arrowSchemaToBigQuerySchema(arrowSchema); assertEquals(4, bqSchema.getFields().size()); assertEquals("int_col", bqSchema.getFields().get(0).getName()); assertEquals(LegacySQLTypeName.INTEGER, bqSchema.getFields().get(0).getType()); - assertEquals( - com.google.cloud.bigquery.Field.Mode.NULLABLE, bqSchema.getFields().get(0).getMode()); + assertEquals(Field.Mode.NULLABLE, bqSchema.getFields().get(0).getMode()); assertEquals("str_col", bqSchema.getFields().get(1).getName()); assertEquals(LegacySQLTypeName.STRING, bqSchema.getFields().get(1).getType()); - assertEquals( - com.google.cloud.bigquery.Field.Mode.REQUIRED, bqSchema.getFields().get(1).getMode()); + assertEquals(Field.Mode.REQUIRED, bqSchema.getFields().get(1).getMode()); assertEquals("bool_col", bqSchema.getFields().get(2).getName()); assertEquals(LegacySQLTypeName.BOOLEAN, bqSchema.getFields().get(2).getType()); @@ -128,9 +131,8 @@ public void testDeserializeRecordBatchPrimitives() throws IOException { ImmutableList.of(intVector, nameVector, scoreVector, activeVector, bytesVector, tsVector); try (VectorSchemaRoot root = new VectorSchemaRoot(vectors)) { - Schema arrowSchema = root.getSchema(); - com.google.cloud.bigquery.Schema bqSchema = - ArrowDeserializer.arrowSchemaToBigQuerySchema(arrowSchema); + Object arrowSchema = root.getSchema(); + Schema bqSchema = ArrowDeserializer.arrowSchemaToBigQuerySchema(arrowSchema); byte[] recordBatchBytes = serializeVectorSchemaRoot(root, allocator); @@ -155,11 +157,7 @@ public void testDeserializeRecordBatchPrimitives() throws IOException { assertEquals("102", row1.get("id").getStringValue()); assertEquals("Bob", row1.get("name").getStringValue()); assertNull(row1.get("score").getValue()); - assertEquals( - "false", - row1.get("false".equals("false") ? "active" : "score") != null - ? row1.get("active").getStringValue() - : "false"); + assertEquals("false", row1.get("active").getStringValue()); assertNull(row1.get("data").getValue()); assertNull(row1.get("ts").getValue()); } finally { @@ -179,10 +177,10 @@ public void testSchemaMismatchThrowsException() { intVector.setValueCount(1); try (VectorSchemaRoot root = new VectorSchemaRoot(ImmutableList.of(intVector))) { - com.google.cloud.bigquery.Schema mismatchedSchema = - com.google.cloud.bigquery.Schema.of( - com.google.cloud.bigquery.Field.of("col1", LegacySQLTypeName.INTEGER), - com.google.cloud.bigquery.Field.of("col2", LegacySQLTypeName.STRING)); + Schema mismatchedSchema = + Schema.of( + Field.of("col1", LegacySQLTypeName.INTEGER), + Field.of("col2", LegacySQLTypeName.STRING)); try { ArrowDeserializer.arrowRootToFieldValueList(root, 0, mismatchedSchema); From 88530f72e9c9454c31e222a8bf15f6968a15c55f Mon Sep 17 00:00:00 2001 From: Jin Seop Kim Date: Fri, 7 Aug 2026 16:15:19 -0400 Subject: [PATCH 7/7] refactor(bigquery): eliminate FQCNs in ArrowSerializationOptionsConverter and BigQueryImpl --- .../src/main/java/com/google/cloud/bigquery/BigQueryImpl.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 cc4ca5412e9f..d40724363f6c 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 @@ -45,7 +45,6 @@ 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.ArrowRecordBatch; import com.google.cloud.bigquery.storage.v1.BigQueryReadClient; import com.google.cloud.bigquery.storage.v1.BigQueryReadSettings; import com.google.cloud.bigquery.storage.v1.ReadRowsRequest; @@ -357,7 +356,8 @@ public Page getNextPage() { while (rowBatch.size() < pageSize && streamIterator.hasNext()) { ReadRowsResponse response = streamIterator.next(); if (response.hasArrowRecordBatch()) { - ArrowRecordBatch batch = response.getArrowRecordBatch(); + com.google.cloud.bigquery.storage.v1.ArrowRecordBatch batch = + response.getArrowRecordBatch(); ArrowRecordBatch deserializedBatch = MessageSerializer.deserializeRecordBatch( new ReadChannel(