github-actions[bot] commented on code in PR #68786:
URL: https://github.com/apache/doris/pull/68786#discussion_r4235675539
##########
fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergTypeMapping.java:
##########
@@ -124,28 +124,24 @@ private static ConnectorType
fromPrimitive(Type.PrimitiveType primitive,
case STRING:
return ConnectorType.of("STRING");
case UUID:
- return enableMappingVarbinary
- ? ConnectorType.of("VARBINARY", 16, 0) :
ConnectorType.of("STRING");
+ return ConnectorType.of("VARBINARY", 16, 0);
case BINARY:
Review Comment:
[P1] Keep V1 equality deletes readable after UUID becomes VARBINARY. With
`enable_file_scanner_v2=false`, an Iceberg table having a single UUID
equality-delete field reaches `SimpleEqualityDelete::_build_set`; its VARBINARY
key is passed to `create_set`, whose predicate factory throws
`NOT_IMPLEMENTED_ERROR` for that type. The changed UUID compatibility suite now
runs only V2 and explicitly notes the V1 set cannot consume VARBINARY. This
breaks scans that worked with the former default STRING mapping. Add byte-aware
V1 delete-key support or route these splits to V2, and retain a V1 regression.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/connector/converter/ConnectorWriteValueConverter.java:
##########
@@ -0,0 +1,131 @@
+// 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.doris.datasource.connector.converter;
+
+import org.apache.doris.catalog.Column;
+import org.apache.doris.nereids.trees.expressions.ArrayItemReference;
+import org.apache.doris.nereids.trees.expressions.Cast;
+import org.apache.doris.nereids.trees.expressions.Expression;
+import org.apache.doris.nereids.trees.expressions.IsNull;
+import org.apache.doris.nereids.trees.expressions.functions.scalar.ArrayMap;
+import
org.apache.doris.nereids.trees.expressions.functions.scalar.CreateNamedStruct;
+import org.apache.doris.nereids.trees.expressions.functions.scalar.ElementAt;
+import org.apache.doris.nereids.trees.expressions.functions.scalar.If;
+import org.apache.doris.nereids.trees.expressions.functions.scalar.Lambda;
+import org.apache.doris.nereids.trees.expressions.functions.scalar.Replace;
+import org.apache.doris.nereids.trees.expressions.functions.scalar.Unhex;
+import org.apache.doris.nereids.trees.expressions.literal.IntegerLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.NullLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.StringLikeLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.StringLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.StructLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.UuidLiteral;
+import org.apache.doris.nereids.trees.expressions.literal.VarBinaryLiteral;
+import org.apache.doris.nereids.types.ArrayType;
+import org.apache.doris.nereids.types.DataType;
+import org.apache.doris.nereids.types.StringType;
+import org.apache.doris.nereids.types.StructField;
+import org.apache.doris.nereids.types.StructType;
+import org.apache.doris.nereids.types.UuidType;
+import org.apache.doris.nereids.types.VarBinaryType;
+import org.apache.doris.nereids.util.TypeCoercionUtils;
+
+import java.nio.ByteBuffer;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.UUID;
+
+/** Applies connector-declared textual input semantics before ordinary sink
type coercion. */
+public final class ConnectorWriteValueConverter {
+ private ConnectorWriteValueConverter() {
+ }
+
+ public static Expression convert(Column column, Expression input) {
+ if (column.getConnectorStringWriteType() == null) {
+ return input;
+ }
+ return convert(input,
DataType.fromCatalogType(column.getConnectorStringWriteType()),
+ DataType.fromCatalogType(column.getType()));
+ }
+
+ private static Expression convert(Expression input, DataType semanticType,
DataType targetType) {
+ if (!needsConversion(input.getDataType(), semanticType)) {
+ return input;
+ }
+ if (semanticType instanceof ArrayType) {
+ DataType elementSemantic = ((ArrayType)
semanticType).getItemType();
+ DataType elementTarget = ((ArrayType) targetType).getItemType();
+ ArrayItemReference item = new ArrayItemReference("item", input);
+ return new ArrayMap(new
Lambda(Collections.singletonList(item.getName()),
+ convert(item.toSlot(), elementSemantic, elementTarget),
Collections.singletonList(item)));
+ }
+ if (semanticType instanceof StructType) {
+ List<StructField> semanticFields = ((StructType)
semanticType).getFields();
+ List<StructField> targetFields = ((StructType)
targetType).getFields();
+ List<Expression> fields = new ArrayList<>();
+ for (int i = 0; i < semanticFields.size(); i++) {
+ Expression field = input instanceof StructLiteral ?
((StructLiteral) input).getValue().get(i)
+ : TypeCoercionUtils.processBoundFunction(new
ElementAt(input, new IntegerLiteral(i + 1)));
+ fields.add(new StringLiteral(targetFields.get(i).getName()));
+ fields.add(convert(field, semanticFields.get(i).getDataType(),
targetFields.get(i).getDataType()));
+ }
+ Expression result = TypeCoercionUtils.processBoundFunction(
+ new CreateNamedStruct(fields.toArray(new Expression[0])));
+ // Reconstructing a nullable struct must not turn NULL into a
non-null struct of NULL fields.
+ return input.nullable() ? TypeCoercionUtils.processBoundFunction(
+ new If(new IsNull(input), new
NullLiteral(result.getDataType()), result)) : result;
+ }
+ if (input instanceof StringLikeLiteral) {
+ UUID uuid = new UuidLiteral(((StringLikeLiteral)
input).getStringValue()).getValue();
+ return new VarBinaryLiteral((VarBinaryType) targetType,
ByteBuffer.allocate(16)
+
.putLong(uuid.getMostSignificantBits()).putLong(uuid.getLeastSignificantBits()).array());
+ }
+ // Validate UUID text strictly before decoding canonical hex. A raw
string-to-binary cast
+ // would write 36 text bytes, and a permissive UUID cast would
silently turn bad input into NULL.
+ Expression uuidText = new Cast(new Cast(input, semanticType, false,
true), StringType.INSTANCE);
+ Expression hex = TypeCoercionUtils.processBoundFunction(
+ new Replace(uuidText, new StringLiteral("-"), new
StringLiteral("")));
+ Expression bytes = TypeCoercionUtils.processBoundFunction(new
Unhex(hex));
+ return new Cast(bytes, targetType);
+ }
+
+ private static boolean needsConversion(DataType input, DataType semantic) {
+ if (semantic instanceof UuidType) {
+ return input.isStringLikeType();
+ }
Review Comment:
[P2] Convert UUID-typed source values before Iceberg UUID sink coercion. For
`INSERT INTO iceberg_uuid SELECT id, u FROM native_uuid` (where `u` is UUID),
this returns false, so the Iceberg full-schema sink calls
`castIfNotSameType(UUID, VARBINARY(16))`; FE rejects that cast. The former
default STRING mapping accepted UUID-to-STRING and its writer parsed canonical
UUID text. Handle UUID inputs as well as strings, including nested UUID leaves,
and cover an INSERT SELECT from a UUID column.
##########
fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergTypeMapping.java:
##########
@@ -124,28 +124,24 @@ private static ConnectorType
fromPrimitive(Type.PrimitiveType primitive,
case STRING:
return ConnectorType.of("STRING");
case UUID:
- return enableMappingVarbinary
- ? ConnectorType.of("VARBINARY", 16, 0) :
ConnectorType.of("STRING");
+ return ConnectorType.of("VARBINARY", 16, 0);
case BINARY:
// Iceberg BINARY is unbounded. Emit VARBINARY with NO
explicit length so
// ConnectorColumnConverter applies
ScalarType.MAX_VARBINARY_LENGTH — byte-identical to
// legacy IcebergUtils
createVarbinaryType(VarBinaryType.MAX_VARBINARY_LENGTH). A
// concrete length (e.g. 65535) would render a different
DESCRIBE / SHOW CREATE type.
- return enableMappingVarbinary
- ? ConnectorType.of("VARBINARY") :
ConnectorType.of("STRING");
+ // Binary payloads need not be valid UTF-8.
+ return ConnectorType.of("VARBINARY");
case FIXED:
Review Comment:
[P2] Keep `to_json` usable on nested binary values after this mapping
change. For an Iceberg `ARRAY<BINARY>` column, FE accepts `to_json(items)`, but
BE's array JSONB serializer delegates each non-null VARBINARY element to
`DataTypeVarbinarySerDe`, which has no `serialize_column_to_jsonb`
implementation and returns NotSupported. The former nested STRING mapping
serialized successfully. Add a defined byte encoding for JSON or reject this
signature during analysis, and cover a non-null nested binary value.
##########
regression-test/suites/external_table_p0/iceberg/test_iceberg_transform_partitions.groovy:
##########
@@ -130,14 +133,15 @@ suite("test_iceberg_transform_partitions", "p0,external")
{
"""
// Bucket by BINARY
+ // These binary fixtures contain UTF-8 text; compare their explicit
STRING representation.
qt_bucket_binary_4_cnt1 """
- select count(*) from bucket_binary_4 where partition_key = 'abc';
+ select count(*) from bucket_binary_4 where cast(partition_key as
string) = 'abc';
"""
Review Comment:
[P2] Preserve direct predicates on formerly string-mapped binary columns.
Replacing `partition_key = 'abc'` with this cast hides a compatibility
regression: Iceberg BINARY now maps unconditionally to VARBINARY, and
`processComparisonPredicateInternal` rejects every direct VARBINARY comparison
during FE analysis. Existing catalog queries using `WHERE partition_key =
'abc'` therefore fail; the UUID suite similarly switched to `HEX(u)`. Add
byte-aware comparison support or a compatible direct-predicate path, then
retain the uncast predicate in regression coverage.
##########
fe/fe-connector/fe-connector-jdbc/src/main/java/org/apache/doris/connector/jdbc/JdbcQueryBuilder.java:
##########
@@ -190,6 +191,222 @@ public String buildQuery(String remoteDbName, String
remoteTableName,
return sql.toString();
}
+ private String timestampProjection(String expression,
org.apache.doris.connector.spi.ConnectorType type,
+ int depth) {
+ if ("TIMESTAMPTZ".equals(type.getTypeName())) {
+ if (dbType == JdbcDbType.MYSQL || dbType == JdbcDbType.OCEANBASE) {
+ // MySQL drivers can apply a cached timezone even in
getString() for binary results.
+ // A server-side text projection preserves the UTC session
fields and microseconds.
+ return "CAST(" + expression + " AS CHAR)";
+ }
+ if (dbType == JdbcDbType.CLICKHOUSE) {
+ return "toUnixTimestamp64Micro(toDateTime64(" + expression +
", 6))";
+ }
+ if (dbType == JdbcDbType.TRINO || dbType == JdbcDbType.PRESTO) {
+ return "(" + expression + " AT TIME ZONE 'UTC')";
+ }
+ }
+ if ("ARRAY".equals(type.getTypeName()) && containsInstant(type)) {
+ String element = "doris_ts_" + depth;
+ String converted = timestampProjection(element,
type.getChildren().get(0), depth + 1);
+ if (dbType == JdbcDbType.CLICKHOUSE) {
+ return "arrayMap(" + element + " -> " + converted + ", " +
expression + ")";
+ }
+ if (dbType == JdbcDbType.TRINO || dbType == JdbcDbType.PRESTO) {
+ return "transform(" + expression + ", " + element + " -> " +
converted + ")";
+ }
+ }
+ return expression;
+ }
+
+ private static boolean
containsInstant(org.apache.doris.connector.spi.ConnectorType type) {
+ return "TIMESTAMPTZ".equals(type.getTypeName()) ||
type.getChildren().stream().anyMatch(
+ JdbcQueryBuilder::containsInstant);
+ }
+
+ public String wrapPassthroughQuery(String query,
List<ConnectorColumnHandle> columns) {
+ if (columns.stream().noneMatch(c -> c instanceof JdbcColumnHandle
+ && containsInstant(((JdbcColumnHandle) c).getType()))
+ || (dbType != JdbcDbType.CLICKHOUSE && dbType !=
JdbcDbType.TRINO && dbType != JdbcDbType.PRESTO
+ && dbType != JdbcDbType.MYSQL && dbType !=
JdbcDbType.OCEANBASE)) {
+ return query;
+ }
+ // Project before driver decoding: a named-zone DST fold has already
lost its offset afterward.
+ StringJoiner projections = new StringJoiner(", ");
+ for (ConnectorColumnHandle column : columns) {
+ JdbcColumnHandle jdbcColumn = (JdbcColumnHandle) column;
+ String name = JdbcIdentifierQuoter.quoteRemoteIdentifier(dbType,
jdbcColumn.getRemoteName());
+ projections.add(timestampProjection(name, jdbcColumn.getType(), 0)
+ " AS " + name);
+ }
+ String inner = stripTerminalDelimiter(query.trim());
+ // WITH SESSION belongs to the Trino statement, not to a derived-table
query.
+ int queryStart = dbType == JdbcDbType.TRINO ?
trinoSessionQueryStart(inner) : 0;
+ String prefix = inner.substring(0, queryStart);
+ // A trailing SQL line comment must end before the wrapper closes its
derived table.
+ return prefix + "SELECT " + projections + " FROM (" +
inner.substring(queryStart)
+ + "\n) doris_jdbc_query";
+ }
+
+ private String stripTerminalDelimiter(String sql) {
+ List<Integer> delimiters = new java.util.ArrayList<>();
+ for (int i = 0; i < sql.length();) {
+ char c = sql.charAt(i);
+ if (Character.isWhitespace(c)) {
+ i++;
+ } else if (sql.startsWith("--", i) || (c == '#'
+ && (dbType == JdbcDbType.MYSQL || dbType ==
JdbcDbType.OCEANBASE))) {
+ while (i < sql.length() && sql.charAt(i) != '\n' &&
sql.charAt(i) != '\r') {
+ i++;
+ }
+ } else if (sql.startsWith("/*", i)) {
+ int depth = 1;
+ i += 2;
+ while (i < sql.length() && depth > 0) {
+ if (sql.startsWith("/*", i)) {
+ depth++;
+ i += 2;
+ } else if (sql.startsWith("*/", i)) {
+ depth--;
+ i += 2;
+ } else {
+ i++;
+ }
+ }
+ } else if (c == '\'' || c == '"' || c == '`') {
+ char quote = c;
+ for (i++; i < sql.length(); i++) {
+ if (sql.charAt(i) == '\\' && dbType != JdbcDbType.TRINO &&
dbType != JdbcDbType.PRESTO) {
+ i++;
Review Comment:
[P2] Honor MySQL NO_BACKSLASH_ESCAPES while stripping TVF delimiters. With
that SQL mode, `SELECT ts FROM t WHERE label='a\';` is valid (one literal
backslash), but this scanner skips its closing quote and retains the terminal
semicolon. The TIMESTAMPTZ projection then sends a derived table with that
semicolon inside it, which MySQL rejects although metadata discovery prepared
the original query. Use the connection's SQL mode or a SQL-aware delimiter
strategy, and cover this mode with a TIMESTAMP query TVF.
##########
fe/be-java-extensions/jdbc-scanner/src/main/java/org/apache/doris/jdbc/OracleTypeHandler.java:
##########
@@ -117,6 +120,11 @@ public PreparedStatement initializeStatement(Connection
conn, String sql,
@Override
public Object getColumnValue(ResultSet rs, int columnIndex, ColumnType
type,
ResultSetMetaData metadata) throws
SQLException {
+ if (type.getType() == ColumnType.Type.TIMESTAMPTZ) {
+ // Both driver paths preserve the instant instead of passing local
wall-clock fields to JNI.
+ Timestamp value = rs.getTimestamp(columnIndex);
+ return value == null ? null :
checkedUtcTimestamp(value.toInstant());
Review Comment:
[P2] Initialize Oracle's session time zone before reading TSLTZ. Both Oracle
mappings now expose `TIMESTAMP WITH LOCAL TIME ZONE` as TIMESTAMPTZ, and this
`getTimestamp` call is reached without `OracleConnection.setSessionTimeZone`;
`initializeStatement` only prepares the query. Oracle documents that the driver
call is required before TSLTZ access, so ordinary scans can fail with `Session
Time Zone not set!` or use the wrong zone. Set the driver session zone on each
scan connection and cover a non-UTC TSLTZ read.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]