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]

Reply via email to