This is an automated email from the ASF dual-hosted git repository.
xiangfu0 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git
The following commit(s) were added to refs/heads/master by this push:
new cd504afe901 Add JDBC support for MAP, arrays and extended scalar types
cd504afe901 is described below
commit cd504afe901b008671fe77b60ad8e3c393d1515c
Author: Vojtech Mucha <[email protected]>
AuthorDate: Wed Sep 30 11:21:46 2026 +0200
Add JDBC support for MAP, arrays and extended scalar types
Use getObject("map", Map.class) for MAP columns and getObject("items",
List.class) for arrays; extended scalars return their corresponding Java types.
Reproduce the fixed edge cases with null or nested MAP values over JSON
gRPC, or results without column types.
Validated 29 focused tests on JDK 25, scoped formatting/license/Checkstyle
checks, and GitHub CI on the final commit.
---
.../org/apache/pinot/client/PinotResultSet.java | 281 ++++++-----------
.../pinot/client/base/AbstractBaseResultSet.java | 342 ++++++++++++++++++++-
.../pinot/client/grpc/PinotGrpcResultSet.java | 293 +++++++-----------
.../apache/pinot/client/utils/BigDecimalUtils.java | 36 +++
.../apache/pinot/client/utils/DateTimeUtils.java | 16 +-
.../org/apache/pinot/client/utils/DriverUtils.java | 20 ++
.../apache/pinot/client/PinotResultSetTest.java | 164 ++++++++--
.../pinot/client/grpc/PinotGrpcResultSetTest.java | 260 ++++++++++++++++
.../pinot/client/utils/BigDecimalUtilsTest.java | 48 +++
.../response/encoder/JsonResponseEncoder.java | 95 +-----
.../org/apache/pinot/common/utils/DataSchema.java | 19 ++
11 files changed, 1077 insertions(+), 497 deletions(-)
diff --git
a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/PinotResultSet.java
b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/PinotResultSet.java
index f62b52e497e..7ed5227015e 100644
---
a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/PinotResultSet.java
+++
b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/PinotResultSet.java
@@ -18,26 +18,22 @@
*/
package org.apache.pinot.client;
+import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode;
-import java.io.ByteArrayInputStream;
-import java.io.InputStream;
-import java.io.InputStreamReader;
-import java.io.Reader;
+import com.fasterxml.jackson.databind.ObjectReader;
+import java.io.IOException;
import java.math.BigDecimal;
-import java.net.URL;
-import java.nio.charset.StandardCharsets;
-import java.sql.Date;
import java.sql.ResultSetMetaData;
import java.sql.SQLDataException;
import java.sql.SQLException;
-import java.sql.Time;
import java.sql.Timestamp;
-import java.util.Calendar;
import java.util.HashMap;
+import java.util.List;
import java.util.Map;
-import org.apache.commons.codec.binary.Hex;
+import java.util.UUID;
+import javax.annotation.Nullable;
import org.apache.pinot.client.base.AbstractBaseResultSet;
-import org.apache.pinot.client.utils.DateTimeUtils;
+import org.apache.pinot.common.utils.DataSchema.ColumnDataType;
import org.apache.pinot.spi.utils.JsonUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -46,12 +42,29 @@ import org.slf4j.LoggerFactory;
public class PinotResultSet extends AbstractBaseResultSet {
public static final String NULL_STRING = "null";
private static final Logger LOG =
LoggerFactory.getLogger(PinotResultSet.class);
+ private static final ObjectReader MAP_READER =
+ JsonUtils.DEFAULT_READER.forType(new TypeReference<Map<?, ?>>() { });
+ private static final ObjectReader BOOLEAN_LIST_READER =
+ JsonUtils.DEFAULT_READER.forType(new TypeReference<List<Boolean>>() { });
+ private static final ObjectReader INT_LIST_READER =
+ JsonUtils.DEFAULT_READER.forType(new TypeReference<List<Integer>>() { });
+ private static final ObjectReader LONG_LIST_READER =
+ JsonUtils.DEFAULT_READER.forType(new TypeReference<List<Long>>() { });
+ private static final ObjectReader FLOAT_LIST_READER =
+ JsonUtils.DEFAULT_READER.forType(new TypeReference<List<Float>>() { });
+ private static final ObjectReader DOUBLE_LIST_READER =
+ JsonUtils.DEFAULT_READER.forType(new TypeReference<List<Double>>() { });
+ private static final ObjectReader BIG_DECIMAL_LIST_READER =
+ JsonUtils.DEFAULT_READER.forType(new TypeReference<List<BigDecimal>>() {
});
+ private static final ObjectReader STRING_LIST_READER =
+ JsonUtils.DEFAULT_READER.forType(new TypeReference<List<String>>() { });
private org.apache.pinot.client.ResultSet _resultSet;
private int _totalRows;
private int _currentRow;
private final int _totalColumns;
private final Map<String, Integer> _columns = new HashMap<>();
private final Map<Integer, String> _columnDataTypes = new HashMap<>();
+ private final Map<Integer, ColumnDataType> _resolvedColumnDataTypes = new
HashMap<>();
private boolean _closed;
private boolean _wasNull = false;
@@ -63,7 +76,9 @@ public class PinotResultSet extends AbstractBaseResultSet {
_closed = false;
for (int i = 0; i < _totalColumns; i++) {
_columns.put(_resultSet.getColumnName(i), i + 1);
- _columnDataTypes.put(i + 1, _resultSet.getColumnDataType(i));
+ String columnTypeName = _resultSet.getColumnDataType(i);
+ _columnDataTypes.put(i + 1, columnTypeName);
+ _resolvedColumnDataTypes.put(i + 1, columnTypeName == null ? null :
ColumnDataType.forName(columnTypeName));
}
}
@@ -97,13 +112,6 @@ public class PinotResultSet extends AbstractBaseResultSet {
}
}
- protected void validateState()
- throws SQLException {
- if (isClosed()) {
- throw new SQLException("Not possible to operate on closed or empty
result sets");
- }
- }
-
protected void validateColumn(int columnIndex)
throws SQLException {
validateState();
@@ -172,6 +180,12 @@ public class PinotResultSet extends AbstractBaseResultSet {
return new PinotResultMetadata(_totalColumns, _columns, _columnDataTypes);
}
+ @Nullable
+ @Override
+ protected ColumnDataType getColumnType(int columnIndex) {
+ return _resolvedColumnDataTypes.get(columnIndex);
+ }
+
@Override
public boolean first()
throws SQLException {
@@ -181,101 +195,6 @@ public class PinotResultSet extends AbstractBaseResultSet
{
return true;
}
- @Override
- public InputStream getAsciiStream(int columnIndex)
- throws SQLException {
- String value = getString(columnIndex);
- InputStream in = new
ByteArrayInputStream(value.getBytes(StandardCharsets.US_ASCII));
- return in;
- }
-
- @Override
- public BigDecimal getBigDecimal(int columnIndex, int scale)
- throws SQLException {
- try {
- String value = this.getString(columnIndex);
- int calculatedScale = getCalculatedScale(value);
- return value == null ? null : new
BigDecimal(value).setScale(calculatedScale);
- } catch (Exception e) {
- throw new SQLException("Unable to fetch BigDecimal value", e);
- }
- }
-
- int getCalculatedScale(String value) {
- int index = value.indexOf(".");
- return index == -1 ? 0 : value.length() - index - 1;
- }
-
- @Override
- public boolean getBoolean(int columnIndex)
- throws SQLException {
- validateColumn(columnIndex);
- String value = getString(columnIndex);
- return value == null ? false : Boolean.parseBoolean(value);
- }
-
- @Override
- public byte[] getBytes(int columnIndex)
- throws SQLException {
- try {
- String value = getString(columnIndex);
- return value == null ? null : Hex.decodeHex(value.toCharArray());
- } catch (Exception e) {
- throw new SQLException(String.format("Unable to fetch value for column
%d", columnIndex), e);
- }
- }
-
- @Override
- public Reader getCharacterStream(int columnIndex)
- throws SQLException {
- InputStream in = getUnicodeStream(columnIndex);
- Reader reader = new InputStreamReader(in, StandardCharsets.UTF_8);
- return reader;
- }
-
- @Override
- public Date getDate(int columnIndex, Calendar cal)
- throws SQLException {
- try {
- String value = getString(columnIndex);
- return value == null ? null : DateTimeUtils.getDateFromString(value,
cal);
- } catch (Exception e) {
- throw new SQLException("Unable to fetch date", e);
- }
- }
-
- @Override
- public double getDouble(int columnIndex)
- throws SQLException {
- validateColumn(columnIndex);
- String value = getString(columnIndex);
- return value == null ? 0.0 : Double.parseDouble(value);
- }
-
- @Override
- public float getFloat(int columnIndex)
- throws SQLException {
- validateColumn(columnIndex);
- String value = getString(columnIndex);
- return value == null ? 0.0f : Float.parseFloat(value);
- }
-
- @Override
- public int getInt(int columnIndex)
- throws SQLException {
- validateColumn(columnIndex);
- String value = getString(columnIndex);
- return value == null ? 0 : Integer.parseInt(value);
- }
-
- @Override
- public long getLong(int columnIndex)
- throws SQLException {
- validateColumn(columnIndex);
- String value = getString(columnIndex);
- return value == null ? 0 : Long.parseLong(value);
- }
-
@Override
public int getRow()
throws SQLException {
@@ -284,13 +203,7 @@ public class PinotResultSet extends AbstractBaseResultSet {
return _currentRow;
}
- @Override
- public short getShort(int columnIndex)
- throws SQLException {
- Integer value = getInt(columnIndex);
- return value == null ? null : value.shortValue();
- }
-
+ @Nullable
@Override
public String getString(int columnIndex)
throws SQLException {
@@ -304,101 +217,81 @@ public class PinotResultSet extends
AbstractBaseResultSet {
return val;
}
- @Override
- public Object getObject(int columnIndex)
- throws SQLException {
-
- String dataType = _columnDataTypes.getOrDefault(columnIndex, "");
-
- if (dataType.isEmpty()) {
- throw new SQLDataException("Data type not supported for " + dataType);
- }
-
- switch (dataType) {
- case "STRING":
- return getString(columnIndex);
- case "INT":
- return getInt(columnIndex);
- case "LONG":
- return getLong(columnIndex);
- case "FLOAT":
- return getFloat(columnIndex);
- case "DOUBLE":
- return getDouble(columnIndex);
- case "BOOLEAN":
- return getBoolean(columnIndex);
- case "BYTES":
- return getBytes(columnIndex);
- default:
- throw new SQLDataException("Data type not supported for " + dataType);
+ private boolean checkIsNull(String val) {
+ if (val == null || val.toLowerCase().contentEquals(NULL_STRING)) {
+ _wasNull = true;
+ return true;
}
+ return false;
}
+ @Nullable
@Override
- public <T> T getObject(int columnIndex, Class<T> type)
+ protected Map<?, ?> getMap(int columnIndex)
throws SQLException {
- Object value = getObject(columnIndex);
-
- try {
- return type.cast(value);
- } catch (ClassCastException e) {
- throw new SQLDataException("Data type conversion is not supported from
:" + value.getClass() + " to: " + type);
- }
+ return parseJson(columnIndex, MAP_READER, "map");
}
+ @Nullable
@Override
- public <T> T getObject(String columnLabel, Class<T> type)
+ protected List<?> getList(int columnIndex, ColumnDataType dataType)
throws SQLException {
- return super.getObject(columnLabel, type);
- }
-
- private boolean checkIsNull(String val) {
- if (val == null || val.toLowerCase().contentEquals(NULL_STRING)) {
- _wasNull = true;
- return true;
+ switch (dataType) {
+ case BOOLEAN_ARRAY:
+ return parseJson(columnIndex, BOOLEAN_LIST_READER, "boolean array");
+ case INT_ARRAY:
+ return parseJson(columnIndex, INT_LIST_READER, "int array");
+ case LONG_ARRAY:
+ return parseJson(columnIndex, LONG_LIST_READER, "long array");
+ case FLOAT_ARRAY:
+ return parseJson(columnIndex, FLOAT_LIST_READER, "float array");
+ case DOUBLE_ARRAY:
+ return parseJson(columnIndex, DOUBLE_LIST_READER, "double array");
+ case BIG_DECIMAL_ARRAY:
+ return parseJson(columnIndex, BIG_DECIMAL_LIST_READER, "big decimal
array");
+ case TIMESTAMP_ARRAY:
+ return getTimestampList(columnIndex);
+ case STRING_ARRAY:
+ return parseJson(columnIndex, STRING_LIST_READER, "string array");
+ case BYTES_ARRAY:
+ return getBytesList(columnIndex);
+ case UUID_ARRAY:
+ return getUuidList(columnIndex);
+ default:
+ throw new SQLDataException("Data type is not an array: " + dataType);
}
- return false;
}
- @Override
- public Time getTime(int columnIndex, Calendar cal)
+ @Nullable
+ private List<Timestamp> getTimestampList(int columnIndex)
throws SQLException {
- try {
- String value = getString(columnIndex);
- return value == null ? null : DateTimeUtils.getTimeFromString(value,
cal);
- } catch (Exception e) {
- throw new SQLException("Unable to fetch date", e);
- }
+ List<String> values = parseJson(columnIndex, STRING_LIST_READER,
"timestamp array");
+ return toTimestampList(values);
}
- @Override
- public Timestamp getTimestamp(int columnIndex, Calendar cal)
+ @Nullable
+ private List<byte[]> getBytesList(int columnIndex)
throws SQLException {
- try {
- String value = getString(columnIndex);
- return value == null ? null :
DateTimeUtils.getTimestampFromString(value, cal);
- } catch (Exception e) {
- throw new SQLException("Unable to fetch date", e);
- }
+ List<String> values = parseJson(columnIndex, STRING_LIST_READER, "bytes
array");
+ return toBytesList(values);
}
- @Override
- public URL getURL(int columnIndex)
+ @Nullable
+ private List<UUID> getUuidList(int columnIndex)
throws SQLException {
- try {
- URL url = new URL(getString(columnIndex));
- return url;
- } catch (Exception e) {
- throw new SQLException("Unable to fetch URL", e);
- }
+ List<String> values = parseJson(columnIndex, STRING_LIST_READER, "UUID
array");
+ return toUuidList(values);
}
- @Override
- public InputStream getUnicodeStream(int columnIndex)
+ @Nullable
+ private <T> T parseJson(int columnIndex, ObjectReader reader, String type)
throws SQLException {
- String value = getString(columnIndex);
- InputStream in = new
ByteArrayInputStream(value.getBytes(StandardCharsets.UTF_8));
- return in;
+ try {
+ String stringVal = getString(columnIndex);
+ return (stringVal == null) ? null : reader.readValue(stringVal);
+ } catch (IOException e) {
+ throw new SQLDataException("Error parsing " + type, e);
+ }
}
@Override
diff --git
a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/base/AbstractBaseResultSet.java
b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/base/AbstractBaseResultSet.java
index 8370bea07b9..d5eefd0ef54 100644
---
a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/base/AbstractBaseResultSet.java
+++
b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/base/AbstractBaseResultSet.java
@@ -18,10 +18,13 @@
*/
package org.apache.pinot.client.base;
+import java.io.ByteArrayInputStream;
import java.io.InputStream;
+import java.io.InputStreamReader;
import java.io.Reader;
import java.math.BigDecimal;
import java.net.URL;
+import java.nio.charset.StandardCharsets;
import java.sql.Array;
import java.sql.Blob;
import java.sql.Clob;
@@ -30,6 +33,7 @@ import java.sql.NClob;
import java.sql.Ref;
import java.sql.ResultSet;
import java.sql.RowId;
+import java.sql.SQLDataException;
import java.sql.SQLException;
import java.sql.SQLFeatureNotSupportedException;
import java.sql.SQLType;
@@ -38,18 +42,49 @@ import java.sql.SQLXML;
import java.sql.Statement;
import java.sql.Time;
import java.sql.Timestamp;
+import java.text.ParseException;
+import java.util.ArrayList;
import java.util.Calendar;
+import java.util.List;
import java.util.Map;
+import java.util.UUID;
+import javax.annotation.Nullable;
+import org.apache.commons.codec.DecoderException;
+import org.apache.commons.codec.binary.Hex;
+import org.apache.pinot.client.utils.BigDecimalUtils;
+import org.apache.pinot.client.utils.DateTimeUtils;
+import org.apache.pinot.common.utils.DataSchema.ColumnDataType;
public abstract class AbstractBaseResultSet implements ResultSet {
- protected abstract void validateState()
- throws SQLException;
-
protected abstract void validateColumn(int columnIndex)
throws SQLException;
+ @Nullable
+ protected Map<?, ?> getMap(int columnIndex)
+ throws SQLException {
+ throw new SQLFeatureNotSupportedException("Map is not supported by the
ResultSet");
+ }
+
+ @Nullable
+ protected List<?> getList(int columnIndex, ColumnDataType dataType)
+ throws SQLException {
+ throw new SQLFeatureNotSupportedException("Array is not supported by the
ResultSet");
+ }
+
+ @Nullable
+ protected ColumnDataType getColumnType(int columnIndex) throws SQLException {
+ return null;
+ }
+
+ protected void validateState()
+ throws SQLException {
+ if (isClosed()) {
+ throw new SQLException("Not possible to operate on closed or empty
result sets");
+ }
+ }
+
@Override
public void cancelRowUpdates()
throws SQLException {
@@ -86,18 +121,40 @@ public abstract class AbstractBaseResultSet implements
ResultSet {
return getAsciiStream(findColumn(columnLabel));
}
+ @Override
+ public InputStream getAsciiStream(int columnIndex)
+ throws SQLException {
+ String value = getString(columnIndex);
+ InputStream in = value == null ? null : new
ByteArrayInputStream(value.getBytes(StandardCharsets.US_ASCII));
+ return in;
+ }
+
@Override
public BigDecimal getBigDecimal(int columnIndex)
throws SQLException {
return getBigDecimal(columnIndex, 0);
}
+ @Nullable
@Override
public BigDecimal getBigDecimal(String columnLabel, int scale)
throws SQLException {
return getBigDecimal(findColumn(columnLabel), scale);
}
+ @Nullable
+ @Override
+ public BigDecimal getBigDecimal(int columnIndex, int scale)
+ throws SQLException {
+ try {
+ String value = getString(columnIndex);
+ return BigDecimalUtils.getBigDecimalFromString(value);
+ } catch (Exception e) {
+ throw new SQLDataException("Unable to fetch BigDecimal value", e);
+ }
+ }
+
+ @Nullable
@Override
public BigDecimal getBigDecimal(String columnLabel)
throws SQLException {
@@ -122,6 +179,14 @@ public abstract class AbstractBaseResultSet implements
ResultSet {
return getBoolean(findColumn(columnLabel));
}
+ @Override
+ public boolean getBoolean(int columnIndex)
+ throws SQLException {
+ validateColumn(columnIndex);
+ String value = getString(columnIndex);
+ return value == null ? false : Boolean.parseBoolean(value);
+ }
+
@Override
public Blob getBlob(int columnIndex)
throws SQLException {
@@ -146,18 +211,41 @@ public abstract class AbstractBaseResultSet implements
ResultSet {
throw new SQLFeatureNotSupportedException();
}
+ @Nullable
@Override
public byte[] getBytes(String columnLabel)
throws SQLException {
return getBytes(findColumn(columnLabel));
}
+ @Nullable
+ @Override
+ public byte[] getBytes(int columnIndex)
+ throws SQLException {
+ try {
+ String value = getString(columnIndex);
+ return value == null ? null : Hex.decodeHex(value.toCharArray());
+ } catch (Exception e) {
+ throw new SQLDataException(String.format("Unable to fetch value for
column %d", columnIndex), e);
+ }
+ }
+
+ @Nullable
@Override
public Reader getCharacterStream(String columnLabel)
throws SQLException {
return getCharacterStream(findColumn(columnLabel));
}
+ @Nullable
+ @Override
+ public Reader getCharacterStream(int columnIndex)
+ throws SQLException {
+ InputStream in = getUnicodeStream(columnIndex);
+ Reader reader = in == null ? null : new InputStreamReader(in,
StandardCharsets.UTF_8);
+ return reader;
+ }
+
@Override
public Clob getClob(int columnIndex)
throws SQLException {
@@ -183,30 +271,53 @@ public abstract class AbstractBaseResultSet implements
ResultSet {
throw new SQLFeatureNotSupportedException();
}
+ @Nullable
@Override
public Date getDate(int columnIndex)
throws SQLException {
return getDate(columnIndex, Calendar.getInstance());
}
+ @Nullable
@Override
public Date getDate(String columnLabel)
throws SQLException {
return getDate(findColumn(columnLabel), Calendar.getInstance());
}
+ @Nullable
@Override
public Date getDate(String columnLabel, Calendar cal)
throws SQLException {
return getDate(findColumn(columnLabel), cal);
}
+ @Nullable
+ @Override
+ public Date getDate(int columnIndex, Calendar cal)
+ throws SQLException {
+ try {
+ String value = getString(columnIndex);
+ return value == null ? null : DateTimeUtils.getDateFromString(value,
cal);
+ } catch (Exception e) {
+ throw new SQLException("Unable to fetch date", e);
+ }
+ }
+
@Override
public double getDouble(String columnLabel)
throws SQLException {
return getDouble(findColumn(columnLabel));
}
+ @Override
+ public double getDouble(int columnIndex)
+ throws SQLException {
+ validateColumn(columnIndex);
+ String value = getString(columnIndex);
+ return value == null ? 0.0 : Double.parseDouble(value);
+ }
+
@Override
public int getFetchDirection()
throws SQLException {
@@ -238,6 +349,14 @@ public abstract class AbstractBaseResultSet implements
ResultSet {
return getFloat(findColumn(columnLabel));
}
+ @Override
+ public float getFloat(int columnIndex)
+ throws SQLException {
+ validateColumn(columnIndex);
+ String value = getString(columnIndex);
+ return value == null ? 0.0f : Float.parseFloat(value);
+ }
+
@Override
public int getHoldability()
throws SQLException {
@@ -251,16 +370,38 @@ public abstract class AbstractBaseResultSet implements
ResultSet {
return getInt(findColumn(columnLabel));
}
+ @Override
+ public int getInt(int columnIndex)
+ throws SQLException {
+ validateColumn(columnIndex);
+ String value = getString(columnIndex);
+ return value == null ? 0 : Integer.parseInt(value);
+ }
+
@Override
public long getLong(String columnLabel)
throws SQLException {
return getLong(findColumn(columnLabel));
}
+ @Override
+ public long getLong(int columnIndex)
+ throws SQLException {
+ validateColumn(columnIndex);
+ String value = getString(columnIndex);
+ return value == null ? 0 : Long.parseLong(value);
+ }
+
@Override
public short getShort(String columnLabel)
throws SQLException {
- return getShort(columnLabel);
+ return getShort(findColumn(columnLabel));
+ }
+
+ @Override
+ public short getShort(int columnIndex)
+ throws SQLException {
+ return (short) getInt(columnIndex);
}
@Override
@@ -305,30 +446,47 @@ public abstract class AbstractBaseResultSet implements
ResultSet {
throw new SQLFeatureNotSupportedException();
}
+ @Nullable
@Override
public Object getObject(String columnLabel)
throws SQLException {
return getObject(findColumn(columnLabel));
}
+ @Nullable
@Override
public Object getObject(int columnIndex, Map<String, Class<?>> map)
throws SQLException {
return getObject(columnIndex);
}
+ @Nullable
@Override
public Object getObject(String columnLabel, Map<String, Class<?>> map)
throws SQLException {
return getObject(findColumn(columnLabel), map);
}
+ @Nullable
@Override
public <T> T getObject(String columnLabel, Class<T> type)
throws SQLException {
return getObject(findColumn(columnLabel), type);
}
+ @Nullable
+ @Override
+ public <T> T getObject(int columnIndex, Class<T> type)
+ throws SQLException {
+ Object value = getObject(columnIndex);
+
+ try {
+ return type.cast(value);
+ } catch (ClassCastException e) {
+ throw new SQLDataException("Data type conversion is not supported from
:" + value.getClass() + " to: " + type);
+ }
+ }
+
@Override
public Ref getRef(int columnIndex)
throws SQLException {
@@ -365,48 +523,79 @@ public abstract class AbstractBaseResultSet implements
ResultSet {
throw new SQLFeatureNotSupportedException();
}
+ @Nullable
@Override
public Statement getStatement()
throws SQLException {
return null;
}
+ @Nullable
@Override
public Time getTime(int columnIndex)
throws SQLException {
return getTime(columnIndex, Calendar.getInstance());
}
+ @Nullable
@Override
public Time getTime(String columnLabel)
throws SQLException {
return getTime(findColumn(columnLabel), Calendar.getInstance());
}
+ @Nullable
@Override
public Time getTime(String columnLabel, Calendar cal)
throws SQLException {
return getTime(findColumn(columnLabel), cal);
}
+ @Nullable
+ @Override
+ public Time getTime(int columnIndex, Calendar cal)
+ throws SQLException {
+ try {
+ String value = getString(columnIndex);
+ return value == null ? null : DateTimeUtils.getTimeFromString(value,
cal);
+ } catch (Exception e) {
+ throw new SQLException("Unable to fetch date", e);
+ }
+ }
+
+ @Nullable
@Override
public Timestamp getTimestamp(int columnIndex)
throws SQLException {
return getTimestamp(columnIndex, Calendar.getInstance());
}
+ @Nullable
@Override
public Timestamp getTimestamp(String columnLabel)
throws SQLException {
return getTimestamp(findColumn(columnLabel), Calendar.getInstance());
}
+ @Nullable
@Override
public Timestamp getTimestamp(String columnLabel, Calendar cal)
throws SQLException {
return getTimestamp(findColumn(columnLabel), cal);
}
+ @Nullable
+ @Override
+ public Timestamp getTimestamp(int columnIndex, Calendar cal)
+ throws SQLException {
+ try {
+ String value = getString(columnIndex);
+ return value == null ? null :
DateTimeUtils.getTimestampFromString(value, cal);
+ } catch (Exception e) {
+ throw new SQLDataException("Unable to fetch date", e);
+ }
+ }
+
@Override
public int getType()
throws SQLException {
@@ -414,18 +603,53 @@ public abstract class AbstractBaseResultSet implements
ResultSet {
return ResultSet.TYPE_FORWARD_ONLY;
}
+ @Nullable
@Override
public URL getURL(String columnLabel)
throws SQLException {
return getURL(findColumn(columnLabel));
}
+ @Nullable
+ @Override
+ public URL getURL(int columnIndex)
+ throws SQLException {
+ try {
+ String value = getString(columnIndex);
+ URL url = value == null ? null : new URL(value);
+ return url;
+ } catch (Exception e) {
+ throw new SQLException("Unable to fetch URL", e);
+ }
+ }
+
+ @Nullable
@Override
public InputStream getUnicodeStream(String columnLabel)
throws SQLException {
return getUnicodeStream(findColumn(columnLabel));
}
+ @Nullable
+ @Override
+ public InputStream getUnicodeStream(int columnIndex)
+ throws SQLException {
+ String value = getString(columnIndex);
+ InputStream in = value == null ? null : new
ByteArrayInputStream(value.getBytes(StandardCharsets.UTF_8));
+ return in;
+ }
+
+ @Nullable
+ private UUID getUuid(int columnIndex)
+ throws SQLException {
+ String value = getString(columnIndex);
+ try {
+ return value == null ? null : UUID.fromString(value);
+ } catch (IllegalArgumentException e) {
+ throw new SQLDataException("Error parsing UUID", e);
+ }
+ }
+
@Override
public SQLWarning getWarnings()
throws SQLException {
@@ -1007,4 +1231,114 @@ public abstract class AbstractBaseResultSet implements
ResultSet {
throws SQLException {
throw new SQLFeatureNotSupportedException();
}
+
+ @Nullable
+ @Override
+ public Object getObject(int columnIndex)
+ throws SQLException {
+
+ ColumnDataType dataType = getColumnType(columnIndex);
+ if (dataType == null) {
+ throw new SQLDataException("Data type not supported for column " +
columnIndex);
+ }
+
+ Object value = getObject(columnIndex, dataType);
+ if (wasNull()) {
+ return null;
+ }
+ return value;
+ }
+
+ @Nullable
+ private Object getObject(int columnIndex, ColumnDataType dataType)
+ throws SQLException {
+ if (dataType.isArray()) {
+ return getList(columnIndex, dataType);
+ }
+ return getScalar(columnIndex, dataType);
+ }
+
+ @Nullable
+ private Object getScalar(int columnIndex, ColumnDataType dataType)
+ throws SQLException {
+ switch (dataType) {
+ case STRING:
+ case JSON:
+ return getString(columnIndex);
+ case INT:
+ return getInt(columnIndex);
+ case LONG:
+ return getLong(columnIndex);
+ case FLOAT:
+ return getFloat(columnIndex);
+ case DOUBLE:
+ return getDouble(columnIndex);
+ case BIG_DECIMAL:
+ return getBigDecimal(columnIndex);
+ case BOOLEAN:
+ return getBoolean(columnIndex);
+ case BYTES:
+ return getBytes(columnIndex);
+ case UUID:
+ return getUuid(columnIndex);
+ case TIMESTAMP:
+ return getTimestamp(columnIndex);
+ case MAP:
+ return getMap(columnIndex);
+ default:
+ throw new SQLDataException("Data type not supported for " + dataType);
+ }
+ }
+
+ @Nullable
+ protected List<Timestamp> toTimestampList(List<String> values)
+ throws SQLException {
+ if (values == null) {
+ return null;
+ }
+ List<Timestamp> list = new ArrayList<>(values.size());
+ Calendar calendar = Calendar.getInstance();
+ try {
+ for (String value : values) {
+ list.add(value == null ? null :
DateTimeUtils.getTimestampFromString(value, calendar));
+ }
+ return list;
+ } catch (ParseException e) {
+ throw new SQLDataException("Error parsing timestamp array", e);
+ }
+ }
+
+ @Nullable
+ protected List<byte[]> toBytesList(List<String> values)
+ throws SQLException {
+ if (values == null) {
+ return null;
+ }
+ List<byte[]> bytes = new ArrayList<>(values.size());
+ try {
+ for (String value : values) {
+ bytes.add(value == null ? null : Hex.decodeHex(value));
+ }
+ return bytes;
+ } catch (DecoderException e) {
+ throw new SQLDataException("Error parsing bytes array", e);
+ }
+ }
+
+ @Nullable
+ protected List<UUID> toUuidList(List<String> values)
+ throws SQLException {
+ if (values == null) {
+ return null;
+ }
+ List<UUID> uuids = new ArrayList<>(values.size());
+ try {
+ for (String value : values) {
+ uuids.add(value == null ? null : UUID.fromString(value));
+ }
+ return uuids;
+ } catch (IllegalArgumentException e) {
+ throw new SQLDataException("Error parsing UUID array", e);
+ }
+ }
}
diff --git
a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/grpc/PinotGrpcResultSet.java
b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/grpc/PinotGrpcResultSet.java
index a47b3fcfc44..b88a65b1fca 100644
---
a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/grpc/PinotGrpcResultSet.java
+++
b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/grpc/PinotGrpcResultSet.java
@@ -19,31 +19,24 @@
package org.apache.pinot.client.grpc;
import com.fasterxml.jackson.databind.node.ObjectNode;
-import java.io.ByteArrayInputStream;
import java.io.IOException;
-import java.io.InputStream;
-import java.io.InputStreamReader;
-import java.io.Reader;
import java.math.BigDecimal;
-import java.net.URL;
-import java.nio.charset.StandardCharsets;
-import java.sql.Date;
import java.sql.ResultSetMetaData;
import java.sql.SQLDataException;
import java.sql.SQLException;
-import java.sql.Time;
-import java.sql.Timestamp;
-import java.util.Calendar;
+import java.util.ArrayList;
+import java.util.Arrays;
import java.util.HashMap;
import java.util.Iterator;
+import java.util.List;
import java.util.Map;
-import org.apache.commons.codec.binary.Hex;
+import javax.annotation.Nullable;
import org.apache.pinot.client.PinotResultMetadata;
import org.apache.pinot.client.base.AbstractBaseResultSet;
-import org.apache.pinot.client.utils.DateTimeUtils;
import org.apache.pinot.common.proto.Broker;
import org.apache.pinot.common.response.broker.ResultTable;
import org.apache.pinot.common.utils.DataSchema;
+import org.apache.pinot.common.utils.DataSchema.ColumnDataType;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -90,13 +83,6 @@ public class PinotGrpcResultSet extends
AbstractBaseResultSet {
return new PinotGrpcResultSet();
}
- protected void validateState()
- throws SQLException {
- if (isClosed()) {
- throw new SQLException("Not possible to operate on closed or empty
result sets");
- }
- }
-
protected void validateColumn(int columnIndex)
throws SQLException {
validateState();
@@ -152,6 +138,13 @@ public class PinotGrpcResultSet extends
AbstractBaseResultSet {
return new PinotResultMetadata(_totalColumns, _columns, _columnDataTypes);
}
+ @Nullable
+ @Override
+ protected ColumnDataType getColumnType(int columnIndex) throws SQLException {
+ validateColumn(columnIndex);
+ return _dataSchema.getColumnDataType(columnIndex - 1);
+ }
+
@Override
public boolean first()
throws SQLException {
@@ -160,173 +153,140 @@ public class PinotGrpcResultSet extends
AbstractBaseResultSet {
}
@Override
- public InputStream getAsciiStream(int columnIndex)
+ public int getRow()
throws SQLException {
- String value = getString(columnIndex);
- InputStream in = new
ByteArrayInputStream(value.getBytes(StandardCharsets.US_ASCII));
- return in;
+ validateState();
+ return _currentRow;
}
+ @Nullable
@Override
- public BigDecimal getBigDecimal(int columnIndex, int scale)
+ public String getString(int columnIndex)
throws SQLException {
- try {
- String value = this.getString(columnIndex);
- int calculatedScale = getCalculatedScale(value);
- return value == null ? null : new
BigDecimal(value).setScale(calculatedScale);
- } catch (Exception e) {
- throw new SQLException("Unable to fetch BigDecimal value", e);
+ Object value = getValue(columnIndex);
+ String val = value == null ? null : value.toString();
+ if (checkIsNull(val)) {
+ return null;
}
+ return val;
}
- int getCalculatedScale(String value) {
- int index = value.indexOf(".");
- return index == -1 ? 0 : value.length() - index - 1;
- }
-
- @Override
- public boolean getBoolean(int columnIndex)
+ @Nullable
+ private Object getValue(int columnIndex)
throws SQLException {
validateColumn(columnIndex);
- String value = getString(columnIndex);
- return value == null ? false : Boolean.parseBoolean(value);
- }
-
- @Override
- public byte[] getBytes(int columnIndex)
- throws SQLException {
- try {
- String value = getString(columnIndex);
- return value == null ? null : Hex.decodeHex(value.toCharArray());
- } catch (Exception e) {
- throw new SQLException(String.format("Unable to fetch value for column
%d", columnIndex), e);
+ Object value =
_currentRowBatch.getRows().get(_currentBatchIndex)[columnIndex - 1];
+ if (value == null) {
+ _wasNull = true;
}
+ return value;
}
+ @Nullable
@Override
- public Reader getCharacterStream(int columnIndex)
+ protected Map<?, ?> getMap(int columnIndex)
throws SQLException {
- InputStream in = getUnicodeStream(columnIndex);
- Reader reader = new InputStreamReader(in, StandardCharsets.UTF_8);
- return reader;
+ Object value = getValue(columnIndex);
+ if (value == null) {
+ return null;
+ }
+ if (!(value instanceof Map)) {
+ throw new SQLDataException("Expected map value, found: " +
value.getClass());
+ }
+ return (Map<?, ?>) value;
}
+ @Nullable
@Override
- public Date getDate(int columnIndex, Calendar cal)
+ protected List<?> getList(int columnIndex, ColumnDataType dataType)
throws SQLException {
+ Object value = getValue(columnIndex);
+ if (value == null) {
+ return null;
+ }
try {
- String value = getString(columnIndex);
- return value == null ? null : DateTimeUtils.getDateFromString(value,
cal);
- } catch (Exception e) {
- throw new SQLException("Unable to fetch date", e);
+ switch (dataType) {
+ case BOOLEAN_ARRAY:
+ return toList((boolean[]) value);
+ case INT_ARRAY:
+ return toList((int[]) value);
+ case LONG_ARRAY:
+ return toList((long[]) value);
+ case FLOAT_ARRAY:
+ return toList((float[]) value);
+ case DOUBLE_ARRAY:
+ return toList((double[]) value);
+ case BIG_DECIMAL_ARRAY:
+ return toBigDecimalList((String[]) value);
+ case TIMESTAMP_ARRAY:
+ return toTimestampList(Arrays.asList((String[]) value));
+ case STRING_ARRAY:
+ return new ArrayList<>(Arrays.asList((String[]) value));
+ case BYTES_ARRAY:
+ return toBytesList(Arrays.asList((String[]) value));
+ case UUID_ARRAY:
+ return toUuidList(Arrays.asList((String[]) value));
+ default:
+ throw new SQLDataException("Data type is not an array: " + dataType);
+ }
+ } catch (ClassCastException e) {
+ throw new SQLDataException("Unexpected value type for " + dataType + ":
" + value.getClass(), e);
}
}
- @Override
- public double getDouble(int columnIndex)
- throws SQLException {
- validateColumn(columnIndex);
- String value = getString(columnIndex);
- return value == null ? 0.0 : Double.parseDouble(value);
- }
-
- @Override
- public float getFloat(int columnIndex)
- throws SQLException {
- validateColumn(columnIndex);
- String value = getString(columnIndex);
- return value == null ? 0.0f : Float.parseFloat(value);
- }
-
- @Override
- public int getInt(int columnIndex)
- throws SQLException {
- validateColumn(columnIndex);
- String value = getString(columnIndex);
- return value == null ? 0 : Integer.parseInt(value);
- }
-
- @Override
- public long getLong(int columnIndex)
- throws SQLException {
- validateColumn(columnIndex);
- String value = getString(columnIndex);
- return value == null ? 0 : Long.parseLong(value);
- }
-
- @Override
- public int getRow()
- throws SQLException {
- validateState();
- return _currentRow;
+ private static List<Boolean> toList(boolean[] values) {
+ List<Boolean> list = new ArrayList<>(values.length);
+ for (boolean value : values) {
+ list.add(value);
+ }
+ return list;
}
- @Override
- public short getShort(int columnIndex)
- throws SQLException {
- Integer value = getInt(columnIndex);
- return value == null ? null : value.shortValue();
+ private static List<Integer> toList(int[] values) {
+ List<Integer> list = new ArrayList<>(values.length);
+ for (int value : values) {
+ list.add(value);
+ }
+ return list;
}
- @Override
- public String getString(int columnIndex)
- throws SQLException {
- validateColumn(columnIndex);
- String val =
_currentRowBatch.getRows().get(_currentBatchIndex)[columnIndex - 1].toString();
- if (checkIsNull(val)) {
- return null;
+ private static List<Long> toList(long[] values) {
+ List<Long> list = new ArrayList<>(values.length);
+ for (long value : values) {
+ list.add(value);
}
- return val;
+ return list;
}
- @Override
- public Object getObject(int columnIndex)
- throws SQLException {
-
- String dataType = _columnDataTypes.getOrDefault(columnIndex, "");
-
- if (dataType.isEmpty()) {
- throw new SQLDataException("Data type not supported for " + dataType);
+ private static List<Float> toList(float[] values) {
+ List<Float> list = new ArrayList<>(values.length);
+ for (float value : values) {
+ list.add(value);
}
+ return list;
+ }
- switch (dataType) {
- case "STRING":
- return getString(columnIndex);
- case "INT":
- return getInt(columnIndex);
- case "LONG":
- return getLong(columnIndex);
- case "FLOAT":
- return getFloat(columnIndex);
- case "DOUBLE":
- return getDouble(columnIndex);
- case "BOOLEAN":
- return getBoolean(columnIndex);
- case "BYTES":
- return getBytes(columnIndex);
- default:
- throw new SQLDataException("Data type not supported for " + dataType);
+ private static List<Double> toList(double[] values) {
+ List<Double> list = new ArrayList<>(values.length);
+ for (double value : values) {
+ list.add(value);
}
+ return list;
}
- @Override
- public <T> T getObject(int columnIndex, Class<T> type)
+ private static List<BigDecimal> toBigDecimalList(String[] values)
throws SQLException {
- Object value = getObject(columnIndex);
-
+ List<BigDecimal> list = new ArrayList<>(values.length);
try {
- return type.cast(value);
- } catch (ClassCastException e) {
- throw new SQLDataException("Data type conversion is not supported from
:" + value.getClass() + " to: " + type);
+ for (String value : values) {
+ list.add(value == null ? null : new BigDecimal(value));
+ }
+ return list;
+ } catch (NumberFormatException e) {
+ throw new SQLDataException("Error parsing big decimal array", e);
}
}
- @Override
- public <T> T getObject(String columnLabel, Class<T> type)
- throws SQLException {
- return super.getObject(columnLabel, type);
- }
-
private boolean checkIsNull(String val) {
if (val == null || val.toLowerCase().contentEquals(NULL_STRING)) {
_wasNull = true;
@@ -335,47 +295,6 @@ public class PinotGrpcResultSet extends
AbstractBaseResultSet {
return false;
}
- @Override
- public Time getTime(int columnIndex, Calendar cal)
- throws SQLException {
- try {
- String value = getString(columnIndex);
- return value == null ? null : DateTimeUtils.getTimeFromString(value,
cal);
- } catch (Exception e) {
- throw new SQLException("Unable to fetch date", e);
- }
- }
-
- @Override
- public Timestamp getTimestamp(int columnIndex, Calendar cal)
- throws SQLException {
- try {
- String value = getString(columnIndex);
- return value == null ? null :
DateTimeUtils.getTimestampFromString(value, cal);
- } catch (Exception e) {
- throw new SQLException("Unable to fetch date", e);
- }
- }
-
- @Override
- public URL getURL(int columnIndex)
- throws SQLException {
- try {
- URL url = new URL(getString(columnIndex));
- return url;
- } catch (Exception e) {
- throw new SQLException("Unable to fetch URL", e);
- }
- }
-
- @Override
- public InputStream getUnicodeStream(int columnIndex)
- throws SQLException {
- String value = getString(columnIndex);
- InputStream in = new
ByteArrayInputStream(value.getBytes(StandardCharsets.UTF_8));
- return in;
- }
-
@Override
public boolean isAfterLast()
throws SQLException {
diff --git
a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/utils/BigDecimalUtils.java
b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/utils/BigDecimalUtils.java
new file mode 100644
index 00000000000..b059e3dd232
--- /dev/null
+++
b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/utils/BigDecimalUtils.java
@@ -0,0 +1,36 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.client.utils;
+
+import java.math.BigDecimal;
+
+public class BigDecimalUtils {
+
+ private BigDecimalUtils() {
+ }
+
+ public static BigDecimal getBigDecimalFromString(String value) {
+ return value == null ? null : new
BigDecimal(value).setScale(getCalculatedScale(value));
+ }
+
+ static int getCalculatedScale(String value) {
+ int index = value.indexOf(".");
+ return index == -1 ? 0 : value.length() - index - 1;
+ }
+}
diff --git
a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/utils/DateTimeUtils.java
b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/utils/DateTimeUtils.java
index 7e9b4df1523..4849c6e48c1 100644
---
a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/utils/DateTimeUtils.java
+++
b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/utils/DateTimeUtils.java
@@ -57,8 +57,20 @@ public class DateTimeUtils {
throws ParseException {
SimpleDateFormat timestampFormat = TIMESTAMP_FORMAT.get();
timestampFormat.setTimeZone(cal.getTimeZone());
- java.util.Date date = timestampFormat.parse(value);
- return new Timestamp(date.getTime());
+ // Response timestamps may include variable-width fractional seconds;
parse them separately so the shared
+ // seconds-only formatter remains compatible with time parsing and
formatting.
+ int decimalIndex = value.indexOf('.');
+ String timestampValue = decimalIndex >= 0 ? value.substring(0,
decimalIndex) : value;
+ java.util.Date date = timestampFormat.parse(timestampValue);
+ Timestamp timestamp = new Timestamp(date.getTime());
+ if (decimalIndex >= 0) {
+ String fractionalSeconds = value.substring(decimalIndex + 1);
+ if (fractionalSeconds.isEmpty() || fractionalSeconds.length() > 9 ||
!fractionalSeconds.matches("\\d+")) {
+ throw new ParseException("Invalid fractional seconds in timestamp: " +
value, decimalIndex + 1);
+ }
+ timestamp.setNanos(Integer.parseInt((fractionalSeconds +
"000000000").substring(0, 9)));
+ }
+ return timestamp;
}
public static Timestamp getTimestampFromLong(Long value) {
diff --git
a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/utils/DriverUtils.java
b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/utils/DriverUtils.java
index 2e99c58bd17..51b7160d7b3 100644
---
a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/utils/DriverUtils.java
+++
b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/utils/DriverUtils.java
@@ -28,6 +28,7 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Properties;
+import java.util.UUID;
import java.util.regex.Pattern;
import javax.net.ssl.SSLContext;
import org.apache.commons.configuration2.MapConfiguration;
@@ -36,6 +37,7 @@ import org.apache.hc.core5.http.NameValuePair;
import org.apache.hc.core5.net.URLEncodedUtils;
import org.apache.pinot.common.auth.BasicAuthTokenUtils;
import org.apache.pinot.common.config.TlsConfig;
+import org.apache.pinot.common.utils.DataSchema.ColumnDataType;
import org.apache.pinot.common.utils.tls.TlsUtils;
import org.apache.pinot.spi.env.PinotConfiguration;
import org.slf4j.Logger;
@@ -152,6 +154,12 @@ public class DriverUtils {
if (columnDataType == null) {
return Types.VARCHAR;
}
+ if (ColumnDataType.isArray(columnDataType)) {
+ return Types.JAVA_OBJECT;
+ }
+ if (ColumnDataType.MAP.name().equals(columnDataType)) {
+ return Types.JAVA_OBJECT;
+ }
Integer columnsSQLDataType;
switch (columnDataType) {
case "STRING":
@@ -181,6 +189,9 @@ public class DriverUtils {
case "TIMESTAMP":
columnsSQLDataType = Types.TIMESTAMP;
break;
+ case "UUID":
+ columnsSQLDataType = Types.OTHER;
+ break;
default:
columnsSQLDataType = Types.NULL;
break;
@@ -192,6 +203,12 @@ public class DriverUtils {
if (columnDataType == null) {
return String.class.getTypeName();
}
+ if (ColumnDataType.isArray(columnDataType)) {
+ return List.class.getTypeName();
+ }
+ if (ColumnDataType.MAP.name().equals(columnDataType)) {
+ return Map.class.getTypeName();
+ }
String columnsJavaClassName;
switch (columnDataType) {
case "STRING":
@@ -221,6 +238,9 @@ public class DriverUtils {
case "TIMESTAMP":
columnsJavaClassName = Timestamp.class.getTypeName();
break;
+ case "UUID":
+ columnsJavaClassName = UUID.class.getTypeName();
+ break;
default:
columnsJavaClassName = String.class.getTypeName();
break;
diff --git
a/pinot-clients/pinot-jdbc-client/src/test/java/org/apache/pinot/client/PinotResultSetTest.java
b/pinot-clients/pinot-jdbc-client/src/test/java/org/apache/pinot/client/PinotResultSetTest.java
index 144256a7fca..9965bf7e8e3 100644
---
a/pinot-clients/pinot-jdbc-client/src/test/java/org/apache/pinot/client/PinotResultSetTest.java
+++
b/pinot-clients/pinot-jdbc-client/src/test/java/org/apache/pinot/client/PinotResultSetTest.java
@@ -22,10 +22,16 @@ import java.io.InputStream;
import java.math.BigDecimal;
import java.nio.charset.StandardCharsets;
import java.sql.ResultSetMetaData;
+import java.sql.SQLDataException;
import java.sql.SQLException;
+import java.sql.Timestamp;
+import java.sql.Types;
+import java.util.Arrays;
import java.util.Calendar;
import java.util.Date;
import java.util.List;
+import java.util.Map;
+import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
@@ -37,6 +43,9 @@ import org.apache.pinot.spi.utils.JsonUtils;
import org.testng.Assert;
import org.testng.annotations.Test;
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertTrue;
+
/// Tests deserialization of a ResultSet given hardcoded Pinot results.
public class PinotResultSetTest {
@@ -53,7 +62,9 @@ public class PinotResultSetTest {
int currentRow = 0;
while (pinotResultSet.next()) {
Assert.assertEquals(pinotResultSet.getInt(1),
resultSet.getInt(currentRow, 0));
+ Assert.assertEquals(pinotResultSet.getShort(1),
resultSet.getInt(currentRow, 0));
Assert.assertEquals(pinotResultSet.getLong(2),
resultSet.getLong(currentRow, 1));
+ Assert.assertEquals(pinotResultSet.getShort(2), (short)
resultSet.getLong(currentRow, 1));
Assert.assertEquals(pinotResultSet.getFloat(3),
resultSet.getFloat(currentRow, 2));
Assert.assertEquals(pinotResultSet.getDouble(4),
resultSet.getDouble(currentRow, 3));
Assert.assertEquals(pinotResultSet.getString(5),
resultSet.getString(currentRow, 4));
@@ -61,6 +72,135 @@ public class PinotResultSetTest {
}
}
+ @Test
+ public void testFetchValuesWithoutColumnTypes()
+ throws Exception {
+ PinotResultSet resultSet = new PinotResultSet(new AggregationResultSet(
+
JsonUtils.stringToJsonNode("{\"function\":\"sum(value)\",\"value\":\"42\"}")));
+
+ assertTrue(resultSet.next());
+ assertEquals(resultSet.getInt(1), 42);
+ assertEquals(resultSet.getString(1), "42");
+ assertEquals(resultSet.getMetaData().getColumnType(1), Types.VARCHAR);
+ }
+
+ @Test
+ public void testGetArraysAsTypedLists()
+ throws Exception {
+ PinotResultSet resultSet = PinotResultSet.fromJson("{\"resultTable\":{"
+ +
"\"dataSchema\":{\"columnNames\":[\"booleans\",\"ints\",\"longs\",\"floats\",\"doubles\","
+ + "\"decimals\",\"timestamps\",\"strings\",\"bytes\",\"uuids\"],"
+ +
"\"columnDataTypes\":[\"BOOLEAN_ARRAY\",\"INT_ARRAY\",\"LONG_ARRAY\",\"FLOAT_ARRAY\","
+ +
"\"DOUBLE_ARRAY\",\"BIG_DECIMAL_ARRAY\",\"TIMESTAMP_ARRAY\",\"STRING_ARRAY\",\"BYTES_ARRAY\","
+ + "\"UUID_ARRAY\"]},"
+ +
"\"rows\":[[[true,false],[1,null,2],[2147483648,3],[1.25,2.5],[1.5,2.75],"
+ + "[\"1.20\",\"3.4\"],[\"2020-01-01 12:00:00\",\"2021-02-03
04:05:06.123\"],"
+ + "[\"first\",\"second\"],[\"00ff\",\"1020\"],"
+ +
"[\"00000000-0000-0000-0000-000000000001\",\"00000000-0000-0000-0000-000000000002\"]]]}}");
+
+ Assert.assertTrue(resultSet.next());
+ Assert.assertEquals(resultSet.getObject(1), List.of(true, false));
+ Assert.assertEquals(resultSet.getObject(2), Arrays.asList(1, null, 2));
+ Assert.assertEquals(resultSet.getObject(3), List.of(2147483648L, 3L));
+ Assert.assertEquals(resultSet.getObject(4), List.of(1.25f, 2.5f));
+ Assert.assertEquals(resultSet.getObject(5), List.of(1.5, 2.75));
+ Assert.assertEquals(resultSet.getObject(6), List.of(new
BigDecimal("1.20"), new BigDecimal("3.4")));
+ Assert.assertEquals(resultSet.getObject(7),
+ List.of(Timestamp.valueOf("2020-01-01 12:00:00"),
Timestamp.valueOf("2021-02-03 04:05:06.123")));
+ Assert.assertEquals(resultSet.getObject(8), List.of("first", "second"));
+ List<?> bytes = (List<?>) resultSet.getObject(9);
+ Assert.assertEquals(bytes.get(0), new byte[]{0, (byte) 0xff});
+ Assert.assertEquals(bytes.get(1), new byte[]{0x10, 0x20});
+ Assert.assertEquals(resultSet.getObject(10), List.of(
+ UUID.fromString("00000000-0000-0000-0000-000000000001"),
+ UUID.fromString("00000000-0000-0000-0000-000000000002")));
+
+ ResultSetMetaData metadata = resultSet.getMetaData();
+ for (int columnIndex = 1; columnIndex <= metadata.getColumnCount();
columnIndex++) {
+ Assert.assertEquals(metadata.getColumnType(columnIndex),
Types.JAVA_OBJECT);
+ Assert.assertEquals(metadata.getColumnClassName(columnIndex),
List.class.getTypeName());
+ }
+ }
+
+ @Test
+ public void testGetUuid()
+ throws Exception {
+ UUID uuid = UUID.fromString("00000000-0000-0000-0000-000000000001");
+ PinotResultSet resultSet = PinotResultSet.fromJson("{\"resultTable\":{"
+ +
"\"dataSchema\":{\"columnNames\":[\"uuid\"],\"columnDataTypes\":[\"UUID\"]},"
+ + "\"rows\":[[\"" + uuid + "\"]]}}");
+
+ Assert.assertTrue(resultSet.next());
+ Assert.assertEquals(resultSet.getObject(1), uuid);
+ Assert.assertEquals(resultSet.getObject("uuid", UUID.class), uuid);
+
+ ResultSetMetaData metadata = resultSet.getMetaData();
+ Assert.assertEquals(metadata.getColumnType(1), Types.OTHER);
+ Assert.assertEquals(metadata.getColumnClassName(1),
UUID.class.getTypeName());
+ }
+
+ @Test
+ public void testGetMapWithMixedNumericTypes()
+ throws Exception {
+ PinotResultSet resultSet = PinotResultSet.fromJson("{\"resultTable\":{"
+ +
"\"dataSchema\":{\"columnNames\":[\"map\"],\"columnDataTypes\":[\"MAP\"]},"
+ +
"\"rows\":[[{\"small\":1,\"large\":2147483648,\"float\":1.25,\"decimal\":2.5,"
+ + "\"nested\":{\"values\":[2,3.5]}}]]}}");
+
+ Assert.assertTrue(resultSet.next());
+ Assert.assertEquals(resultSet.getObject(1), Map.of(
+ "small", 1,
+ "large", 2147483648L,
+ "float", 1.25d,
+ "decimal", 2.5d,
+ "nested", Map.of("values", List.of(2, 3.5d))));
+ }
+
+ @Test
+ public void testGetObjectErrors()
+ throws Exception {
+ PinotResultSet resultSet = PinotResultSet.fromJson("{\"resultTable\":{"
+ +
"\"dataSchema\":{\"columnNames\":[\"map\",\"bytes\",\"timestamp\",\"decimal\",\"unknown\"],"
+ +
"\"columnDataTypes\":[\"MAP\",\"BYTES\",\"TIMESTAMP\",\"BIG_DECIMAL\",\"UNKNOWN\"]},"
+ +
"\"rows\":[[\"not-json\",\"zz\",\"not-a-timestamp\",\"not-a-decimal\",\"value\"]]}}");
+
+ Assert.assertTrue(resultSet.next());
+ Assert.expectThrows(SQLDataException.class, () -> resultSet.getObject(1));
+ Assert.expectThrows(SQLDataException.class, () -> resultSet.getObject(2));
+ Assert.expectThrows(SQLDataException.class, () -> resultSet.getObject(3));
+ Assert.expectThrows(SQLDataException.class, () -> resultSet.getObject(4));
+ Assert.expectThrows(SQLDataException.class, () -> resultSet.getObject(5));
+ }
+
+ @Test
+ public void testGetAdditionalScalarTypes()
+ throws Exception {
+ PinotResultSet resultSet = PinotResultSet.fromJson("{\"resultTable\":{"
+ +
"\"dataSchema\":{\"columnNames\":[\"json\",\"decimal\",\"timestamp\"],"
+ + "\"columnDataTypes\":[\"JSON\",\"BIG_DECIMAL\",\"TIMESTAMP\"]},"
+ + "\"rows\":[[\"{\\\"key\\\":1}\",\"123.450\",\"2020-01-01
12:00:00\"]]}}");
+
+ Assert.assertTrue(resultSet.next());
+ Assert.assertEquals(resultSet.getObject(1), "{\"key\":1}");
+ Assert.assertEquals(resultSet.getObject(2), new BigDecimal("123.450"));
+ Assert.assertEquals(resultSet.getObject(3), Timestamp.valueOf("2020-01-01
12:00:00"));
+ }
+
+ @Test
+ public void testNullAndEmptyArrays()
+ throws Exception {
+ PinotResultSet resultSet = PinotResultSet.fromJson("{\"resultTable\":{"
+ + "\"dataSchema\":{\"columnNames\":[\"nulls\",\"empty\"],"
+ +
"\"columnDataTypes\":[\"INT_ARRAY\",\"STRING_ARRAY\"]},\"rows\":[[null,[]]]}}");
+
+ Assert.assertTrue(resultSet.next());
+ Assert.assertNull(resultSet.getObject(1, List.class));
+ Assert.assertTrue(resultSet.wasNull());
+ Assert.assertNull(resultSet.getObject(1));
+ Assert.assertTrue(resultSet.wasNull());
+ Assert.assertEquals(resultSet.getObject(2), List.of());
+ }
+
@Test
public void testCursorMovement()
throws Exception {
@@ -187,30 +327,6 @@ public class PinotResultSetTest {
}
}
- @Test
- public void testGetCalculatedScale() {
- PinotResultSet pinotResultSet = new PinotResultSet();
- int calculatedResult;
-
- calculatedResult = pinotResultSet.getCalculatedScale("1");
- Assert.assertEquals(calculatedResult, 0);
-
- calculatedResult = pinotResultSet.getCalculatedScale("1.0");
- Assert.assertEquals(calculatedResult, 1);
-
- calculatedResult = pinotResultSet.getCalculatedScale("1.2");
- Assert.assertEquals(calculatedResult, 1);
-
- calculatedResult = pinotResultSet.getCalculatedScale("1.23");
- Assert.assertEquals(calculatedResult, 2);
-
- calculatedResult = pinotResultSet.getCalculatedScale("1.234");
- Assert.assertEquals(calculatedResult, 3);
-
- calculatedResult = pinotResultSet.getCalculatedScale("-1.234");
- Assert.assertEquals(calculatedResult, 3);
- }
-
@Test
public void testDateFromStringConcurrent()
throws Throwable {
diff --git
a/pinot-clients/pinot-jdbc-client/src/test/java/org/apache/pinot/client/grpc/PinotGrpcResultSetTest.java
b/pinot-clients/pinot-jdbc-client/src/test/java/org/apache/pinot/client/grpc/PinotGrpcResultSetTest.java
new file mode 100644
index 00000000000..bf87ee87352
--- /dev/null
+++
b/pinot-clients/pinot-jdbc-client/src/test/java/org/apache/pinot/client/grpc/PinotGrpcResultSetTest.java
@@ -0,0 +1,260 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.client.grpc;
+
+import com.google.protobuf.ByteString;
+import java.lang.reflect.Field;
+import java.math.BigDecimal;
+import java.sql.ResultSetMetaData;
+import java.sql.SQLDataException;
+import java.sql.Timestamp;
+import java.sql.Types;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.UUID;
+import org.apache.pinot.client.PinotResultSet;
+import org.apache.pinot.common.proto.Broker;
+import org.apache.pinot.common.response.broker.ResultTable;
+import org.apache.pinot.common.response.encoder.JsonResponseEncoder;
+import org.apache.pinot.common.utils.DataSchema;
+import org.apache.pinot.common.utils.DataSchema.ColumnDataType;
+import org.apache.pinot.spi.utils.CommonConstants;
+import org.apache.pinot.spi.utils.JsonUtils;
+import org.testng.annotations.Test;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertFalse;
+import static org.testng.Assert.assertNull;
+import static org.testng.Assert.assertTrue;
+import static org.testng.Assert.expectThrows;
+
+
+/// Tests collection-valued results returned by the gRPC JDBC result set.
+public class PinotGrpcResultSetTest {
+
+ @Test
+ public void testGetMapAndArrays()
+ throws Exception {
+ String[] columnNames = {
+ "map", "booleans", "ints", "longs", "floats", "doubles", "decimals",
"timestamps", "strings", "bytes",
+ "uuids"
+ };
+ ColumnDataType[] columnTypes = {
+ ColumnDataType.MAP, ColumnDataType.BOOLEAN_ARRAY,
ColumnDataType.INT_ARRAY, ColumnDataType.LONG_ARRAY,
+ ColumnDataType.FLOAT_ARRAY, ColumnDataType.DOUBLE_ARRAY,
ColumnDataType.BIG_DECIMAL_ARRAY,
+ ColumnDataType.TIMESTAMP_ARRAY, ColumnDataType.STRING_ARRAY,
ColumnDataType.BYTES_ARRAY,
+ ColumnDataType.UUID_ARRAY
+ };
+ Object[] row = {
+ Map.of("name", "pinot", "count", 2),
+ new boolean[]{true, false},
+ new int[]{1, 2},
+ new long[]{2147483648L, 3L},
+ new float[]{1.25f, 2.5f},
+ new double[]{1.5, 2.75},
+ new BigDecimal[]{new BigDecimal("1.20"), new BigDecimal("3.4")},
+ new Timestamp[]{Timestamp.valueOf("2020-01-01 12:00:00"),
Timestamp.valueOf("2021-02-03 04:05:06")},
+ new String[]{"first", "second"},
+ new byte[][]{new byte[]{0, (byte) 0xff}, new byte[]{0x10, 0x20}},
+ new UUID[]{UUID.fromString("00000000-0000-0000-0000-000000000001"),
+ UUID.fromString("00000000-0000-0000-0000-000000000002")}
+ };
+
+ PinotGrpcResultSet resultSet = createResultSet(columnNames, columnTypes,
row);
+
+ assertTrue(resultSet.next());
+ assertEquals(resultSet.getObject("map"), Map.of("name", "pinot", "count",
2));
+ assertEquals(resultSet.getObject(2), List.of(true, false));
+ assertEquals(resultSet.getObject(3), List.of(1, 2));
+ assertEquals(resultSet.getObject(4), List.of(2147483648L, 3L));
+ assertEquals(resultSet.getObject(5), List.of(1.25f, 2.5f));
+ assertEquals(resultSet.getObject(6), List.of(1.5, 2.75));
+ assertEquals(resultSet.getObject(7), List.of(new BigDecimal("1.20"), new
BigDecimal("3.4")));
+ assertEquals(resultSet.getObject(8),
+ List.of(Timestamp.valueOf("2020-01-01 12:00:00"),
Timestamp.valueOf("2021-02-03 04:05:06")));
+ assertEquals(resultSet.getObject(9), List.of("first", "second"));
+ List<?> bytes = (List<?>) resultSet.getObject(10);
+ assertEquals(bytes.get(0), new byte[]{0, (byte) 0xff});
+ assertEquals(bytes.get(1), new byte[]{0x10, 0x20});
+ assertEquals(resultSet.getObject(11), List.of(
+ UUID.fromString("00000000-0000-0000-0000-000000000001"),
+ UUID.fromString("00000000-0000-0000-0000-000000000002")));
+
+ ResultSetMetaData metadata = resultSet.getMetaData();
+ assertEquals(metadata.getColumnType(1), Types.JAVA_OBJECT);
+ assertEquals(metadata.getColumnClassName(1), Map.class.getTypeName());
+ for (int columnIndex = 2; columnIndex <= metadata.getColumnCount();
columnIndex++) {
+ assertEquals(metadata.getColumnType(columnIndex), Types.JAVA_OBJECT);
+ assertEquals(metadata.getColumnClassName(columnIndex),
List.class.getTypeName());
+ }
+ }
+
+ @Test
+ public void testNullAndEmptyValues()
+ throws Exception {
+ PinotGrpcResultSet resultSet = createResultSet(
+ new String[]{"nullMap", "nullInts", "nullStrings", "emptyMap",
"emptyInts", "emptyStrings"},
+ new ColumnDataType[]{ColumnDataType.MAP, ColumnDataType.INT_ARRAY,
ColumnDataType.STRING_ARRAY,
+ ColumnDataType.MAP, ColumnDataType.INT_ARRAY,
ColumnDataType.STRING_ARRAY},
+ new Object[]{null, null, null, Map.of(), new int[0], new String[0]});
+
+ assertTrue(resultSet.next());
+ assertNull(resultSet.getObject(1, List.class));
+ assertTrue(resultSet.wasNull());
+ for (int columnIndex = 1; columnIndex <= 3; columnIndex++) {
+ assertNull(resultSet.getObject(columnIndex));
+ assertTrue(resultSet.wasNull());
+ }
+ assertEquals(resultSet.getObject(4), Map.of());
+ assertFalse(resultSet.wasNull());
+ for (int columnIndex = 5; columnIndex <= 6; columnIndex++) {
+ assertEquals(resultSet.getObject(columnIndex), List.of());
+ assertFalse(resultSet.wasNull());
+ }
+ }
+
+ @Test
+ public void testGetStringForNullScalar()
+ throws Exception {
+ PinotGrpcResultSet resultSet = createResultSet(
+ new String[]{"value"}, new ColumnDataType[]{ColumnDataType.STRING},
new Object[]{null});
+
+ assertTrue(resultSet.next());
+ assertNull(resultSet.getString(1));
+ assertTrue(resultSet.wasNull());
+ }
+
+ @Test
+ public void testGetUuid()
+ throws Exception {
+ UUID uuid = UUID.fromString("00000000-0000-0000-0000-000000000001");
+ PinotGrpcResultSet resultSet = createResultSet(
+ new String[]{"uuid"}, new ColumnDataType[]{ColumnDataType.UUID}, new
Object[]{uuid});
+
+ assertTrue(resultSet.next());
+ assertEquals(resultSet.getObject(1), uuid);
+ assertEquals(resultSet.getObject("uuid", UUID.class), uuid);
+
+ ResultSetMetaData metadata = resultSet.getMetaData();
+ assertEquals(metadata.getColumnType(1), Types.OTHER);
+ assertEquals(metadata.getColumnClassName(1), UUID.class.getTypeName());
+ }
+
+ @Test
+ public void testGetMapWithMixedNumericTypes()
+ throws Exception {
+ Map<String, Object> map = new HashMap<>(Map.of(
+ "small", 1,
+ "large", 2147483648L,
+ "float", 1.25d,
+ "nested", Map.of("value", 2),
+ "lists", List.of(List.of(1, true), List.of("pinot")),
+ "maps", List.of(Map.of("name", "pinot"), Map.of("count", 2)),
+ "arrays", new Object[]{new int[]{1, 2}, new Object[]{true, null,
Map.of("value", 3)}}));
+ map.put("null", null);
+ DataSchema schema = new DataSchema(new String[]{"map"}, new
ColumnDataType[]{ColumnDataType.MAP});
+ Object[] row = {map};
+ PinotGrpcResultSet resultSet = createResultSetFromFormattedRow(schema,
row);
+ List<Object[]> rows = new ArrayList<>();
+ rows.add(row);
+ PinotResultSet httpResultSet = PinotResultSet.fromJson(
+ JsonUtils.objectToString(Map.of("resultTable", new ResultTable(schema,
rows))));
+
+ assertTrue(resultSet.next());
+ assertTrue(httpResultSet.next());
+ assertEquals(resultSet.getObject(1), httpResultSet.getObject(1));
+ }
+
+ @Test
+ public void testGetAdditionalScalarTypes()
+ throws Exception {
+ PinotGrpcResultSet resultSet = createResultSet(
+ new String[]{"json", "decimal", "timestamp"},
+ new ColumnDataType[]{ColumnDataType.JSON, ColumnDataType.BIG_DECIMAL,
ColumnDataType.TIMESTAMP},
+ new Object[]{"{\"key\":1}", new BigDecimal("123.450"),
Timestamp.valueOf("2020-01-01 12:00:00")});
+
+ assertTrue(resultSet.next());
+ assertEquals(resultSet.getObject(1), "{\"key\":1}");
+ assertEquals(resultSet.getObject(2), new BigDecimal("123.450"));
+ assertEquals(resultSet.getObject(3), Timestamp.valueOf("2020-01-01
12:00:00"));
+ }
+
+ @Test
+ public void testGetObjectErrors()
+ throws Exception {
+ DataSchema schema = new DataSchema(
+ new String[]{"map", "ints", "bytes", "timestamp", "decimal"},
+ new ColumnDataType[]{ColumnDataType.MAP, ColumnDataType.INT_ARRAY,
ColumnDataType.BYTES,
+ ColumnDataType.TIMESTAMP, ColumnDataType.BIG_DECIMAL});
+ PinotGrpcResultSet resultSet = createResultSetFromFormattedRow(
+ schema, new Object[]{Map.of(), new int[0], "zz", "not-a-timestamp",
"not-a-decimal"});
+
+ assertTrue(resultSet.next());
+ setCurrentRowValue(resultSet, 0, "not-a-map");
+ setCurrentRowValue(resultSet, 1, "not-an-array");
+ expectThrows(SQLDataException.class, () -> resultSet.getObject(1));
+ expectThrows(SQLDataException.class, () -> resultSet.getObject(2));
+ expectThrows(SQLDataException.class, () -> resultSet.getObject(3));
+ expectThrows(SQLDataException.class, () -> resultSet.getObject(4));
+ expectThrows(SQLDataException.class, () -> resultSet.getObject(5));
+ }
+
+ private static void setCurrentRowValue(PinotGrpcResultSet resultSet, int
columnIndex, Object value)
+ throws Exception {
+ Field currentRowBatchField =
PinotGrpcResultSet.class.getDeclaredField("_currentRowBatch");
+ currentRowBatchField.setAccessible(true);
+ ResultTable currentRowBatch = (ResultTable)
currentRowBatchField.get(resultSet);
+ currentRowBatch.getRows().get(0)[columnIndex] = value;
+ }
+
+ private static PinotGrpcResultSet createResultSet(String[] columnNames,
ColumnDataType[] columnTypes, Object[] row)
+ throws Exception {
+ DataSchema schema = new DataSchema(columnNames, columnTypes);
+ Object[] formattedRow = row.clone();
+ for (int i = 0; i < formattedRow.length; i++) {
+ if (formattedRow[i] != null) {
+ formattedRow[i] = columnTypes[i].format(formattedRow[i]);
+ }
+ }
+ return createResultSetFromFormattedRow(schema, formattedRow);
+ }
+
+ private static PinotGrpcResultSet createResultSetFromFormattedRow(DataSchema
schema, Object[] formattedRow)
+ throws Exception {
+ List<Object[]> rows = new ArrayList<>();
+ rows.add(formattedRow);
+ byte[] encodedRows = new JsonResponseEncoder().encodeResultTable(new
ResultTable(schema, rows), 0, rows.size());
+
+ Broker.BrokerResponse metadataResponse = Broker.BrokerResponse.newBuilder()
+ .setPayload(ByteString.copyFromUtf8("{}"))
+ .build();
+ Broker.BrokerResponse schemaResponse = Broker.BrokerResponse.newBuilder()
+ .setPayload(ByteString.copyFrom(schema.toBytes()))
+ .build();
+ Broker.BrokerResponse rowsResponse = Broker.BrokerResponse.newBuilder()
+ .setPayload(ByteString.copyFrom(encodedRows))
+ .putMetadata("rowSize", Integer.toString(rows.size()))
+ .putMetadata(CommonConstants.Broker.Grpc.COMPRESSION, "NONE")
+ .putMetadata(CommonConstants.Broker.Grpc.ENCODING, "JSON")
+ .build();
+ return new PinotGrpcResultSet(List.of(metadataResponse, schemaResponse,
rowsResponse).iterator());
+ }
+}
diff --git
a/pinot-clients/pinot-jdbc-client/src/test/java/org/apache/pinot/client/utils/BigDecimalUtilsTest.java
b/pinot-clients/pinot-jdbc-client/src/test/java/org/apache/pinot/client/utils/BigDecimalUtilsTest.java
new file mode 100644
index 00000000000..2bd40c2a2d2
--- /dev/null
+++
b/pinot-clients/pinot-jdbc-client/src/test/java/org/apache/pinot/client/utils/BigDecimalUtilsTest.java
@@ -0,0 +1,48 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.client.utils;
+
+import org.testng.Assert;
+import org.testng.annotations.Test;
+
+public class BigDecimalUtilsTest {
+
+ @Test
+ public void testGetCalculatedScale() {
+ int calculatedResult;
+
+ calculatedResult = BigDecimalUtils.getCalculatedScale("1");
+ Assert.assertEquals(calculatedResult, 0);
+
+ calculatedResult = BigDecimalUtils.getCalculatedScale("1.0");
+ Assert.assertEquals(calculatedResult, 1);
+
+ calculatedResult = BigDecimalUtils.getCalculatedScale("1.2");
+ Assert.assertEquals(calculatedResult, 1);
+
+ calculatedResult = BigDecimalUtils.getCalculatedScale("1.23");
+ Assert.assertEquals(calculatedResult, 2);
+
+ calculatedResult = BigDecimalUtils.getCalculatedScale("1.234");
+ Assert.assertEquals(calculatedResult, 3);
+
+ calculatedResult = BigDecimalUtils.getCalculatedScale("-1.234");
+ Assert.assertEquals(calculatedResult, 3);
+ }
+}
diff --git
a/pinot-common/src/main/java/org/apache/pinot/common/response/encoder/JsonResponseEncoder.java
b/pinot-common/src/main/java/org/apache/pinot/common/response/encoder/JsonResponseEncoder.java
index aa86400de85..6c7f146cd51 100644
---
a/pinot-common/src/main/java/org/apache/pinot/common/response/encoder/JsonResponseEncoder.java
+++
b/pinot-common/src/main/java/org/apache/pinot/common/response/encoder/JsonResponseEncoder.java
@@ -24,11 +24,10 @@ import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
-import java.util.HashMap;
import java.util.List;
-import java.util.Map;
import org.apache.pinot.common.response.broker.ResultTable;
import org.apache.pinot.common.utils.DataSchema;
+import org.apache.pinot.common.utils.DataSchema.ColumnDataType;
import org.apache.pinot.spi.utils.JsonUtils;
@@ -65,12 +64,14 @@ public class JsonResponseEncoder implements ResponseEncoder
{
JsonNode jsonRow = JsonUtils.stringToJsonNode(rowString);
Object[] row = new Object[jsonRow.size()];
for (int columnIdx = 0; columnIdx < jsonRow.size(); columnIdx++) {
- DataSchema.ColumnDataType columnDataType =
schema.getColumnDataType(columnIdx);
+ ColumnDataType columnDataType = schema.getColumnDataType(columnIdx);
JsonNode jsonValue = jsonRow.get(columnIdx);
- if (columnDataType.isArray()) {
+ if (jsonValue.isNull()) {
+ row[columnIdx] = null;
+ } else if (columnDataType.isArray()) {
row[columnIdx] = extractArray(columnDataType, jsonValue);
- } else if (columnDataType == DataSchema.ColumnDataType.MAP) {
- row[columnIdx] = extractMap(jsonValue);
+ } else if (columnDataType == ColumnDataType.MAP) {
+ row[columnIdx] = JsonUtils.jsonNodeToMap(jsonValue);
} else {
row[columnIdx] = extractValue(columnDataType, jsonValue);
}
@@ -80,85 +81,7 @@ public class JsonResponseEncoder implements ResponseEncoder {
return new ResultTable(schema, rows);
}
- private Object[] extractArray(JsonNode jsonValue) {
- Object[] array = new Object[jsonValue.size()];
- for (int k = 0; k < jsonValue.size(); k++) {
- if (jsonValue.get(k).isNull()) {
- array[k] = null;
- } else if (jsonValue.get(k).isBoolean()) {
- array[k] = jsonValue.get(k).asBoolean();
- } else if (jsonValue.get(k).isInt()) {
- array[k] = jsonValue.get(k).asInt();
- } else if (jsonValue.get(k).isLong()) {
- array[k] = jsonValue.get(k).asLong();
- } else if (jsonValue.get(k).isFloat()) {
- array[k] = jsonValue.get(k).floatValue();
- } else if (jsonValue.get(k).isDouble()) {
- array[k] = jsonValue.get(k).asDouble();
- } else if (jsonValue.get(k).isTextual()) {
- array[k] = jsonValue.get(k).textValue();
- } else if (jsonValue.isArray()) {
- array[k] = extractArray(jsonValue.get(k));
- } else if (jsonValue.isObject()) {
- array[k] = extractMap(jsonValue.get(k));
- } else {
- array[k] = jsonValue.get(k).toString();
- }
- }
- return array;
- }
-
- private Object extractValue(JsonNode jsonValue) {
- if (jsonValue.isNull()) {
- return null;
- }
- if (jsonValue.isBoolean()) {
- return jsonValue.asBoolean();
- }
- if (jsonValue.isShort()) {
- return jsonValue.shortValue();
- }
- if (jsonValue.isBigInteger()) {
- return jsonValue.bigIntegerValue();
- }
- if (jsonValue.isBigDecimal()) {
- return jsonValue.decimalValue();
- }
- if (jsonValue.isInt()) {
- return jsonValue.asInt();
- }
- if (jsonValue.isLong()) {
- return jsonValue.asLong();
- }
- if (jsonValue.isFloat()) {
- return jsonValue.floatValue();
- }
- if (jsonValue.isDouble()) {
- return jsonValue.asDouble();
- }
- if (jsonValue.isTextual()) {
- return jsonValue.textValue();
- }
- if (jsonValue.isArray()) {
- return extractArray(jsonValue);
- }
- if (jsonValue.isObject()) {
- return extractMap(jsonValue);
- }
- return jsonValue.toString();
- }
-
- private Map<String, Object> extractMap(JsonNode jsonValue) {
- Map<String, Object> map = new HashMap<>();
- jsonValue.properties().forEach(entry -> {
- String key = entry.getKey();
- Object value = extractValue(entry.getValue());
- map.put(key, value);
- });
- return map;
- }
-
- private static Object extractArray(DataSchema.ColumnDataType columnDataType,
JsonNode jsonValue) {
+ private static Object extractArray(ColumnDataType columnDataType, JsonNode
jsonValue) {
switch (columnDataType) {
case BOOLEAN_ARRAY:
boolean[] booleanArray = new boolean[jsonValue.size()];
@@ -205,7 +128,7 @@ public class JsonResponseEncoder implements ResponseEncoder
{
}
}
- private static Object extractValue(DataSchema.ColumnDataType columnDataType,
JsonNode jsonValue) {
+ private static Object extractValue(ColumnDataType columnDataType, JsonNode
jsonValue) {
if (jsonValue.isNull()) {
return null;
}
diff --git
a/pinot-common/src/main/java/org/apache/pinot/common/utils/DataSchema.java
b/pinot-common/src/main/java/org/apache/pinot/common/utils/DataSchema.java
index 83fc755230f..70458405884 100644
--- a/pinot-common/src/main/java/org/apache/pinot/common/utils/DataSchema.java
+++ b/pinot-common/src/main/java/org/apache/pinot/common/utils/DataSchema.java
@@ -41,6 +41,9 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.UUID;
+import java.util.function.Function;
+import java.util.stream.Collectors;
+import javax.annotation.Nullable;
import org.apache.calcite.rel.type.RelDataType;
import org.apache.calcite.rel.type.RelDataTypeFactory;
import org.apache.calcite.sql.type.SqlTypeName;
@@ -435,6 +438,9 @@ public class DataSchema {
private static final EnumSet<ColumnDataType> NUMERIC_ARRAY_TYPES =
EnumSet.of(INT_ARRAY, LONG_ARRAY, FLOAT_ARRAY, DOUBLE_ARRAY,
BIG_DECIMAL_ARRAY);
private static final EnumSet<ColumnDataType> INTEGRAL_ARRAY_TYPES =
EnumSet.of(INT_ARRAY, LONG_ARRAY);
+ private static final Map<String, ColumnDataType> NAME_TO_TYPES = Arrays
+ .stream(ColumnDataType.values())
+ .collect(Collectors.toUnmodifiableMap(ColumnDataType::name,
Function.identity()));
// stored data type.
private final ColumnDataType _storedColumnDataType;
@@ -1167,6 +1173,19 @@ public class DataSchema {
}
}
+ @Nullable
+ public static ColumnDataType forName(String typeName) {
+ return NAME_TO_TYPES.get(typeName);
+ }
+
+ public static boolean isArray(String typeName) {
+ ColumnDataType type = NAME_TO_TYPES.get(typeName);
+ if (type == null) {
+ return false;
+ }
+ return type.isArray();
+ }
+
/// Renders a single UUID as its canonical lowercase RFC 4122 string.
Accepts every representation
/// [UuidUtils#toUUID(Object)] does: `UUID`, `byte[]`, [ByteArray] and
`CharSequence`. A non-canonical (e.g.
/// upper-case) string input is re-canonicalized rather than passed
through.
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]