From baa875505e979432308d68dae894f2efda21f01e Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Mon, 6 Jul 2026 13:28:24 -0400 Subject: [PATCH 01/11] feat(bigquery-jdbc): optimize Standard API performance with streaming row parser --- .../jdbc/BigQueryFieldValueListWrapper.java | 45 +++- .../bigquery/jdbc/BigQueryJsonResultSet.java | 2 +- .../jdbc/BigQueryJsonStreamParser.java | 248 ++++++++++++++++++ .../bigquery/jdbc/BigQueryStatement.java | 4 +- .../jdbc/BigQueryJsonStreamParserTest.java | 225 ++++++++++++++++ 5 files changed, 518 insertions(+), 6 deletions(-) create mode 100644 java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java create mode 100644 java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java index 39740e021773..c21561ab7297 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java @@ -18,12 +18,13 @@ import com.google.cloud.bigquery.FieldList; import com.google.cloud.bigquery.FieldValue; +import com.google.cloud.bigquery.FieldValue.Attribute; import com.google.cloud.bigquery.FieldValueList; import java.util.List; /** * Package-private, This class acts as a facade layer and wraps the FieldList(schema) and - * FieldValueList + * FieldValueList or lightweight Object[] row buffers. */ class BigQueryFieldValueListWrapper { @@ -37,6 +38,9 @@ class BigQueryFieldValueListWrapper { // reference as a List in case of an Array private final List arrayFieldValueList; + // Lightweight row values buffer (Object[]) for streaming parsing + private final Object[] rowValues; + // This flag marks the end of the stream for the ResultSet private boolean isLast = false; private final Exception exception; @@ -44,29 +48,38 @@ class BigQueryFieldValueListWrapper { static BigQueryFieldValueListWrapper of( FieldList fieldList, FieldValueList fieldValueList, boolean... isLast) { boolean isLastFlag = isLast != null && isLast.length == 1 && isLast[0]; - return new BigQueryFieldValueListWrapper(fieldList, fieldValueList, null, isLastFlag, null); + return new BigQueryFieldValueListWrapper( + fieldList, fieldValueList, null, null, isLastFlag, null); + } + + static BigQueryFieldValueListWrapper ofRow( + FieldList fieldList, Object[] rowValues, boolean... isLast) { + boolean isLastFlag = isLast != null && isLast.length == 1 && isLast[0]; + return new BigQueryFieldValueListWrapper(fieldList, null, null, rowValues, isLastFlag, null); } static BigQueryFieldValueListWrapper getNestedFieldValueListWrapper( FieldList fieldList, List arrayFieldValueList, boolean... isLast) { boolean isLastFlag = isLast != null && isLast.length == 1 && isLast[0]; return new BigQueryFieldValueListWrapper( - fieldList, null, arrayFieldValueList, isLastFlag, null); + fieldList, null, arrayFieldValueList, null, isLastFlag, null); } static BigQueryFieldValueListWrapper ofError(Exception exception) { - return new BigQueryFieldValueListWrapper(null, null, null, true, exception); + return new BigQueryFieldValueListWrapper(null, null, null, null, true, exception); } private BigQueryFieldValueListWrapper( FieldList fieldList, FieldValueList fieldValueList, List arrayFieldValueList, + Object[] rowValues, boolean isLast, Exception exception) { this.fieldList = fieldList; this.fieldValueList = fieldValueList; this.arrayFieldValueList = arrayFieldValueList; + this.rowValues = rowValues; this.isLast = isLast; this.exception = exception; } @@ -83,6 +96,30 @@ public List getArrayFieldValueList() { return this.arrayFieldValueList; } + public Object[] getRowValues() { + return this.rowValues; + } + + public FieldValue get(int index) { + if (this.fieldValueList != null) { + return this.fieldValueList.get(index); + } + if (this.arrayFieldValueList != null) { + return this.arrayFieldValueList.get(index); + } + if (this.rowValues != null && index >= 0 && index < this.rowValues.length) { + Object val = this.rowValues[index]; + if (val == null) { + return FieldValue.of(Attribute.PRIMITIVE, null); + } + if (val instanceof FieldValue) { + return (FieldValue) val; + } + return FieldValue.of(Attribute.PRIMITIVE, val.toString()); + } + return null; + } + public boolean isLast() { return this.isLast; } diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonResultSet.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonResultSet.java index 998a189eae2b..0dbda843d1e1 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonResultSet.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonResultSet.java @@ -299,7 +299,7 @@ private FieldValue getObjectInternal(int columnIndex) throws SQLException { // non nested, return the value else { // SQL Index to 0 based index - value = this.cursor.getFieldValueList().get(columnIndex - 1); + value = this.cursor.get(columnIndex - 1); } setWasNull(value.getValue()); return value; diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java new file mode 100644 index 000000000000..9cd11f57fd5e --- /dev/null +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java @@ -0,0 +1,248 @@ +/* + * 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. + */ + +package com.google.cloud.bigquery.jdbc; + +import static com.google.cloud.bigquery.jdbc.BigQueryBaseArray.isArray; +import static com.google.cloud.bigquery.jdbc.BigQueryBaseStruct.isStruct; + +import com.fasterxml.jackson.core.JsonFactory; +import com.fasterxml.jackson.core.JsonParser; +import com.fasterxml.jackson.core.JsonToken; +import com.google.api.core.InternalApi; +import com.google.cloud.bigquery.Field; +import com.google.cloud.bigquery.FieldList; +import com.google.cloud.bigquery.FieldValue; +import com.google.cloud.bigquery.FieldValue.Attribute; +import com.google.cloud.bigquery.FieldValueList; +import com.google.cloud.bigquery.Schema; +import java.io.IOException; +import java.io.InputStream; +import java.util.ArrayList; +import java.util.List; + +/** + * Package-private streaming parser for BigQuery REST JSON responses. + * + *

This class extracts cell data into compact primitive row arrays ({@code Object[]}), bypassing + * intermediate {@link FieldValueList} / {@link FieldValue} POJO allocations for primitive scalar + * columns to drastically reduce JVM heap overhead and GC pause times. + */ +@InternalApi +class BigQueryJsonStreamParser { + private static final JsonFactory JSON_FACTORY = new JsonFactory(); + + private final FieldList fieldList; + private final boolean[] isComplexColumn; + + BigQueryJsonStreamParser(Schema schema) { + this.fieldList = schema == null ? null : schema.getFields(); + if (this.fieldList != null) { + int size = this.fieldList.size(); + this.isComplexColumn = new boolean[size]; + for (int i = 0; i < size; i++) { + Field field = this.fieldList.get(i); + this.isComplexColumn[i] = isArray(field) || isStruct(field); + } + } else { + this.isComplexColumn = new boolean[0]; + } + } + + BigQueryJsonStreamParser(FieldList fieldList) { + this.fieldList = fieldList; + if (this.fieldList != null) { + int size = this.fieldList.size(); + this.isComplexColumn = new boolean[size]; + for (int i = 0; i < size; i++) { + Field field = this.fieldList.get(i); + this.isComplexColumn[i] = isArray(field) || isStruct(field); + } + } else { + this.isComplexColumn = new boolean[0]; + } + } + + /** + * Unpacks a {@link FieldValueList} row into a lightweight {@code Object[]} array, extracting raw + * string/primitive values and immediately discarding {@link FieldValueList} references to reduce + * heap memory pressure. + */ + public Object[] unpackRow(FieldValueList fieldValueList) { + if (fieldValueList == null) { + return null; + } + int size = fieldValueList.size(); + Object[] row = new Object[size]; + for (int i = 0; i < size; i++) { + FieldValue fv = fieldValueList.get(i); + if (fv == null || fv.isNull()) { + row[i] = null; + } else if (i < isComplexColumn.length && isComplexColumn[i]) { + // Retain FieldValue wrapper for complex ARRAY / STRUCT types + row[i] = fv; + } else { + // Extract raw scalar string representation for primitives + Object val = fv.getValue(); + row[i] = val == null ? null : val.toString(); + } + } + return row; + } + + /** Parses a raw JSON InputStream returning a list of row arrays ({@code Object[]}). */ + public List parseStream(InputStream inputStream) throws IOException { + List rows = new ArrayList<>(); + if (inputStream == null || fieldList == null) { + return rows; + } + + try (JsonParser parser = JSON_FACTORY.createParser(inputStream)) { + while (!parser.isClosed()) { + JsonToken token = parser.nextToken(); + if (token == JsonToken.FIELD_NAME && "rows".equals(parser.currentName())) { + // Found "rows" array + if (parser.nextToken() == JsonToken.START_ARRAY) { + while (parser.nextToken() != JsonToken.END_ARRAY && !parser.isClosed()) { + Object[] row = parseSingleRow(parser); + if (row != null) { + rows.add(row); + } + } + } + break; + } + } + } + return rows; + } + + private Object[] parseSingleRow(JsonParser parser) throws IOException { + // Expecting START_OBJECT for row: { "f": [...] } + if (parser.currentToken() != JsonToken.START_OBJECT) { + return null; + } + + int colCount = fieldList.size(); + Object[] row = new Object[colCount]; + + while (parser.nextToken() != JsonToken.END_OBJECT && !parser.isClosed()) { + String fieldName = parser.currentName(); + if ("f".equals(fieldName)) { + if (parser.nextToken() == JsonToken.START_ARRAY) { + int colIdx = 0; + while (parser.nextToken() != JsonToken.END_ARRAY && !parser.isClosed()) { + if (colIdx < colCount) { + row[colIdx] = parseCell(parser, colIdx); + } else { + skipValue(parser); + } + colIdx++; + } + } + } else { + skipValue(parser); + } + } + return row; + } + + private Object parseCell(JsonParser parser, int colIdx) throws IOException { + // Cell structure: { "v": ... } + if (parser.currentToken() != JsonToken.START_OBJECT) { + return null; + } + + Object cellValue = null; + while (parser.nextToken() != JsonToken.END_OBJECT && !parser.isClosed()) { + String name = parser.currentName(); + if ("v".equals(name)) { + parser.nextToken(); // move to value token + if (parser.currentToken() == JsonToken.VALUE_NULL) { + cellValue = null; + } else if (isComplexColumn[colIdx]) { + // Complex type (ARRAY or STRUCT) fallback to FieldValue + cellValue = parseComplexFieldValue(parser, fieldList.get(colIdx)); + } else { + // Primitive scalar token + cellValue = parser.getText(); + } + } else { + skipValue(parser); + } + } + return cellValue; + } + + private FieldValue parseComplexFieldValue(JsonParser parser, Field field) throws IOException { + if (parser.currentToken() == JsonToken.VALUE_NULL) { + return FieldValue.of(Attribute.PRIMITIVE, null); + } + if (isArray(field)) { + List elements = new ArrayList<>(); + if (parser.currentToken() == JsonToken.START_ARRAY) { + Field elementField = field.toBuilder().setMode(Field.Mode.REQUIRED).build(); + while (parser.nextToken() != JsonToken.END_ARRAY && !parser.isClosed()) { + // element is { "v": ... } + elements.add(parseComplexFieldValue(parser, elementField)); + } + } + return FieldValue.of(Attribute.REPEATED, elements); + } else if (isStruct(field)) { + List fields = new ArrayList<>(); + if (parser.currentToken() == JsonToken.START_OBJECT) { + FieldList subFields = field.getSubFields(); + int subIdx = 0; + while (parser.nextToken() != JsonToken.END_OBJECT && !parser.isClosed()) { + if ("f".equals(parser.currentName()) && parser.nextToken() == JsonToken.START_ARRAY) { + while (parser.nextToken() != JsonToken.END_ARRAY && !parser.isClosed()) { + Field subField = + subFields != null && subIdx < subFields.size() ? subFields.get(subIdx) : null; + fields.add( + subField != null + ? parseComplexFieldValue(parser, subField) + : FieldValue.of(Attribute.PRIMITIVE, null)); + subIdx++; + } + } + } + } + return FieldValue.of(Attribute.RECORD, FieldValueList.of(fields, field.getSubFields())); + } else { + if (parser.currentToken() == JsonToken.START_OBJECT) { + String val = null; + while (parser.nextToken() != JsonToken.END_OBJECT && !parser.isClosed()) { + if ("v".equals(parser.currentName())) { + parser.nextToken(); + val = parser.currentToken() == JsonToken.VALUE_NULL ? null : parser.getText(); + } else { + skipValue(parser); + } + } + return FieldValue.of(Attribute.PRIMITIVE, val); + } else { + return FieldValue.of(Attribute.PRIMITIVE, parser.getText()); + } + } + } + + private void skipValue(JsonParser parser) throws IOException { + JsonToken token = parser.currentToken(); + if (token == JsonToken.START_OBJECT || token == JsonToken.START_ARRAY) { + parser.skipChildren(); + } + } +} diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java index 6f8a5d71deb0..78340fd2a2ed 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java @@ -1280,15 +1280,17 @@ Future parseAndPopulateRpcDataAsync( long startTime = System.nanoTime(); long results = 0; + BigQueryJsonStreamParser streamParser = new BigQueryJsonStreamParser(schema); for (FieldValueList fieldValueList : fieldValueLists) { if (Thread.currentThread().isInterrupted() || executor.isShutdown()) { // do not process further pages and shutdown (inner loop) break; } + Object[] rowArray = streamParser.unpackRow(fieldValueList); Uninterruptibles.putUninterruptibly( bigQueryFieldValueListWrapperBlockingQueue, - BigQueryFieldValueListWrapper.of(schema.getFields(), fieldValueList)); + BigQueryFieldValueListWrapper.ofRow(schema.getFields(), rowArray)); results += 1; } LOG.fine( diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java new file mode 100644 index 000000000000..43f9c45982be --- /dev/null +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java @@ -0,0 +1,225 @@ +/* + * 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. + */ + +package com.google.cloud.bigquery.jdbc; + +import static com.google.common.truth.Truth.assertThat; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; + +import com.google.cloud.bigquery.Field; +import com.google.cloud.bigquery.FieldList; +import com.google.cloud.bigquery.FieldValue; +import com.google.cloud.bigquery.FieldValue.Attribute; +import com.google.cloud.bigquery.FieldValueList; +import com.google.cloud.bigquery.Schema; +import com.google.cloud.bigquery.StandardSQLTypeName; +import com.google.common.collect.ImmutableList; +import com.sun.management.ThreadMXBean; +import java.io.ByteArrayInputStream; +import java.io.InputStream; +import java.lang.management.ManagementFactory; +import java.math.BigDecimal; +import java.nio.charset.StandardCharsets; +import java.sql.Array; +import java.sql.Date; +import java.sql.Struct; +import java.sql.Time; +import java.util.List; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.Future; +import java.util.concurrent.LinkedBlockingDeque; +import org.junit.jupiter.api.Test; + +public class BigQueryJsonStreamParserTest { + + @Test + public void testStreamParsingAndCoercionAllTypes() throws Exception { + FieldList profileSchema = + FieldList.of( + Field.of("pname", StandardSQLTypeName.STRING), + Field.of("page", StandardSQLTypeName.INT64)); + + FieldList innerStructSchema = + FieldList.of( + Field.of("city", StandardSQLTypeName.STRING), + Field.of("country", StandardSQLTypeName.STRING)); + + FieldList outerStructSchema = + FieldList.of( + Field.of("company", StandardSQLTypeName.STRING), + Field.of("location", StandardSQLTypeName.STRUCT, innerStructSchema)); + + FieldList schemaFields = + FieldList.of( + Field.of("id", StandardSQLTypeName.INT64), + Field.of("name", StandardSQLTypeName.STRING), + Field.of("score", StandardSQLTypeName.FLOAT64), + Field.of("active", StandardSQLTypeName.BOOL), + Field.of("geo", StandardSQLTypeName.GEOGRAPHY), + Field.of("json_col", StandardSQLTypeName.JSON), + Field.of("num", StandardSQLTypeName.NUMERIC), + Field.of("bignum", StandardSQLTypeName.BIGNUMERIC), + Field.of("interval_col", StandardSQLTypeName.INTERVAL), + Field.of("bytes_col", StandardSQLTypeName.BYTES), + Field.of("ts", StandardSQLTypeName.TIMESTAMP), + Field.of("dt", StandardSQLTypeName.DATE), + Field.of("tm", StandardSQLTypeName.TIME), + Field.newBuilder("tags", StandardSQLTypeName.STRING) + .setMode(Field.Mode.REPEATED) + .build(), + Field.newBuilder("profiles", StandardSQLTypeName.STRUCT, profileSchema) + .setMode(Field.Mode.REPEATED) + .build(), + Field.of("org", StandardSQLTypeName.STRUCT, outerStructSchema)); + + Schema schema = Schema.of(schemaFields); + BigQueryJsonStreamParser parser = new BigQueryJsonStreamParser(schema); + + String json = + "{\n" + + " \"rows\": [\n" + + " {\n" + + " \"f\": [\n" + + " { \"v\": \"101\" },\n" + + " { \"v\": \"Alice\" },\n" + + " { \"v\": \"98.5\" },\n" + + " { \"v\": \"true\" },\n" + + " { \"v\": \"POINT(-122.084 37.422)\" },\n" + + " { \"v\": \"{\\\"key\\\": \\\"value\\\"}\" },\n" + + " { \"v\": \"123456789.987654321\" },\n" + + " { \"v\": \"99999999999999999999999999999999999999.999999999\" },\n" + + " { \"v\": \"0-0 0 0:0:0\" },\n" + + " { \"v\": \"SGVsbG8gV29ybGQ=\" },\n" + + " { \"v\": \"1408452095.22\" },\n" + + " { \"v\": \"2023-03-13\" },\n" + + " { \"v\": \"23:59:59\" },\n" + + " { \"v\": [ { \"v\": \"tag1\" }, { \"v\": \"tag2\" } ] },\n" + + " { \"v\": [ { \"v\": { \"f\": [ { \"v\": \"Bob\" }, { \"v\": \"30\" } ] } } ] },\n" + + " { \"v\": { \"f\": [ { \"v\": \"Acme Corp\" }, { \"v\": { \"f\": [ { \"v\": \"London\" }, { \"v\": \"UK\" } ] } } ] } }\n" + + " ]\n" + + " }\n" + + " ]\n" + + "}"; + + InputStream stream = new ByteArrayInputStream(json.getBytes(StandardCharsets.UTF_8)); + List rows = parser.parseStream(stream); + + assertThat(rows).hasSize(1); + Object[] row = rows.get(0); + assertThat(row[0]).isEqualTo("101"); + assertThat(row[1]).isEqualTo("Alice"); + assertThat(row[2]).isEqualTo("98.5"); + assertThat(row[3]).isEqualTo("true"); + assertThat(row[4]).isEqualTo("POINT(-122.084 37.422)"); + assertThat(row[5]).isEqualTo("{\"key\": \"value\"}"); + assertThat(row[6]).isEqualTo("123456789.987654321"); + assertThat(row[7]).isEqualTo("99999999999999999999999999999999999999.999999999"); + assertThat(row[8]).isEqualTo("0-0 0 0:0:0"); + assertThat(row[9]).isEqualTo("SGVsbG8gV29ybGQ="); + assertThat(row[10]).isEqualTo("1408452095.22"); + assertThat(row[11]).isEqualTo("2023-03-13"); + assertThat(row[12]).isEqualTo("23:59:59"); + assertThat(row[13]).isInstanceOf(FieldValue.class); + assertThat(row[14]).isInstanceOf(FieldValue.class); + assertThat(row[15]).isInstanceOf(FieldValue.class); + + // Verify ResultSet end-to-end JDBC Coercion + BlockingQueue queue = new LinkedBlockingDeque<>(); + queue.add(BigQueryFieldValueListWrapper.ofRow(schemaFields, row)); + queue.add(BigQueryFieldValueListWrapper.ofRow(null, null, true)); + BigQueryStatement statement = mock(BigQueryStatement.class); + Future[] workerTasks = {mock(Future.class)}; + + BigQueryJsonResultSet rs = BigQueryJsonResultSet.of(schema, 1L, queue, statement, workerTasks); + assertTrue(rs.next()); + + assertThat(rs.getLong("id")).isEqualTo(101L); + assertThat(rs.getString("name")).isEqualTo("Alice"); + assertThat(rs.getDouble("score")).isEqualTo(98.5D); + assertThat(rs.getBoolean("active")).isTrue(); + assertThat(rs.getString("geo")).isEqualTo("POINT(-122.084 37.422)"); + assertThat(rs.getString("json_col")).isEqualTo("{\"key\": \"value\"}"); + assertThat(rs.getBigDecimal("num")).isEqualTo(new BigDecimal("123456789.987654321")); + assertThat(rs.getBigDecimal("bignum")) + .isEqualTo(new BigDecimal("99999999999999999999999999999999999999.999999999")); + assertThat(rs.getString("interval_col")).isEqualTo("0-0 0 0:0:0"); + assertThat(rs.getBytes("bytes_col")).isEqualTo("Hello World".getBytes(StandardCharsets.UTF_8)); + assertThat(rs.getTimestamp("ts")).isNotNull(); + assertThat(rs.getDate("dt")).isEqualTo(Date.valueOf("2023-03-13")); + assertThat(rs.getTime("tm")).isEqualTo(Time.valueOf("23:59:59")); + + // Array of Primitives + Array tagsArray = rs.getArray("tags"); + assertThat((String[]) tagsArray.getArray()).isEqualTo(new String[] {"tag1", "tag2"}); + + // Array of Structs + Array profilesArray = rs.getArray("profiles"); + Object[] profileStructs = (Object[]) profilesArray.getArray(); + assertThat(profileStructs.length).isEqualTo(1); + assertThat(((Struct) profileStructs[0]).getAttributes()).isEqualTo(new Object[] {"Bob", 30L}); + + // Nested Structs + Struct orgStruct = (Struct) rs.getObject("org"); + Object[] orgAttributes = orgStruct.getAttributes(); + assertThat(orgAttributes[0]).isEqualTo("Acme Corp"); + assertThat(((Struct) orgAttributes[1]).getAttributes()) + .isEqualTo(new Object[] {"London", "UK"}); + } + + @Test + public void testMemoryAllocationReduction() throws Exception { + java.lang.management.ThreadMXBean baseBean = ManagementFactory.getThreadMXBean(); + if (!(baseBean instanceof ThreadMXBean)) { + return; + } + ThreadMXBean threadBean = (ThreadMXBean) baseBean; + long threadId = Thread.currentThread().getId(); + + FieldList simpleFieldList = + FieldList.of( + Field.of("col1", StandardSQLTypeName.INT64), + Field.of("col2", StandardSQLTypeName.STRING), + Field.of("col3", StandardSQLTypeName.FLOAT64), + Field.of("col4", StandardSQLTypeName.BOOL), + Field.of("col5", StandardSQLTypeName.TIMESTAMP)); + + Schema simpleSchema = Schema.of(simpleFieldList); + BigQueryJsonStreamParser parser = new BigQueryJsonStreamParser(simpleSchema); + + int rowCount = 10000; + FieldValue fv1 = FieldValue.of(Attribute.PRIMITIVE, "100"); + FieldValue fv2 = FieldValue.of(Attribute.PRIMITIVE, "test_string"); + FieldValue fv3 = FieldValue.of(Attribute.PRIMITIVE, "123.456"); + FieldValue fv4 = FieldValue.of(Attribute.PRIMITIVE, "true"); + FieldValue fv5 = FieldValue.of(Attribute.PRIMITIVE, "1680174859.820000"); + + List fvItems = ImmutableList.of(fv1, fv2, fv3, fv4, fv5); + FieldValueList sampleFvl = FieldValueList.of(fvItems, simpleFieldList); + + long bytesBefore = threadBean.getThreadAllocatedBytes(threadId); + Object[][] rowBuffers = new Object[rowCount][]; + for (int i = 0; i < rowCount; i++) { + rowBuffers[i] = parser.unpackRow(sampleFvl); + } + long allocatedBytes = threadBean.getThreadAllocatedBytes(threadId) - bytesBefore; + + assertThat(rowBuffers.length).isEqualTo(rowCount); + assertTrue( + allocatedBytes < 5 * 1024 * 1024, + "Allocated bytes " + allocatedBytes + " exceeded threshold"); + } +} From 9a32437c446c133980c0d547c396297a22934425 Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Mon, 6 Jul 2026 13:30:26 -0400 Subject: [PATCH 02/11] lint --- .../cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java index 43f9c45982be..759ff21b5021 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java @@ -182,11 +182,11 @@ public void testStreamParsingAndCoercionAllTypes() throws Exception { @Test public void testMemoryAllocationReduction() throws Exception { - java.lang.management.ThreadMXBean baseBean = ManagementFactory.getThreadMXBean(); - if (!(baseBean instanceof ThreadMXBean)) { + ThreadMXBean baseBean = ManagementFactory.getThreadMXBean(); + if (!(baseBean instanceof com.sun.management.ThreadMXBean)) { return; } - ThreadMXBean threadBean = (ThreadMXBean) baseBean; + com.sun.management.ThreadMXBean threadBean = (com.sun.management.ThreadMXBean) baseBean; long threadId = Thread.currentThread().getId(); FieldList simpleFieldList = From e0392b0bd2232b4ed2c280478eadddb5cfd582fa Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Mon, 6 Jul 2026 14:51:19 -0400 Subject: [PATCH 03/11] check token --- .../jdbc/BigQueryJsonStreamParser.java | 73 ++++++++++++++----- .../jdbc/BigQueryJsonStreamParserTest.java | 6 +- 2 files changed, 57 insertions(+), 22 deletions(-) diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java index 9cd11f57fd5e..a6c0a2fcb9bd 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java @@ -111,12 +111,15 @@ public List parseStream(InputStream inputStream) throws IOException { } try (JsonParser parser = JSON_FACTORY.createParser(inputStream)) { - while (!parser.isClosed()) { - JsonToken token = parser.nextToken(); + JsonToken token; + while ((token = parser.nextToken()) != null && !parser.isClosed()) { if (token == JsonToken.FIELD_NAME && "rows".equals(parser.currentName())) { // Found "rows" array if (parser.nextToken() == JsonToken.START_ARRAY) { - while (parser.nextToken() != JsonToken.END_ARRAY && !parser.isClosed()) { + JsonToken rowToken; + while ((rowToken = parser.nextToken()) != null + && rowToken != JsonToken.END_ARRAY + && !parser.isClosed()) { Object[] row = parseSingleRow(parser); if (row != null) { rows.add(row); @@ -139,12 +142,18 @@ private Object[] parseSingleRow(JsonParser parser) throws IOException { int colCount = fieldList.size(); Object[] row = new Object[colCount]; - while (parser.nextToken() != JsonToken.END_OBJECT && !parser.isClosed()) { + JsonToken token; + while ((token = parser.nextToken()) != null + && token != JsonToken.END_OBJECT + && !parser.isClosed()) { String fieldName = parser.currentName(); if ("f".equals(fieldName)) { if (parser.nextToken() == JsonToken.START_ARRAY) { int colIdx = 0; - while (parser.nextToken() != JsonToken.END_ARRAY && !parser.isClosed()) { + JsonToken cellToken; + while ((cellToken = parser.nextToken()) != null + && cellToken != JsonToken.END_ARRAY + && !parser.isClosed()) { if (colIdx < colCount) { row[colIdx] = parseCell(parser, colIdx); } else { @@ -167,7 +176,10 @@ private Object parseCell(JsonParser parser, int colIdx) throws IOException { } Object cellValue = null; - while (parser.nextToken() != JsonToken.END_OBJECT && !parser.isClosed()) { + JsonToken token; + while ((token = parser.nextToken()) != null + && token != JsonToken.END_OBJECT + && !parser.isClosed()) { String name = parser.currentName(); if ("v".equals(name)) { parser.nextToken(); // move to value token @@ -195,7 +207,10 @@ private FieldValue parseComplexFieldValue(JsonParser parser, Field field) throws List elements = new ArrayList<>(); if (parser.currentToken() == JsonToken.START_ARRAY) { Field elementField = field.toBuilder().setMode(Field.Mode.REQUIRED).build(); - while (parser.nextToken() != JsonToken.END_ARRAY && !parser.isClosed()) { + JsonToken elemToken; + while ((elemToken = parser.nextToken()) != null + && elemToken != JsonToken.END_ARRAY + && !parser.isClosed()) { // element is { "v": ... } elements.add(parseComplexFieldValue(parser, elementField)); } @@ -206,25 +221,45 @@ private FieldValue parseComplexFieldValue(JsonParser parser, Field field) throws if (parser.currentToken() == JsonToken.START_OBJECT) { FieldList subFields = field.getSubFields(); int subIdx = 0; - while (parser.nextToken() != JsonToken.END_OBJECT && !parser.isClosed()) { - if ("f".equals(parser.currentName()) && parser.nextToken() == JsonToken.START_ARRAY) { - while (parser.nextToken() != JsonToken.END_ARRAY && !parser.isClosed()) { - Field subField = - subFields != null && subIdx < subFields.size() ? subFields.get(subIdx) : null; - fields.add( - subField != null - ? parseComplexFieldValue(parser, subField) - : FieldValue.of(Attribute.PRIMITIVE, null)); - subIdx++; + int depth = 1; + while (depth > 0 && !parser.isClosed()) { + JsonToken token = parser.nextToken(); + if (token == null) { + break; + } + if (token == JsonToken.START_OBJECT) { + depth++; + } else if (token == JsonToken.END_OBJECT) { + depth--; + } else if (token == JsonToken.FIELD_NAME) { + if ("f".equals(parser.currentName()) && parser.nextToken() == JsonToken.START_ARRAY) { + JsonToken subToken; + while ((subToken = parser.nextToken()) != null + && subToken != JsonToken.END_ARRAY + && !parser.isClosed()) { + Field subField = + subFields != null && subIdx < subFields.size() ? subFields.get(subIdx) : null; + fields.add( + subField != null + ? parseComplexFieldValue(parser, subField) + : FieldValue.of(Attribute.PRIMITIVE, null)); + subIdx++; + } } } } + if (subFields != null && fields.size() == subFields.size()) { + return FieldValue.of(Attribute.RECORD, FieldValueList.of(fields, subFields)); + } } - return FieldValue.of(Attribute.RECORD, FieldValueList.of(fields, field.getSubFields())); + return FieldValue.of(Attribute.RECORD, FieldValueList.of(fields)); } else { if (parser.currentToken() == JsonToken.START_OBJECT) { String val = null; - while (parser.nextToken() != JsonToken.END_OBJECT && !parser.isClosed()) { + JsonToken token; + while ((token = parser.nextToken()) != null + && token != JsonToken.END_OBJECT + && !parser.isClosed()) { if ("v".equals(parser.currentName())) { parser.nextToken(); val = parser.currentToken() == JsonToken.VALUE_NULL ? null : parser.getText(); diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java index 759ff21b5021..43f9c45982be 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java @@ -182,11 +182,11 @@ public void testStreamParsingAndCoercionAllTypes() throws Exception { @Test public void testMemoryAllocationReduction() throws Exception { - ThreadMXBean baseBean = ManagementFactory.getThreadMXBean(); - if (!(baseBean instanceof com.sun.management.ThreadMXBean)) { + java.lang.management.ThreadMXBean baseBean = ManagementFactory.getThreadMXBean(); + if (!(baseBean instanceof ThreadMXBean)) { return; } - com.sun.management.ThreadMXBean threadBean = (com.sun.management.ThreadMXBean) baseBean; + ThreadMXBean threadBean = (ThreadMXBean) baseBean; long threadId = Thread.currentThread().getId(); FieldList simpleFieldList = From c0809b25d3c17f8bb2a33619928865564bbb06d4 Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Mon, 6 Jul 2026 15:13:47 -0400 Subject: [PATCH 04/11] remove parseStream --- .../jdbc/BigQueryJsonStreamParser.java | 190 +----------------- .../jdbc/BigQueryJsonStreamParserTest.java | 85 ++++---- 2 files changed, 52 insertions(+), 223 deletions(-) diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java index a6c0a2fcb9bd..183ce4a877bb 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java @@ -19,23 +19,15 @@ import static com.google.cloud.bigquery.jdbc.BigQueryBaseArray.isArray; import static com.google.cloud.bigquery.jdbc.BigQueryBaseStruct.isStruct; -import com.fasterxml.jackson.core.JsonFactory; -import com.fasterxml.jackson.core.JsonParser; -import com.fasterxml.jackson.core.JsonToken; import com.google.api.core.InternalApi; import com.google.cloud.bigquery.Field; import com.google.cloud.bigquery.FieldList; import com.google.cloud.bigquery.FieldValue; -import com.google.cloud.bigquery.FieldValue.Attribute; import com.google.cloud.bigquery.FieldValueList; import com.google.cloud.bigquery.Schema; -import java.io.IOException; -import java.io.InputStream; -import java.util.ArrayList; -import java.util.List; /** - * Package-private streaming parser for BigQuery REST JSON responses. + * Package-private parser for BigQuery rows into lightweight Object[] buffers. * *

This class extracts cell data into compact primitive row arrays ({@code Object[]}), bypassing * intermediate {@link FieldValueList} / {@link FieldValue} POJO allocations for primitive scalar @@ -43,8 +35,6 @@ */ @InternalApi class BigQueryJsonStreamParser { - private static final JsonFactory JSON_FACTORY = new JsonFactory(); - private final FieldList fieldList; private final boolean[] isComplexColumn; @@ -102,182 +92,4 @@ public Object[] unpackRow(FieldValueList fieldValueList) { } return row; } - - /** Parses a raw JSON InputStream returning a list of row arrays ({@code Object[]}). */ - public List parseStream(InputStream inputStream) throws IOException { - List rows = new ArrayList<>(); - if (inputStream == null || fieldList == null) { - return rows; - } - - try (JsonParser parser = JSON_FACTORY.createParser(inputStream)) { - JsonToken token; - while ((token = parser.nextToken()) != null && !parser.isClosed()) { - if (token == JsonToken.FIELD_NAME && "rows".equals(parser.currentName())) { - // Found "rows" array - if (parser.nextToken() == JsonToken.START_ARRAY) { - JsonToken rowToken; - while ((rowToken = parser.nextToken()) != null - && rowToken != JsonToken.END_ARRAY - && !parser.isClosed()) { - Object[] row = parseSingleRow(parser); - if (row != null) { - rows.add(row); - } - } - } - break; - } - } - } - return rows; - } - - private Object[] parseSingleRow(JsonParser parser) throws IOException { - // Expecting START_OBJECT for row: { "f": [...] } - if (parser.currentToken() != JsonToken.START_OBJECT) { - return null; - } - - int colCount = fieldList.size(); - Object[] row = new Object[colCount]; - - JsonToken token; - while ((token = parser.nextToken()) != null - && token != JsonToken.END_OBJECT - && !parser.isClosed()) { - String fieldName = parser.currentName(); - if ("f".equals(fieldName)) { - if (parser.nextToken() == JsonToken.START_ARRAY) { - int colIdx = 0; - JsonToken cellToken; - while ((cellToken = parser.nextToken()) != null - && cellToken != JsonToken.END_ARRAY - && !parser.isClosed()) { - if (colIdx < colCount) { - row[colIdx] = parseCell(parser, colIdx); - } else { - skipValue(parser); - } - colIdx++; - } - } - } else { - skipValue(parser); - } - } - return row; - } - - private Object parseCell(JsonParser parser, int colIdx) throws IOException { - // Cell structure: { "v": ... } - if (parser.currentToken() != JsonToken.START_OBJECT) { - return null; - } - - Object cellValue = null; - JsonToken token; - while ((token = parser.nextToken()) != null - && token != JsonToken.END_OBJECT - && !parser.isClosed()) { - String name = parser.currentName(); - if ("v".equals(name)) { - parser.nextToken(); // move to value token - if (parser.currentToken() == JsonToken.VALUE_NULL) { - cellValue = null; - } else if (isComplexColumn[colIdx]) { - // Complex type (ARRAY or STRUCT) fallback to FieldValue - cellValue = parseComplexFieldValue(parser, fieldList.get(colIdx)); - } else { - // Primitive scalar token - cellValue = parser.getText(); - } - } else { - skipValue(parser); - } - } - return cellValue; - } - - private FieldValue parseComplexFieldValue(JsonParser parser, Field field) throws IOException { - if (parser.currentToken() == JsonToken.VALUE_NULL) { - return FieldValue.of(Attribute.PRIMITIVE, null); - } - if (isArray(field)) { - List elements = new ArrayList<>(); - if (parser.currentToken() == JsonToken.START_ARRAY) { - Field elementField = field.toBuilder().setMode(Field.Mode.REQUIRED).build(); - JsonToken elemToken; - while ((elemToken = parser.nextToken()) != null - && elemToken != JsonToken.END_ARRAY - && !parser.isClosed()) { - // element is { "v": ... } - elements.add(parseComplexFieldValue(parser, elementField)); - } - } - return FieldValue.of(Attribute.REPEATED, elements); - } else if (isStruct(field)) { - List fields = new ArrayList<>(); - if (parser.currentToken() == JsonToken.START_OBJECT) { - FieldList subFields = field.getSubFields(); - int subIdx = 0; - int depth = 1; - while (depth > 0 && !parser.isClosed()) { - JsonToken token = parser.nextToken(); - if (token == null) { - break; - } - if (token == JsonToken.START_OBJECT) { - depth++; - } else if (token == JsonToken.END_OBJECT) { - depth--; - } else if (token == JsonToken.FIELD_NAME) { - if ("f".equals(parser.currentName()) && parser.nextToken() == JsonToken.START_ARRAY) { - JsonToken subToken; - while ((subToken = parser.nextToken()) != null - && subToken != JsonToken.END_ARRAY - && !parser.isClosed()) { - Field subField = - subFields != null && subIdx < subFields.size() ? subFields.get(subIdx) : null; - fields.add( - subField != null - ? parseComplexFieldValue(parser, subField) - : FieldValue.of(Attribute.PRIMITIVE, null)); - subIdx++; - } - } - } - } - if (subFields != null && fields.size() == subFields.size()) { - return FieldValue.of(Attribute.RECORD, FieldValueList.of(fields, subFields)); - } - } - return FieldValue.of(Attribute.RECORD, FieldValueList.of(fields)); - } else { - if (parser.currentToken() == JsonToken.START_OBJECT) { - String val = null; - JsonToken token; - while ((token = parser.nextToken()) != null - && token != JsonToken.END_OBJECT - && !parser.isClosed()) { - if ("v".equals(parser.currentName())) { - parser.nextToken(); - val = parser.currentToken() == JsonToken.VALUE_NULL ? null : parser.getText(); - } else { - skipValue(parser); - } - } - return FieldValue.of(Attribute.PRIMITIVE, val); - } else { - return FieldValue.of(Attribute.PRIMITIVE, parser.getText()); - } - } - } - - private void skipValue(JsonParser parser) throws IOException { - JsonToken token = parser.currentToken(); - if (token == JsonToken.START_OBJECT || token == JsonToken.START_ARRAY) { - parser.skipChildren(); - } - } } diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java index 43f9c45982be..1ba236a18eee 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java @@ -29,8 +29,6 @@ import com.google.cloud.bigquery.StandardSQLTypeName; import com.google.common.collect.ImmutableList; import com.sun.management.ThreadMXBean; -import java.io.ByteArrayInputStream; -import java.io.InputStream; import java.lang.management.ManagementFactory; import java.math.BigDecimal; import java.nio.charset.StandardCharsets; @@ -38,6 +36,7 @@ import java.sql.Date; import java.sql.Struct; import java.sql.Time; +import java.util.Arrays; import java.util.List; import java.util.concurrent.BlockingQueue; import java.util.concurrent.Future; @@ -47,7 +46,7 @@ public class BigQueryJsonStreamParserTest { @Test - public void testStreamParsingAndCoercionAllTypes() throws Exception { + public void testUnpackRowAndCoercionAllTypes() throws Exception { FieldList profileSchema = FieldList.of( Field.of("pname", StandardSQLTypeName.STRING), @@ -89,37 +88,55 @@ public void testStreamParsingAndCoercionAllTypes() throws Exception { Schema schema = Schema.of(schemaFields); BigQueryJsonStreamParser parser = new BigQueryJsonStreamParser(schema); - String json = - "{\n" - + " \"rows\": [\n" - + " {\n" - + " \"f\": [\n" - + " { \"v\": \"101\" },\n" - + " { \"v\": \"Alice\" },\n" - + " { \"v\": \"98.5\" },\n" - + " { \"v\": \"true\" },\n" - + " { \"v\": \"POINT(-122.084 37.422)\" },\n" - + " { \"v\": \"{\\\"key\\\": \\\"value\\\"}\" },\n" - + " { \"v\": \"123456789.987654321\" },\n" - + " { \"v\": \"99999999999999999999999999999999999999.999999999\" },\n" - + " { \"v\": \"0-0 0 0:0:0\" },\n" - + " { \"v\": \"SGVsbG8gV29ybGQ=\" },\n" - + " { \"v\": \"1408452095.22\" },\n" - + " { \"v\": \"2023-03-13\" },\n" - + " { \"v\": \"23:59:59\" },\n" - + " { \"v\": [ { \"v\": \"tag1\" }, { \"v\": \"tag2\" } ] },\n" - + " { \"v\": [ { \"v\": { \"f\": [ { \"v\": \"Bob\" }, { \"v\": \"30\" } ] } } ] },\n" - + " { \"v\": { \"f\": [ { \"v\": \"Acme Corp\" }, { \"v\": { \"f\": [ { \"v\": \"London\" }, { \"v\": \"UK\" } ] } } ] } }\n" - + " ]\n" - + " }\n" - + " ]\n" - + "}"; - - InputStream stream = new ByteArrayInputStream(json.getBytes(StandardCharsets.UTF_8)); - List rows = parser.parseStream(stream); - - assertThat(rows).hasSize(1); - Object[] row = rows.get(0); + FieldValueList fvl = + FieldValueList.of( + Arrays.asList( + FieldValue.of(Attribute.PRIMITIVE, "101"), + FieldValue.of(Attribute.PRIMITIVE, "Alice"), + FieldValue.of(Attribute.PRIMITIVE, "98.5"), + FieldValue.of(Attribute.PRIMITIVE, "true"), + FieldValue.of(Attribute.PRIMITIVE, "POINT(-122.084 37.422)"), + FieldValue.of(Attribute.PRIMITIVE, "{\"key\": \"value\"}"), + FieldValue.of(Attribute.PRIMITIVE, "123456789.987654321"), + FieldValue.of( + Attribute.PRIMITIVE, "99999999999999999999999999999999999999.999999999"), + FieldValue.of(Attribute.PRIMITIVE, "0-0 0 0:0:0"), + FieldValue.of(Attribute.PRIMITIVE, "SGVsbG8gV29ybGQ="), + FieldValue.of(Attribute.PRIMITIVE, "1408452095.22"), + FieldValue.of(Attribute.PRIMITIVE, "2023-03-13"), + FieldValue.of(Attribute.PRIMITIVE, "23:59:59"), + FieldValue.of( + Attribute.REPEATED, + Arrays.asList( + FieldValue.of(Attribute.PRIMITIVE, "tag1"), + FieldValue.of(Attribute.PRIMITIVE, "tag2"))), + FieldValue.of( + Attribute.REPEATED, + Arrays.asList( + FieldValue.of( + Attribute.RECORD, + FieldValueList.of( + Arrays.asList( + FieldValue.of(Attribute.PRIMITIVE, "Bob"), + FieldValue.of(Attribute.PRIMITIVE, "30")), + profileSchema)))), + FieldValue.of( + Attribute.RECORD, + FieldValueList.of( + Arrays.asList( + FieldValue.of(Attribute.PRIMITIVE, "Acme Corp"), + FieldValue.of( + Attribute.RECORD, + FieldValueList.of( + Arrays.asList( + FieldValue.of(Attribute.PRIMITIVE, "London"), + FieldValue.of(Attribute.PRIMITIVE, "UK")), + innerStructSchema))), + outerStructSchema))), + schemaFields); + + Object[] row = parser.unpackRow(fvl); + assertThat(row[0]).isEqualTo("101"); assertThat(row[1]).isEqualTo("Alice"); assertThat(row[2]).isEqualTo("98.5"); From 595564d383bdddb1c0b9a69b298ff120b856c06d Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Mon, 6 Jul 2026 17:16:32 -0400 Subject: [PATCH 05/11] clean up code --- .../jdbc/BigQueryFieldValueListWrapper.java | 42 ++++++++ .../jdbc/BigQueryJsonStreamParser.java | 95 ------------------- .../bigquery/jdbc/BigQueryStatement.java | 8 +- .../jdbc/BigQueryJsonStreamParserTest.java | 12 ++- 4 files changed, 54 insertions(+), 103 deletions(-) delete mode 100644 java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java index c21561ab7297..cc7cefcde87e 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java @@ -16,6 +16,10 @@ package com.google.cloud.bigquery.jdbc; +import static com.google.cloud.bigquery.jdbc.BigQueryBaseArray.isArray; +import static com.google.cloud.bigquery.jdbc.BigQueryBaseStruct.isStruct; + +import com.google.cloud.bigquery.Field; import com.google.cloud.bigquery.FieldList; import com.google.cloud.bigquery.FieldValue; import com.google.cloud.bigquery.FieldValue.Attribute; @@ -58,6 +62,44 @@ static BigQueryFieldValueListWrapper ofRow( return new BigQueryFieldValueListWrapper(fieldList, null, null, rowValues, isLastFlag, null); } + public static boolean[] createComplexColumnFlags(FieldList fieldList) { + if (fieldList == null) { + return new boolean[0]; + } + int size = fieldList.size(); + boolean[] isComplex = new boolean[size]; + for (int i = 0; i < size; i++) { + Field field = fieldList.get(i); + isComplex[i] = isArray(field) || isStruct(field); + } + return isComplex; + } + + public static Object[] unpackRow(FieldValueList fieldValueList, boolean[] isComplexColumn) { + if (fieldValueList == null) { + return null; + } + int size = fieldValueList.size(); + Object[] row = new Object[size]; + for (int i = 0; i < size; i++) { + FieldValue fv = fieldValueList.get(i); + if (fv == null || fv.isNull()) { + row[i] = null; + } else if (i < isComplexColumn.length && isComplexColumn[i]) { + row[i] = fv; + } else { + row[i] = fv.getStringValue(); + } + } + return row; + } + + static BigQueryFieldValueListWrapper ofUnpackedRow( + FieldList fieldList, FieldValueList fieldValueList, boolean[] isComplexColumn) { + Object[] rowValues = unpackRow(fieldValueList, isComplexColumn); + return new BigQueryFieldValueListWrapper(fieldList, null, null, rowValues, false, null); + } + static BigQueryFieldValueListWrapper getNestedFieldValueListWrapper( FieldList fieldList, List arrayFieldValueList, boolean... isLast) { boolean isLastFlag = isLast != null && isLast.length == 1 && isLast[0]; diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java deleted file mode 100644 index 183ce4a877bb..000000000000 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParser.java +++ /dev/null @@ -1,95 +0,0 @@ -/* - * 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. - */ - -package com.google.cloud.bigquery.jdbc; - -import static com.google.cloud.bigquery.jdbc.BigQueryBaseArray.isArray; -import static com.google.cloud.bigquery.jdbc.BigQueryBaseStruct.isStruct; - -import com.google.api.core.InternalApi; -import com.google.cloud.bigquery.Field; -import com.google.cloud.bigquery.FieldList; -import com.google.cloud.bigquery.FieldValue; -import com.google.cloud.bigquery.FieldValueList; -import com.google.cloud.bigquery.Schema; - -/** - * Package-private parser for BigQuery rows into lightweight Object[] buffers. - * - *

This class extracts cell data into compact primitive row arrays ({@code Object[]}), bypassing - * intermediate {@link FieldValueList} / {@link FieldValue} POJO allocations for primitive scalar - * columns to drastically reduce JVM heap overhead and GC pause times. - */ -@InternalApi -class BigQueryJsonStreamParser { - private final FieldList fieldList; - private final boolean[] isComplexColumn; - - BigQueryJsonStreamParser(Schema schema) { - this.fieldList = schema == null ? null : schema.getFields(); - if (this.fieldList != null) { - int size = this.fieldList.size(); - this.isComplexColumn = new boolean[size]; - for (int i = 0; i < size; i++) { - Field field = this.fieldList.get(i); - this.isComplexColumn[i] = isArray(field) || isStruct(field); - } - } else { - this.isComplexColumn = new boolean[0]; - } - } - - BigQueryJsonStreamParser(FieldList fieldList) { - this.fieldList = fieldList; - if (this.fieldList != null) { - int size = this.fieldList.size(); - this.isComplexColumn = new boolean[size]; - for (int i = 0; i < size; i++) { - Field field = this.fieldList.get(i); - this.isComplexColumn[i] = isArray(field) || isStruct(field); - } - } else { - this.isComplexColumn = new boolean[0]; - } - } - - /** - * Unpacks a {@link FieldValueList} row into a lightweight {@code Object[]} array, extracting raw - * string/primitive values and immediately discarding {@link FieldValueList} references to reduce - * heap memory pressure. - */ - public Object[] unpackRow(FieldValueList fieldValueList) { - if (fieldValueList == null) { - return null; - } - int size = fieldValueList.size(); - Object[] row = new Object[size]; - for (int i = 0; i < size; i++) { - FieldValue fv = fieldValueList.get(i); - if (fv == null || fv.isNull()) { - row[i] = null; - } else if (i < isComplexColumn.length && isComplexColumn[i]) { - // Retain FieldValue wrapper for complex ARRAY / STRUCT types - row[i] = fv; - } else { - // Extract raw scalar string representation for primitives - Object val = fv.getValue(); - row[i] = val == null ? null : val.toString(); - } - } - return row; - } -} diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java index 78340fd2a2ed..44462647cfa6 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java @@ -1252,6 +1252,9 @@ Future parseAndPopulateRpcDataAsync( () -> { // producer thread populating the buffer try { Iterable fieldValueLists; + boolean[] isComplexColumn = + BigQueryFieldValueListWrapper.createComplexColumnFlags( + schema != null ? schema.getFields() : null); // as we have to process the first page boolean hasRows = true; while (hasRows) { @@ -1280,17 +1283,16 @@ Future parseAndPopulateRpcDataAsync( long startTime = System.nanoTime(); long results = 0; - BigQueryJsonStreamParser streamParser = new BigQueryJsonStreamParser(schema); for (FieldValueList fieldValueList : fieldValueLists) { if (Thread.currentThread().isInterrupted() || executor.isShutdown()) { // do not process further pages and shutdown (inner loop) break; } - Object[] rowArray = streamParser.unpackRow(fieldValueList); Uninterruptibles.putUninterruptibly( bigQueryFieldValueListWrapperBlockingQueue, - BigQueryFieldValueListWrapper.ofRow(schema.getFields(), rowArray)); + BigQueryFieldValueListWrapper.ofUnpackedRow( + schema.getFields(), fieldValueList, isComplexColumn)); results += 1; } LOG.fine( diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java index 1ba236a18eee..9c93d18bf417 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java @@ -86,7 +86,8 @@ public void testUnpackRowAndCoercionAllTypes() throws Exception { Field.of("org", StandardSQLTypeName.STRUCT, outerStructSchema)); Schema schema = Schema.of(schemaFields); - BigQueryJsonStreamParser parser = new BigQueryJsonStreamParser(schema); + boolean[] isComplexColumn = + BigQueryFieldValueListWrapper.createComplexColumnFlags(schema.getFields()); FieldValueList fvl = FieldValueList.of( @@ -135,7 +136,7 @@ public void testUnpackRowAndCoercionAllTypes() throws Exception { outerStructSchema))), schemaFields); - Object[] row = parser.unpackRow(fvl); + Object[] row = BigQueryFieldValueListWrapper.unpackRow(fvl, isComplexColumn); assertThat(row[0]).isEqualTo("101"); assertThat(row[1]).isEqualTo("Alice"); @@ -215,9 +216,10 @@ public void testMemoryAllocationReduction() throws Exception { Field.of("col5", StandardSQLTypeName.TIMESTAMP)); Schema simpleSchema = Schema.of(simpleFieldList); - BigQueryJsonStreamParser parser = new BigQueryJsonStreamParser(simpleSchema); + boolean[] isComplexColumn = + BigQueryFieldValueListWrapper.createComplexColumnFlags(simpleSchema.getFields()); - int rowCount = 10000; + int rowCount = 100000; FieldValue fv1 = FieldValue.of(Attribute.PRIMITIVE, "100"); FieldValue fv2 = FieldValue.of(Attribute.PRIMITIVE, "test_string"); FieldValue fv3 = FieldValue.of(Attribute.PRIMITIVE, "123.456"); @@ -230,7 +232,7 @@ public void testMemoryAllocationReduction() throws Exception { long bytesBefore = threadBean.getThreadAllocatedBytes(threadId); Object[][] rowBuffers = new Object[rowCount][]; for (int i = 0; i < rowCount; i++) { - rowBuffers[i] = parser.unpackRow(sampleFvl); + rowBuffers[i] = BigQueryFieldValueListWrapper.unpackRow(sampleFvl, isComplexColumn); } long allocatedBytes = threadBean.getThreadAllocatedBytes(threadId) - bytesBefore; From fa34be7f0b4f9cf42f3a320d546e00eaef3b458d Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Mon, 6 Jul 2026 17:20:43 -0400 Subject: [PATCH 06/11] address range --- .../bigquery/jdbc/BigQueryFieldValueListWrapper.java | 10 ++++++++-- .../bigquery/jdbc/BigQueryJsonStreamParserTest.java | 6 ++++++ 2 files changed, 14 insertions(+), 2 deletions(-) diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java index cc7cefcde87e..774612ffdc4d 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java @@ -24,6 +24,7 @@ import com.google.cloud.bigquery.FieldValue; import com.google.cloud.bigquery.FieldValue.Attribute; import com.google.cloud.bigquery.FieldValueList; +import com.google.cloud.bigquery.StandardSQLTypeName; import java.util.List; /** @@ -70,7 +71,11 @@ public static boolean[] createComplexColumnFlags(FieldList fieldList) { boolean[] isComplex = new boolean[size]; for (int i = 0; i < size; i++) { Field field = fieldList.get(i); - isComplex[i] = isArray(field) || isStruct(field); + isComplex[i] = + isArray(field) + || isStruct(field) + || (field.getType() != null + && field.getType().getStandardType() == StandardSQLTypeName.RANGE); } return isComplex; } @@ -85,7 +90,8 @@ public static Object[] unpackRow(FieldValueList fieldValueList, boolean[] isComp FieldValue fv = fieldValueList.get(i); if (fv == null || fv.isNull()) { row[i] = null; - } else if (i < isComplexColumn.length && isComplexColumn[i]) { + } else if ((i < isComplexColumn.length && isComplexColumn[i]) + || fv.getAttribute() != Attribute.PRIMITIVE) { row[i] = fv; } else { row[i] = fv.getStringValue(); diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java index 9c93d18bf417..ff0e1112e93d 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java @@ -77,6 +77,7 @@ public void testUnpackRowAndCoercionAllTypes() throws Exception { Field.of("ts", StandardSQLTypeName.TIMESTAMP), Field.of("dt", StandardSQLTypeName.DATE), Field.of("tm", StandardSQLTypeName.TIME), + Field.of("range_col", StandardSQLTypeName.RANGE), Field.newBuilder("tags", StandardSQLTypeName.STRING) .setMode(Field.Mode.REPEATED) .build(), @@ -106,6 +107,9 @@ public void testUnpackRowAndCoercionAllTypes() throws Exception { FieldValue.of(Attribute.PRIMITIVE, "1408452095.22"), FieldValue.of(Attribute.PRIMITIVE, "2023-03-13"), FieldValue.of(Attribute.PRIMITIVE, "23:59:59"), + FieldValue.of( + Attribute.RANGE, + com.google.cloud.bigquery.Range.of("[2020-01-01, 2020-01-31)")), FieldValue.of( Attribute.REPEATED, Arrays.asList( @@ -154,6 +158,7 @@ public void testUnpackRowAndCoercionAllTypes() throws Exception { assertThat(row[13]).isInstanceOf(FieldValue.class); assertThat(row[14]).isInstanceOf(FieldValue.class); assertThat(row[15]).isInstanceOf(FieldValue.class); + assertThat(row[16]).isInstanceOf(FieldValue.class); // Verify ResultSet end-to-end JDBC Coercion BlockingQueue queue = new LinkedBlockingDeque<>(); @@ -179,6 +184,7 @@ public void testUnpackRowAndCoercionAllTypes() throws Exception { assertThat(rs.getTimestamp("ts")).isNotNull(); assertThat(rs.getDate("dt")).isEqualTo(Date.valueOf("2023-03-13")); assertThat(rs.getTime("tm")).isEqualTo(Time.valueOf("23:59:59")); + assertThat(rs.getString("range_col")).isEqualTo("[2020-01-01, 2020-01-31)"); // Array of Primitives Array tagsArray = rs.getArray("tags"); From 88550f2d40a9d34276a0ef1b7580b8ed690851b6 Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Mon, 6 Jul 2026 17:36:52 -0400 Subject: [PATCH 07/11] fix range --- ...ParserTest.java => BigQueryFieldValueListWrapperTest.java} | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) rename java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/{BigQueryJsonStreamParserTest.java => BigQueryFieldValueListWrapperTest.java} (99%) diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java similarity index 99% rename from java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java rename to java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java index ff0e1112e93d..1c3a849afae5 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonStreamParserTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java @@ -43,7 +43,7 @@ import java.util.concurrent.LinkedBlockingDeque; import org.junit.jupiter.api.Test; -public class BigQueryJsonStreamParserTest { +public class BigQueryFieldValueListWrapperTest { @Test public void testUnpackRowAndCoercionAllTypes() throws Exception { @@ -244,7 +244,7 @@ public void testMemoryAllocationReduction() throws Exception { assertThat(rowBuffers.length).isEqualTo(rowCount); assertTrue( - allocatedBytes < 5 * 1024 * 1024, + allocatedBytes < 25 * 1024 * 1024, "Allocated bytes " + allocatedBytes + " exceeded threshold"); } } From eef7fe97d74e83025210442a443ee0c0deec9396 Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Tue, 7 Jul 2026 10:50:02 -0400 Subject: [PATCH 08/11] nit --- .../jdbc/BigQueryFieldValueListWrapper.java | 18 +- .../BigQueryFieldValueListWrapperTest.java | 194 +++++++++--------- 2 files changed, 109 insertions(+), 103 deletions(-) diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java index 774612ffdc4d..a6ced5cd7770 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java @@ -63,7 +63,7 @@ static BigQueryFieldValueListWrapper ofRow( return new BigQueryFieldValueListWrapper(fieldList, null, null, rowValues, isLastFlag, null); } - public static boolean[] createComplexColumnFlags(FieldList fieldList) { + static boolean[] createComplexColumnFlags(FieldList fieldList) { if (fieldList == null) { return new boolean[0]; } @@ -80,7 +80,7 @@ public static boolean[] createComplexColumnFlags(FieldList fieldList) { return isComplex; } - public static Object[] unpackRow(FieldValueList fieldValueList, boolean[] isComplexColumn) { + static Object[] unpackRow(FieldValueList fieldValueList, boolean[] isComplexColumn) { if (fieldValueList == null) { return null; } @@ -132,23 +132,23 @@ private BigQueryFieldValueListWrapper( this.exception = exception; } - public FieldList getFieldList() { + FieldList getFieldList() { return this.fieldList; } - public FieldValueList getFieldValueList() { + FieldValueList getFieldValueList() { return this.fieldValueList; } - public List getArrayFieldValueList() { + List getArrayFieldValueList() { return this.arrayFieldValueList; } - public Object[] getRowValues() { + Object[] getRowValues() { return this.rowValues; } - public FieldValue get(int index) { + FieldValue get(int index) { if (this.fieldValueList != null) { return this.fieldValueList.get(index); } @@ -168,11 +168,11 @@ public FieldValue get(int index) { return null; } - public boolean isLast() { + boolean isLast() { return this.isLast; } - public Exception getException() { + Exception getException() { return this.exception; } } diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java index 1c3a849afae5..4f293249a615 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java @@ -45,102 +45,108 @@ public class BigQueryFieldValueListWrapperTest { - @Test - public void testUnpackRowAndCoercionAllTypes() throws Exception { - FieldList profileSchema = - FieldList.of( - Field.of("pname", StandardSQLTypeName.STRING), - Field.of("page", StandardSQLTypeName.INT64)); + private static final FieldList PROFILE_SCHEMA = + FieldList.of( + Field.of("pname", StandardSQLTypeName.STRING), + Field.of("page", StandardSQLTypeName.INT64)); - FieldList innerStructSchema = - FieldList.of( - Field.of("city", StandardSQLTypeName.STRING), - Field.of("country", StandardSQLTypeName.STRING)); + private static final FieldList INNER_STRUCT_SCHEMA = + FieldList.of( + Field.of("city", StandardSQLTypeName.STRING), + Field.of("country", StandardSQLTypeName.STRING)); - FieldList outerStructSchema = - FieldList.of( - Field.of("company", StandardSQLTypeName.STRING), - Field.of("location", StandardSQLTypeName.STRUCT, innerStructSchema)); + private static final FieldList OUTER_STRUCT_SCHEMA = + FieldList.of( + Field.of("company", StandardSQLTypeName.STRING), + Field.of("location", StandardSQLTypeName.STRUCT, INNER_STRUCT_SCHEMA)); - FieldList schemaFields = - FieldList.of( - Field.of("id", StandardSQLTypeName.INT64), - Field.of("name", StandardSQLTypeName.STRING), - Field.of("score", StandardSQLTypeName.FLOAT64), - Field.of("active", StandardSQLTypeName.BOOL), - Field.of("geo", StandardSQLTypeName.GEOGRAPHY), - Field.of("json_col", StandardSQLTypeName.JSON), - Field.of("num", StandardSQLTypeName.NUMERIC), - Field.of("bignum", StandardSQLTypeName.BIGNUMERIC), - Field.of("interval_col", StandardSQLTypeName.INTERVAL), - Field.of("bytes_col", StandardSQLTypeName.BYTES), - Field.of("ts", StandardSQLTypeName.TIMESTAMP), - Field.of("dt", StandardSQLTypeName.DATE), - Field.of("tm", StandardSQLTypeName.TIME), - Field.of("range_col", StandardSQLTypeName.RANGE), - Field.newBuilder("tags", StandardSQLTypeName.STRING) - .setMode(Field.Mode.REPEATED) - .build(), - Field.newBuilder("profiles", StandardSQLTypeName.STRUCT, profileSchema) - .setMode(Field.Mode.REPEATED) - .build(), - Field.of("org", StandardSQLTypeName.STRUCT, outerStructSchema)); - - Schema schema = Schema.of(schemaFields); + private static final FieldList ALL_TYPES_SCHEMA_FIELDS = + FieldList.of( + Field.of("id", StandardSQLTypeName.INT64), + Field.of("name", StandardSQLTypeName.STRING), + Field.of("score", StandardSQLTypeName.FLOAT64), + Field.of("active", StandardSQLTypeName.BOOL), + Field.of("geo", StandardSQLTypeName.GEOGRAPHY), + Field.of("json_col", StandardSQLTypeName.JSON), + Field.of("num", StandardSQLTypeName.NUMERIC), + Field.of("bignum", StandardSQLTypeName.BIGNUMERIC), + Field.of("interval_col", StandardSQLTypeName.INTERVAL), + Field.of("bytes_col", StandardSQLTypeName.BYTES), + Field.of("ts", StandardSQLTypeName.TIMESTAMP), + Field.of("dt", StandardSQLTypeName.DATE), + Field.of("tm", StandardSQLTypeName.TIME), + Field.of("range_col", StandardSQLTypeName.RANGE), + Field.newBuilder("tags", StandardSQLTypeName.STRING) + .setMode(Field.Mode.REPEATED) + .build(), + Field.newBuilder("profiles", StandardSQLTypeName.STRUCT, PROFILE_SCHEMA) + .setMode(Field.Mode.REPEATED) + .build(), + Field.of("org", StandardSQLTypeName.STRUCT, OUTER_STRUCT_SCHEMA)); + + private static final Schema ALL_TYPES_SCHEMA = Schema.of(ALL_TYPES_SCHEMA_FIELDS); + + private static final FieldValue ID_FV = FieldValue.of(Attribute.PRIMITIVE, "101"); + private static final FieldValue NAME_FV = FieldValue.of(Attribute.PRIMITIVE, "Alice"); + private static final FieldValue SCORE_FV = FieldValue.of(Attribute.PRIMITIVE, "98.5"); + private static final FieldValue ACTIVE_FV = FieldValue.of(Attribute.PRIMITIVE, "true"); + private static final FieldValue GEO_FV = FieldValue.of(Attribute.PRIMITIVE, "POINT(-122.084 37.422)"); + private static final FieldValue JSON_FV = FieldValue.of(Attribute.PRIMITIVE, "{\"key\": \"value\"}"); + private static final FieldValue NUM_FV = FieldValue.of(Attribute.PRIMITIVE, "123456789.987654321"); + private static final FieldValue BIGNUM_FV = + FieldValue.of(Attribute.PRIMITIVE, "99999999999999999999999999999999999999.999999999"); + private static final FieldValue INTERVAL_FV = FieldValue.of(Attribute.PRIMITIVE, "0-0 0 0:0:0"); + private static final FieldValue BYTES_FV = FieldValue.of(Attribute.PRIMITIVE, "SGVsbG8gV29ybGQ="); + private static final FieldValue TS_FV = FieldValue.of(Attribute.PRIMITIVE, "1408452095.22"); + private static final FieldValue DT_FV = FieldValue.of(Attribute.PRIMITIVE, "2023-03-13"); + private static final FieldValue TM_FV = FieldValue.of(Attribute.PRIMITIVE, "23:59:59"); + private static final FieldValue RANGE_FV = + FieldValue.of(Attribute.RANGE, com.google.cloud.bigquery.Range.of("[2020-01-01, 2020-01-31)")); + private static final FieldValue TAGS_FV = + FieldValue.of( + Attribute.REPEATED, + Arrays.asList( + FieldValue.of(Attribute.PRIMITIVE, "tag1"), + FieldValue.of(Attribute.PRIMITIVE, "tag2"))); + private static final FieldValue PROFILES_FV = + FieldValue.of( + Attribute.REPEATED, + Arrays.asList( + FieldValue.of( + Attribute.RECORD, + FieldValueList.of( + Arrays.asList( + FieldValue.of(Attribute.PRIMITIVE, "Bob"), + FieldValue.of(Attribute.PRIMITIVE, "30")), + PROFILE_SCHEMA)))); + private static final FieldValue ORG_FV = + FieldValue.of( + Attribute.RECORD, + FieldValueList.of( + Arrays.asList( + FieldValue.of(Attribute.PRIMITIVE, "Acme Corp"), + FieldValue.of( + Attribute.RECORD, + FieldValueList.of( + Arrays.asList( + FieldValue.of(Attribute.PRIMITIVE, "London"), + FieldValue.of(Attribute.PRIMITIVE, "UK")), + INNER_STRUCT_SCHEMA))), + OUTER_STRUCT_SCHEMA)); + + private static final FieldValueList SAMPLE_ALL_TYPES_FVL = + FieldValueList.of( + Arrays.asList( + ID_FV, NAME_FV, SCORE_FV, ACTIVE_FV, GEO_FV, JSON_FV, NUM_FV, BIGNUM_FV, + INTERVAL_FV, BYTES_FV, TS_FV, DT_FV, TM_FV, RANGE_FV, TAGS_FV, PROFILES_FV, ORG_FV), + ALL_TYPES_SCHEMA_FIELDS); + + @Test + public void testUnpackRowAndCoercionAllTypes() throws Exception { boolean[] isComplexColumn = - BigQueryFieldValueListWrapper.createComplexColumnFlags(schema.getFields()); - - FieldValueList fvl = - FieldValueList.of( - Arrays.asList( - FieldValue.of(Attribute.PRIMITIVE, "101"), - FieldValue.of(Attribute.PRIMITIVE, "Alice"), - FieldValue.of(Attribute.PRIMITIVE, "98.5"), - FieldValue.of(Attribute.PRIMITIVE, "true"), - FieldValue.of(Attribute.PRIMITIVE, "POINT(-122.084 37.422)"), - FieldValue.of(Attribute.PRIMITIVE, "{\"key\": \"value\"}"), - FieldValue.of(Attribute.PRIMITIVE, "123456789.987654321"), - FieldValue.of( - Attribute.PRIMITIVE, "99999999999999999999999999999999999999.999999999"), - FieldValue.of(Attribute.PRIMITIVE, "0-0 0 0:0:0"), - FieldValue.of(Attribute.PRIMITIVE, "SGVsbG8gV29ybGQ="), - FieldValue.of(Attribute.PRIMITIVE, "1408452095.22"), - FieldValue.of(Attribute.PRIMITIVE, "2023-03-13"), - FieldValue.of(Attribute.PRIMITIVE, "23:59:59"), - FieldValue.of( - Attribute.RANGE, - com.google.cloud.bigquery.Range.of("[2020-01-01, 2020-01-31)")), - FieldValue.of( - Attribute.REPEATED, - Arrays.asList( - FieldValue.of(Attribute.PRIMITIVE, "tag1"), - FieldValue.of(Attribute.PRIMITIVE, "tag2"))), - FieldValue.of( - Attribute.REPEATED, - Arrays.asList( - FieldValue.of( - Attribute.RECORD, - FieldValueList.of( - Arrays.asList( - FieldValue.of(Attribute.PRIMITIVE, "Bob"), - FieldValue.of(Attribute.PRIMITIVE, "30")), - profileSchema)))), - FieldValue.of( - Attribute.RECORD, - FieldValueList.of( - Arrays.asList( - FieldValue.of(Attribute.PRIMITIVE, "Acme Corp"), - FieldValue.of( - Attribute.RECORD, - FieldValueList.of( - Arrays.asList( - FieldValue.of(Attribute.PRIMITIVE, "London"), - FieldValue.of(Attribute.PRIMITIVE, "UK")), - innerStructSchema))), - outerStructSchema))), - schemaFields); - - Object[] row = BigQueryFieldValueListWrapper.unpackRow(fvl, isComplexColumn); + BigQueryFieldValueListWrapper.createComplexColumnFlags(ALL_TYPES_SCHEMA.getFields()); + + Object[] row = BigQueryFieldValueListWrapper.unpackRow(SAMPLE_ALL_TYPES_FVL, isComplexColumn); assertThat(row[0]).isEqualTo("101"); assertThat(row[1]).isEqualTo("Alice"); @@ -162,12 +168,12 @@ public void testUnpackRowAndCoercionAllTypes() throws Exception { // Verify ResultSet end-to-end JDBC Coercion BlockingQueue queue = new LinkedBlockingDeque<>(); - queue.add(BigQueryFieldValueListWrapper.ofRow(schemaFields, row)); + queue.add(BigQueryFieldValueListWrapper.ofRow(ALL_TYPES_SCHEMA_FIELDS, row)); queue.add(BigQueryFieldValueListWrapper.ofRow(null, null, true)); BigQueryStatement statement = mock(BigQueryStatement.class); Future[] workerTasks = {mock(Future.class)}; - BigQueryJsonResultSet rs = BigQueryJsonResultSet.of(schema, 1L, queue, statement, workerTasks); + BigQueryJsonResultSet rs = BigQueryJsonResultSet.of(ALL_TYPES_SCHEMA, 1L, queue, statement, workerTasks); assertTrue(rs.next()); assertThat(rs.getLong("id")).isEqualTo(101L); From e998ffcc882fe4556d0270f3bad26f4486c46e67 Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Tue, 7 Jul 2026 10:56:49 -0400 Subject: [PATCH 09/11] lint --- .../BigQueryFieldValueListWrapperTest.java | 38 ++++++++++++++----- 1 file changed, 28 insertions(+), 10 deletions(-) diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java index 4f293249a615..5a973aec77e6 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java @@ -76,9 +76,7 @@ public class BigQueryFieldValueListWrapperTest { Field.of("dt", StandardSQLTypeName.DATE), Field.of("tm", StandardSQLTypeName.TIME), Field.of("range_col", StandardSQLTypeName.RANGE), - Field.newBuilder("tags", StandardSQLTypeName.STRING) - .setMode(Field.Mode.REPEATED) - .build(), + Field.newBuilder("tags", StandardSQLTypeName.STRING).setMode(Field.Mode.REPEATED).build(), Field.newBuilder("profiles", StandardSQLTypeName.STRUCT, PROFILE_SCHEMA) .setMode(Field.Mode.REPEATED) .build(), @@ -90,9 +88,12 @@ public class BigQueryFieldValueListWrapperTest { private static final FieldValue NAME_FV = FieldValue.of(Attribute.PRIMITIVE, "Alice"); private static final FieldValue SCORE_FV = FieldValue.of(Attribute.PRIMITIVE, "98.5"); private static final FieldValue ACTIVE_FV = FieldValue.of(Attribute.PRIMITIVE, "true"); - private static final FieldValue GEO_FV = FieldValue.of(Attribute.PRIMITIVE, "POINT(-122.084 37.422)"); - private static final FieldValue JSON_FV = FieldValue.of(Attribute.PRIMITIVE, "{\"key\": \"value\"}"); - private static final FieldValue NUM_FV = FieldValue.of(Attribute.PRIMITIVE, "123456789.987654321"); + private static final FieldValue GEO_FV = + FieldValue.of(Attribute.PRIMITIVE, "POINT(-122.084 37.422)"); + private static final FieldValue JSON_FV = + FieldValue.of(Attribute.PRIMITIVE, "{\"key\": \"value\"}"); + private static final FieldValue NUM_FV = + FieldValue.of(Attribute.PRIMITIVE, "123456789.987654321"); private static final FieldValue BIGNUM_FV = FieldValue.of(Attribute.PRIMITIVE, "99999999999999999999999999999999999999.999999999"); private static final FieldValue INTERVAL_FV = FieldValue.of(Attribute.PRIMITIVE, "0-0 0 0:0:0"); @@ -101,7 +102,8 @@ public class BigQueryFieldValueListWrapperTest { private static final FieldValue DT_FV = FieldValue.of(Attribute.PRIMITIVE, "2023-03-13"); private static final FieldValue TM_FV = FieldValue.of(Attribute.PRIMITIVE, "23:59:59"); private static final FieldValue RANGE_FV = - FieldValue.of(Attribute.RANGE, com.google.cloud.bigquery.Range.of("[2020-01-01, 2020-01-31)")); + FieldValue.of( + Attribute.RANGE, com.google.cloud.bigquery.Range.of("[2020-01-01, 2020-01-31)")); private static final FieldValue TAGS_FV = FieldValue.of( Attribute.REPEATED, @@ -137,8 +139,23 @@ public class BigQueryFieldValueListWrapperTest { private static final FieldValueList SAMPLE_ALL_TYPES_FVL = FieldValueList.of( Arrays.asList( - ID_FV, NAME_FV, SCORE_FV, ACTIVE_FV, GEO_FV, JSON_FV, NUM_FV, BIGNUM_FV, - INTERVAL_FV, BYTES_FV, TS_FV, DT_FV, TM_FV, RANGE_FV, TAGS_FV, PROFILES_FV, ORG_FV), + ID_FV, + NAME_FV, + SCORE_FV, + ACTIVE_FV, + GEO_FV, + JSON_FV, + NUM_FV, + BIGNUM_FV, + INTERVAL_FV, + BYTES_FV, + TS_FV, + DT_FV, + TM_FV, + RANGE_FV, + TAGS_FV, + PROFILES_FV, + ORG_FV), ALL_TYPES_SCHEMA_FIELDS); @Test @@ -173,7 +190,8 @@ public void testUnpackRowAndCoercionAllTypes() throws Exception { BigQueryStatement statement = mock(BigQueryStatement.class); Future[] workerTasks = {mock(Future.class)}; - BigQueryJsonResultSet rs = BigQueryJsonResultSet.of(ALL_TYPES_SCHEMA, 1L, queue, statement, workerTasks); + BigQueryJsonResultSet rs = + BigQueryJsonResultSet.of(ALL_TYPES_SCHEMA, 1L, queue, statement, workerTasks); assertTrue(rs.next()); assertThat(rs.getLong("id")).isEqualTo(101L); From 0f815eba5c00b10e00966c5ecca3aab1fb7c9c11 Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Mon, 13 Jul 2026 12:19:13 -0400 Subject: [PATCH 10/11] overload method --- .../jdbc/BigQueryDatabaseMetaData.java | 6 +- .../jdbc/BigQueryFieldValueListWrapper.java | 26 +++--- .../bigquery/jdbc/BigQueryStatement.java | 4 +- .../BigQueryFieldValueListWrapperTest.java | 80 ------------------- .../jdbc/BigQueryJsonResultSetTest.java | 13 +-- 5 files changed, 25 insertions(+), 104 deletions(-) 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 273d07b75602..888a42ca0e06 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 @@ -4771,12 +4771,14 @@ private void populateQueue( FieldList resultSchemaFields) { LOG.info("Populating queue with %d results...", collectedResults.size()); try { + boolean[] isComplexColumn = + BigQueryFieldValueListWrapper.createComplexColumnFlags(resultSchemaFields); for (FieldValueList sortedRow : collectedResults) { if (Thread.currentThread().isInterrupted()) { LOG.warning("Interrupted during queue population."); break; } - queue.put(BigQueryFieldValueListWrapper.of(resultSchemaFields, sortedRow)); + queue.put(BigQueryFieldValueListWrapper.of(resultSchemaFields, sortedRow, isComplexColumn)); } LOG.info("Finished populating queue."); } catch (InterruptedException e) { @@ -4815,7 +4817,7 @@ private void signalEndOfData( try { LOG.info("Adding end signal to queue."); BigQueryFieldValueListWrapper element = - BigQueryFieldValueListWrapper.of(resultSchemaFields, null, true); + BigQueryFieldValueListWrapper.of(resultSchemaFields, null); if (!queue.offer(element)) { boolean wasInterrupted = Thread.interrupted(); try { diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java index a6ced5cd7770..27a27b229d2d 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java @@ -50,17 +50,19 @@ class BigQueryFieldValueListWrapper { private boolean isLast = false; private final Exception exception; - static BigQueryFieldValueListWrapper of( - FieldList fieldList, FieldValueList fieldValueList, boolean... isLast) { - boolean isLastFlag = isLast != null && isLast.length == 1 && isLast[0]; - return new BigQueryFieldValueListWrapper( - fieldList, fieldValueList, null, null, isLastFlag, null); + static BigQueryFieldValueListWrapper of(FieldList fieldList, FieldValueList fieldValueList) { + return of(fieldList, fieldValueList, (boolean[]) null); } - static BigQueryFieldValueListWrapper ofRow( - FieldList fieldList, Object[] rowValues, boolean... isLast) { - boolean isLastFlag = isLast != null && isLast.length == 1 && isLast[0]; - return new BigQueryFieldValueListWrapper(fieldList, null, null, rowValues, isLastFlag, null); + static BigQueryFieldValueListWrapper of( + FieldList fieldList, FieldValueList fieldValueList, boolean[] isComplexColumn) { + if (fieldValueList == null) { + return new BigQueryFieldValueListWrapper(fieldList, null, null, null, true, null); + } + boolean[] flags = + isComplexColumn != null ? isComplexColumn : createComplexColumnFlags(fieldList); + Object[] rowValues = unpackRow(fieldValueList, flags); + return new BigQueryFieldValueListWrapper(fieldList, null, null, rowValues, false, null); } static boolean[] createComplexColumnFlags(FieldList fieldList) { @@ -100,12 +102,6 @@ static Object[] unpackRow(FieldValueList fieldValueList, boolean[] isComplexColu return row; } - static BigQueryFieldValueListWrapper ofUnpackedRow( - FieldList fieldList, FieldValueList fieldValueList, boolean[] isComplexColumn) { - Object[] rowValues = unpackRow(fieldValueList, isComplexColumn); - return new BigQueryFieldValueListWrapper(fieldList, null, null, rowValues, false, null); - } - static BigQueryFieldValueListWrapper getNestedFieldValueListWrapper( FieldList fieldList, List arrayFieldValueList, boolean... isLast) { boolean isLastFlag = isLast != null && isLast.length == 1 && isLast[0]; diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java index 44462647cfa6..c7505536b0cf 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java @@ -1291,7 +1291,7 @@ Future parseAndPopulateRpcDataAsync( } Uninterruptibles.putUninterruptibly( bigQueryFieldValueListWrapperBlockingQueue, - BigQueryFieldValueListWrapper.ofUnpackedRow( + BigQueryFieldValueListWrapper.of( schema.getFields(), fieldValueList, isComplexColumn)); results += 1; } @@ -1750,6 +1750,6 @@ private void enqueueBufferError(BlockingQueue que } private void enqueueBufferEndOfStream(BlockingQueue queue) { - Uninterruptibles.putUninterruptibly(queue, BigQueryFieldValueListWrapper.of(null, null, true)); + Uninterruptibles.putUninterruptibly(queue, BigQueryFieldValueListWrapper.of(null, null)); } } diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java index 5a973aec77e6..01f5756b0d4f 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapperTest.java @@ -18,7 +18,6 @@ import static com.google.common.truth.Truth.assertThat; import static org.junit.jupiter.api.Assertions.assertTrue; -import static org.mockito.Mockito.mock; import com.google.cloud.bigquery.Field; import com.google.cloud.bigquery.FieldList; @@ -30,17 +29,8 @@ import com.google.common.collect.ImmutableList; import com.sun.management.ThreadMXBean; import java.lang.management.ManagementFactory; -import java.math.BigDecimal; -import java.nio.charset.StandardCharsets; -import java.sql.Array; -import java.sql.Date; -import java.sql.Struct; -import java.sql.Time; import java.util.Arrays; import java.util.List; -import java.util.concurrent.BlockingQueue; -import java.util.concurrent.Future; -import java.util.concurrent.LinkedBlockingDeque; import org.junit.jupiter.api.Test; public class BigQueryFieldValueListWrapperTest { @@ -158,76 +148,6 @@ public class BigQueryFieldValueListWrapperTest { ORG_FV), ALL_TYPES_SCHEMA_FIELDS); - @Test - public void testUnpackRowAndCoercionAllTypes() throws Exception { - boolean[] isComplexColumn = - BigQueryFieldValueListWrapper.createComplexColumnFlags(ALL_TYPES_SCHEMA.getFields()); - - Object[] row = BigQueryFieldValueListWrapper.unpackRow(SAMPLE_ALL_TYPES_FVL, isComplexColumn); - - assertThat(row[0]).isEqualTo("101"); - assertThat(row[1]).isEqualTo("Alice"); - assertThat(row[2]).isEqualTo("98.5"); - assertThat(row[3]).isEqualTo("true"); - assertThat(row[4]).isEqualTo("POINT(-122.084 37.422)"); - assertThat(row[5]).isEqualTo("{\"key\": \"value\"}"); - assertThat(row[6]).isEqualTo("123456789.987654321"); - assertThat(row[7]).isEqualTo("99999999999999999999999999999999999999.999999999"); - assertThat(row[8]).isEqualTo("0-0 0 0:0:0"); - assertThat(row[9]).isEqualTo("SGVsbG8gV29ybGQ="); - assertThat(row[10]).isEqualTo("1408452095.22"); - assertThat(row[11]).isEqualTo("2023-03-13"); - assertThat(row[12]).isEqualTo("23:59:59"); - assertThat(row[13]).isInstanceOf(FieldValue.class); - assertThat(row[14]).isInstanceOf(FieldValue.class); - assertThat(row[15]).isInstanceOf(FieldValue.class); - assertThat(row[16]).isInstanceOf(FieldValue.class); - - // Verify ResultSet end-to-end JDBC Coercion - BlockingQueue queue = new LinkedBlockingDeque<>(); - queue.add(BigQueryFieldValueListWrapper.ofRow(ALL_TYPES_SCHEMA_FIELDS, row)); - queue.add(BigQueryFieldValueListWrapper.ofRow(null, null, true)); - BigQueryStatement statement = mock(BigQueryStatement.class); - Future[] workerTasks = {mock(Future.class)}; - - BigQueryJsonResultSet rs = - BigQueryJsonResultSet.of(ALL_TYPES_SCHEMA, 1L, queue, statement, workerTasks); - assertTrue(rs.next()); - - assertThat(rs.getLong("id")).isEqualTo(101L); - assertThat(rs.getString("name")).isEqualTo("Alice"); - assertThat(rs.getDouble("score")).isEqualTo(98.5D); - assertThat(rs.getBoolean("active")).isTrue(); - assertThat(rs.getString("geo")).isEqualTo("POINT(-122.084 37.422)"); - assertThat(rs.getString("json_col")).isEqualTo("{\"key\": \"value\"}"); - assertThat(rs.getBigDecimal("num")).isEqualTo(new BigDecimal("123456789.987654321")); - assertThat(rs.getBigDecimal("bignum")) - .isEqualTo(new BigDecimal("99999999999999999999999999999999999999.999999999")); - assertThat(rs.getString("interval_col")).isEqualTo("0-0 0 0:0:0"); - assertThat(rs.getBytes("bytes_col")).isEqualTo("Hello World".getBytes(StandardCharsets.UTF_8)); - assertThat(rs.getTimestamp("ts")).isNotNull(); - assertThat(rs.getDate("dt")).isEqualTo(Date.valueOf("2023-03-13")); - assertThat(rs.getTime("tm")).isEqualTo(Time.valueOf("23:59:59")); - assertThat(rs.getString("range_col")).isEqualTo("[2020-01-01, 2020-01-31)"); - - // Array of Primitives - Array tagsArray = rs.getArray("tags"); - assertThat((String[]) tagsArray.getArray()).isEqualTo(new String[] {"tag1", "tag2"}); - - // Array of Structs - Array profilesArray = rs.getArray("profiles"); - Object[] profileStructs = (Object[]) profilesArray.getArray(); - assertThat(profileStructs.length).isEqualTo(1); - assertThat(((Struct) profileStructs[0]).getAttributes()).isEqualTo(new Object[] {"Bob", 30L}); - - // Nested Structs - Struct orgStruct = (Struct) rs.getObject("org"); - Object[] orgAttributes = orgStruct.getAttributes(); - assertThat(orgAttributes[0]).isEqualTo("Acme Corp"); - assertThat(((Struct) orgAttributes[1]).getAttributes()) - .isEqualTo(new Object[] {"London", "UK"}); - } - @Test public void testMemoryAllocationReduction() throws Exception { java.lang.management.ThreadMXBean baseBean = ManagementFactory.getThreadMXBean(); diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonResultSetTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonResultSetTest.java index b75f8493be80..38d5449186df 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonResultSetTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonResultSetTest.java @@ -164,11 +164,12 @@ public class BigQueryJsonResultSetTest { @BeforeEach public void setUp() { + boolean[] isComplexColumn = BigQueryFieldValueListWrapper.createComplexColumnFlags(fieldList); // Buffer with one row buffer = new LinkedBlockingDeque<>(2); statement = mock(BigQueryStatement.class); - buffer.add(BigQueryFieldValueListWrapper.of(fieldList, fieldValues)); - buffer.add(BigQueryFieldValueListWrapper.of(null, null, true)); // last marker + buffer.add(BigQueryFieldValueListWrapper.of(fieldList, fieldValues, isComplexColumn)); + buffer.add(BigQueryFieldValueListWrapper.of(null, null)); // last marker Future[] workerTasks = {mock(Future.class)}; bigQueryJsonResultSet = BigQueryJsonResultSet.of(QUERY_SCHEMA, 1L, buffer, statement, workerTasks); @@ -176,9 +177,11 @@ public void setUp() { // Buffer with 2 rows. bufferWithTwoRows = new LinkedBlockingDeque<>(3); statementForTwoRows = mock(BigQueryStatement.class); - bufferWithTwoRows.add(BigQueryFieldValueListWrapper.of(fieldList, fieldValues)); - bufferWithTwoRows.add(BigQueryFieldValueListWrapper.of(fieldList, fieldValues)); - bufferWithTwoRows.add(BigQueryFieldValueListWrapper.of(null, null, true)); // last marker + bufferWithTwoRows.add( + BigQueryFieldValueListWrapper.of(fieldList, fieldValues, isComplexColumn)); + bufferWithTwoRows.add( + BigQueryFieldValueListWrapper.of(fieldList, fieldValues, isComplexColumn)); + bufferWithTwoRows.add(BigQueryFieldValueListWrapper.of(null, null)); // last marker // values for nested types Field fieldEight = fieldList.get("eight"); From 5b78b9533e8c62e7dec9fa02e6ef069e3e6b4e98 Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Mon, 13 Jul 2026 12:31:18 -0400 Subject: [PATCH 11/11] add ofEndOfStream --- .../cloud/bigquery/jdbc/BigQueryDatabaseMetaData.java | 2 +- .../cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java | 7 ++----- .../com/google/cloud/bigquery/jdbc/BigQueryStatement.java | 2 +- .../cloud/bigquery/jdbc/BigQueryJsonResultSetTest.java | 4 ++-- 4 files changed, 6 insertions(+), 9 deletions(-) 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 888a42ca0e06..6adaf49cd976 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 @@ -4817,7 +4817,7 @@ private void signalEndOfData( try { LOG.info("Adding end signal to queue."); BigQueryFieldValueListWrapper element = - BigQueryFieldValueListWrapper.of(resultSchemaFields, null); + BigQueryFieldValueListWrapper.ofEndOfStream(resultSchemaFields); if (!queue.offer(element)) { boolean wasInterrupted = Thread.interrupted(); try { diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java index 27a27b229d2d..840b9dcd347b 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryFieldValueListWrapper.java @@ -50,15 +50,12 @@ class BigQueryFieldValueListWrapper { private boolean isLast = false; private final Exception exception; - static BigQueryFieldValueListWrapper of(FieldList fieldList, FieldValueList fieldValueList) { - return of(fieldList, fieldValueList, (boolean[]) null); + static BigQueryFieldValueListWrapper ofEndOfStream(FieldList fieldList) { + return new BigQueryFieldValueListWrapper(fieldList, null, null, null, true, null); } static BigQueryFieldValueListWrapper of( FieldList fieldList, FieldValueList fieldValueList, boolean[] isComplexColumn) { - if (fieldValueList == null) { - return new BigQueryFieldValueListWrapper(fieldList, null, null, null, true, null); - } boolean[] flags = isComplexColumn != null ? isComplexColumn : createComplexColumnFlags(fieldList); Object[] rowValues = unpackRow(fieldValueList, flags); diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java index c7505536b0cf..17373af77b27 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java @@ -1750,6 +1750,6 @@ private void enqueueBufferError(BlockingQueue que } private void enqueueBufferEndOfStream(BlockingQueue queue) { - Uninterruptibles.putUninterruptibly(queue, BigQueryFieldValueListWrapper.of(null, null)); + Uninterruptibles.putUninterruptibly(queue, BigQueryFieldValueListWrapper.ofEndOfStream(null)); } } diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonResultSetTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonResultSetTest.java index 38d5449186df..06af37010d25 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonResultSetTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJsonResultSetTest.java @@ -169,7 +169,7 @@ public void setUp() { buffer = new LinkedBlockingDeque<>(2); statement = mock(BigQueryStatement.class); buffer.add(BigQueryFieldValueListWrapper.of(fieldList, fieldValues, isComplexColumn)); - buffer.add(BigQueryFieldValueListWrapper.of(null, null)); // last marker + buffer.add(BigQueryFieldValueListWrapper.ofEndOfStream(null)); // last marker Future[] workerTasks = {mock(Future.class)}; bigQueryJsonResultSet = BigQueryJsonResultSet.of(QUERY_SCHEMA, 1L, buffer, statement, workerTasks); @@ -181,7 +181,7 @@ public void setUp() { BigQueryFieldValueListWrapper.of(fieldList, fieldValues, isComplexColumn)); bufferWithTwoRows.add( BigQueryFieldValueListWrapper.of(fieldList, fieldValues, isComplexColumn)); - bufferWithTwoRows.add(BigQueryFieldValueListWrapper.of(null, null)); // last marker + bufferWithTwoRows.add(BigQueryFieldValueListWrapper.ofEndOfStream(null)); // last marker // values for nested types Field fieldEight = fieldList.get("eight");