github-actions[bot] commented on code in PR #68786:
URL: https://github.com/apache/doris/pull/68786#discussion_r4218883211
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/CatalogMgr.java:
##########
@@ -475,6 +485,42 @@ private boolean
validateCatalogPropsBeforeUpdate(ExternalCatalog catalog, Map<St
}
}
+ /**
+ * Migrate legacy markers after fenced master replay, before accepting
queries or starting checkpoints.
+ */
+ public void migrateVarbinaryMappingProperties() throws DdlException {
+ writeLock();
+ try {
+ for (CatalogIf catalog : idToCatalog.values()) {
+ if (!(catalog instanceof ExternalCatalog)) {
+ continue;
+ }
+ ExternalCatalog externalCatalog = (ExternalCatalog) catalog;
+ Map<String, String> migratedProperties = Maps.newHashMap();
+ for (String marker : new String[]
{CatalogProperty.ENABLE_MAPPING_VARBINARY,
+ CatalogProperty.ENABLE_MAPPING_TIMESTAMP_TZ}) {
+ if
(!Boolean.parseBoolean(externalCatalog.getProperties().get(marker))) {
Review Comment:
[P1] Preserve reads of Fluss binary partition tables during marker
migration. This loop forces `enable.mapping.varbinary=true` for Fluss catalogs
too, but `FlussPartitionColumnTypes.rejection` explicitly refuses BINARY/BYTES
partition columns in that mode because their partition names contain hex text.
Existing tables that were readable with the default/false marker now fail
partition listing and SELECT after promotion, and ALTER false is rewritten to
true. Decode the partition-name hex for VARBINARY before forcing this marker,
and cover an upgraded partitioned table.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/CatalogMgr.java:
##########
@@ -475,6 +485,42 @@ private boolean
validateCatalogPropsBeforeUpdate(ExternalCatalog catalog, Map<St
}
}
+ /**
+ * Migrate legacy markers after fenced master replay, before accepting
queries or starting checkpoints.
+ */
+ public void migrateVarbinaryMappingProperties() throws DdlException {
+ writeLock();
+ try {
+ for (CatalogIf catalog : idToCatalog.values()) {
+ if (!(catalog instanceof ExternalCatalog)) {
+ continue;
+ }
+ ExternalCatalog externalCatalog = (ExternalCatalog) catalog;
+ Map<String, String> migratedProperties = Maps.newHashMap();
+ for (String marker : new String[]
{CatalogProperty.ENABLE_MAPPING_VARBINARY,
+ CatalogProperty.ENABLE_MAPPING_TIMESTAMP_TZ}) {
+ if
(!Boolean.parseBoolean(externalCatalog.getProperties().get(marker))) {
+ migratedProperties.put(marker, "true");
+ }
+ }
+ if (migratedProperties.isEmpty()) {
+ continue;
+ }
+ CatalogLog log = new CatalogLog();
+ log.setCatalogId(catalog.getId());
+ log.setNewProps(migratedProperties);
+ // Use the existing ALTER format so running older followers
can replay the change.
+ // Journal first: a failed write must leave the marker
eligible for a retry.
+
Env.getCurrentEnv().getEditLog().logCatalogLog(OperationType.OP_ALTER_CATALOG_PROPS,
log);
+ // Migration must not revalidate unrelated legacy connection
properties or contact
+ // the external system while the master is still becoming
ready.
+ replayAlterCatalogProps(log, null, true);
Review Comment:
[P2] Run migration cleanup after releasing the catalog write lock.
`migrateVarbinaryMappingProperties` holds the outer write lock when it calls
`replayAlterCatalogProps`; that method releases only its reentrant hold before
running the detached authorization plugin's `close()` callback. A slow plugin
close therefore holds the global catalog lock through master promotion,
blocking catalog operations and delaying readiness. Preserve the journal/apply
order but defer these callbacks until the outer lock is released, as ordinary
ALTER does.
##########
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);
Review Comment:
[P2] Encode UUID identity partition values for the VARBINARY scan slot.
`getIdentityPartitionInfoMap` still sends a UUID as canonical text in
`columns_from_path`, but the projected UUID slot is now VARBINARY and BE
`DataTypeVarbinarySerDe::from_string` requires `0x` hex. Selecting the UUID
identity partition column therefore fails on an otherwise readable Iceberg
table. Serialize its 16 raw bytes as `0x` hex for this transport, and cover an
identity-partition scan.
##########
fe/fe-connector/fe-connector-jdbc/src/main/java/org/apache/doris/connector/jdbc/JdbcQueryBuilder.java:
##########
@@ -258,6 +326,17 @@ private boolean collectFilters(ConnectorExpression expr,
List<String> clauses,
* Mirrors the old JdbcScanNode.shouldPushDownConjunct() guards.
*/
private boolean shouldPushDownExpression(ConnectorExpression expr) {
+ // Remote NULL handling and calendar operations can differ from
decoded Doris instants.
+ if ((dbType == JdbcDbType.POSTGRESQL && hasInstant(expr)) ||
hasInstantLiteral(expr)
Review Comment:
[P1] Keep mixed instant and wall-clock column comparisons local unless their
cast is preserved. MySQL `TIMESTAMP` now maps to TIMESTAMPTZ while `DATETIME`
stays DATETIMEV2. In a Doris `+08:00` session, `ts > dt` casts `ts` to local
DATETIMEV2, but the connector drops that cast and this guard still pushes the
two-column predicate to the MySQL scan's UTC session. For a row with
`ts=00:00Z` and `dt=04:00`, Doris evaluates `08:00 > 04:00` as true while
remote MySQL evaluates `00:00 > 04:00` as false, so the row is lost before
local filtering. Block this pushdown or preserve the session-zone cast, and
cover the two-column non-UTC case.
##########
fe/be-java-extensions/jdbc-scanner/src/main/java/org/apache/doris/jdbc/OracleTypeHandler.java:
##########
@@ -285,4 +292,12 @@ private static String convertByteArrayToString(byte[]
bytes) {
public void setValidationQuery(HikariDataSource ds) {
ds.setConnectionTestQuery("SELECT 1 FROM dual");
}
+
+ @Override
+ public void setTimestampTz(java.sql.PreparedStatement statement, int
parameterIndex, LocalDateTime value)
+ throws SQLException {
+ // An unzoned TIMESTAMP bind is session-local for both Oracle TZ and
LOCAL TIME ZONE columns.
+ statement.setObject(parameterIndex, value.atOffset(ZoneOffset.UTC),
Types.TIMESTAMP_WITH_TIMEZONE);
Review Comment:
[P2] Provide an old-Oracle-driver bind for zoned writes. This handler
explicitly supports `ojdbc6` for reads, but this new setter always passes JDBC
4.2 `Types.TIMESTAMP_WITH_TIMEZONE` (code 2014). In Oracle's 11.2.0.4 `ojdbc6`,
`OraclePreparedStatement.setObject` has no case for 2014 and throws an SQL
exception; it accepts the proprietary TIMESTAMPTZ code `-101` with
`oracle.sql.TIMESTAMPTZ`. Every non-null TIMESTAMPTZ insert through that
supported driver therefore fails. Bind an old-driver-supported zoned value and
cover both old and current Oracle drivers.
##########
fe/be-java-extensions/jdbc-scanner/src/main/java/org/apache/doris/jdbc/TrinoTypeHandler.java:
##########
@@ -43,6 +50,11 @@ public class TrinoTypeHandler extends DefaultTypeHandler {
@Override
public Object getColumnValue(ResultSet rs, int columnIndex, ColumnType
type,
ResultSetMetaData metadata) throws
SQLException {
+ if (type.getType() == ColumnType.Type.TIMESTAMPTZ) {
+ // The remote projection and driver preserve the instant; JNI
receives UTC fields.
+ ZonedDateTime value = rs.getObject(columnIndex,
ZonedDateTime.class);
+ return value == null ? null :
LocalDateTime.ofInstant(value.toInstant(), ZoneOffset.UTC);
Review Comment:
[P2] Check the UTC range for Trino zoned values before JNI packing. Trino
and its JDBC driver can represent `timestamp with time zone` in year 10000, but
this new TIMESTAMPTZ branch passes that year to `putTimeStampTz`; Doris's
packed datetime accepts at most 9999. The new array-element conversion has the
same gap. Apply the UTC bound to scalar and array values, preserve column
nullability if returning NULL, and cover an out-of-range driver value.
##########
fe/be-java-extensions/jdbc-scanner/src/main/java/org/apache/doris/jdbc/TrinoTypeHandler.java:
##########
@@ -107,7 +119,71 @@ public ColumnValueConverter getOutputConverter(ColumnType
columnType, String rep
}
}
- private Object convertArray(List<?> input, ColumnType childType) {
- return input;
+ private List<?> convertArray(List<?> array, ColumnType type) {
+ if (array == null) {
+ return null;
+ }
+ if (array.isEmpty()) {
+ return Collections.emptyList();
+ }
+ switch (type.getType()) {
+ case DATE:
+ case DATEV2: {
+ List<LocalDate> result = Lists.newArrayList();
+ for (Object element : array) {
+ result.add(element != null ? ((Date)
element).toLocalDate() : null);
+ }
+ return result;
+ }
+ case TIMESTAMPTZ: {
+ List<LocalDateTime> result = Lists.newArrayList();
+ // Trino JDBC exposes timestamp-with-zone array elements as
java.sql.Timestamp.
+ for (Object element : array) {
+ result.add(element == null ? null
+ : LocalDateTime.ofInstant(((Timestamp)
element).toInstant(), ZoneOffset.UTC));
+ }
+ return result;
+ }
+ case DATETIME:
+ case DATETIMEV2: {
+ List<LocalDateTime> result = Lists.newArrayList();
+ for (Object element : array) {
+ result.add(element != null ? ((Timestamp)
element).toLocalDateTime() : null);
+ }
+ return result;
+ }
+ case ARRAY: {
+ List<List<?>> resultArray = Lists.newArrayList();
+ for (Object element : array) {
+ if (element == null) {
+ resultArray.add(null);
+ } else {
+ resultArray.add(
+ Lists.newArrayList(convertArray((List<?>)
element, type.getChildTypes().get(0))));
+ }
+ }
+ return resultArray;
+ }
+ default:
+ return array;
+ }
+ }
+
+ private static final DateTimeFormatter TIMESTAMP_TZ_WRITE_FORMATTER =
+ DateTimeFormatter.ofPattern("uuuu-MM-dd HH:mm:ss.SSSSSS");
+
+ @Override
+ public void setTimestampTz(java.sql.PreparedStatement statement, int
parameterIndex, LocalDateTime value)
+ throws SQLException {
+ // Trino/Presto require a string for typed zoned binds; Timestamp
drops the zone and sub-millisecond digits.
+ statement.setObject(parameterIndex,
value.format(TIMESTAMP_TZ_WRITE_FORMATTER) + " UTC",
Review Comment:
[P2] Bind zoned timestamps with a PrestoDB-supported parameter form. PRESTO
uses this same handler, but the official PrestoDB JDBC
`PrestoPreparedStatement.setObject(..., Types.TIMESTAMP_WITH_TIMEZONE)` throws
`Unsupported target SQL type` because that case is not implemented. The new
PRESTO TIMESTAMPTZ mapping sends non-null inserts here, so these writes fail at
bind time. Handle PrestoDB separately while preserving the UTC instant, and
cover its driver rather than only the PrestoSQL fork.
##########
fe/be-java-extensions/jdbc-scanner/src/main/java/org/apache/doris/jdbc/OracleTypeHandler.java:
##########
@@ -117,6 +119,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 :
LocalDateTime.ofInstant(value.toInstant(), ZoneOffset.UTC);
Review Comment:
[P2] Check the UTC year before packing Oracle zoned timestamps. Oracle
permits a `TIMESTAMP WITH TIME ZONE` near the end of year 9999 with a negative
offset, whose instant falls in UTC year 10000; it also permits BC dates. This
newly mapped TIMESTAMPTZ path passes that year straight to `putTimeStampTz`,
which bit-packs it although Doris accepts only years 0 through 9999. Reject or
null out-of-range instants before JNI conversion, with matching nullability and
boundary tests, as the PostgreSQL handler does.
##########
regression-test/suites/external_table_p0/jdbc/type_test/select/test_pg_all_types_select.groovy:
##########
@@ -38,6 +41,48 @@ suite("test_pg_all_types_select", "p0,external") {
qt_desc_all_types_null """desc catalog_pg_test.extreme_test;"""
+ // PostgreSQL infinities and BC/out-of-range years cannot be packed
into Doris timestamps.
+ assertEquals([[true], [true], [true], [true]],
Review Comment:
[P3] Record these fixed PostgreSQL range results in generated output and
keep the table after the run. The added `assertEquals` checks here and for
`timestamp_range_nullability` have no entries in this suite's `.out`, while the
new `finally` block drops the remote table even when a check fails. The
repository regression rules require ordered `qt` results generated by the
runner and drops before use only. Convert the fixed queries to `order_qt`,
generate the expected output, and remove the post-test drop.
##########
regression-test/suites/datatype_p0/test_varbinary_sql_support.groovy:
##########
@@ -0,0 +1,87 @@
+// 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.
+
+suite("test_varbinary_sql_support") {
+ def source = "binary_sql_input"
+ def table = "binary_sql_source"
+ def view = "binary_sql_view"
+ def ctas = "binary_sql_ctas"
+ def nested = "binary_sql_nested"
+ def mv = "binary_sql_mv"
+ def cleanup = {
+ sql "DROP MATERIALIZED VIEW IF EXISTS ${mv}"
+ sql "DROP VIEW IF EXISTS ${view}"
+ [nested, table].each { sql "DROP VIEW IF EXISTS ${it}" }
+ [ctas, source].each { sql "DROP TABLE IF EXISTS ${it}" }
+ }
+ cleanup()
+ try {
+ // Decode bytes in the execution layer; OLAP storage remains an
ordinary STRING column.
+ sql """CREATE TABLE ${source} (id INT, encoded STRING)
+ DUPLICATE KEY(id) DISTRIBUTED BY HASH(id) BUCKETS 3
+ PROPERTIES ('replication_num'='1')"""
+ sql """INSERT INTO ${source} VALUES
+ (0, NULL), (1, ''), (2, '00'), (3, '7F'), (4, '80'),
+ (5, 'AB'), (6, 'AB00'), (7, 'AB0001'), (8, 'AB'), (9, 'FF')"""
+ // TO_BINARY treats empty input as NULL; use a literal to exercise a
distinct empty binary value.
+ sql """CREATE VIEW ${table} AS SELECT id,
+ CASE WHEN encoded = '' THEN X'' ELSE to_binary(encoded) END AS
payload FROM ${source}"""
+ def bytes = [[0, null], [1, ""], [2, "00"], [3, "7F"], [4, "80"],
+ [5, "AB"], [6, "AB00"], [7, "AB0001"], [8, "AB"], [9,
"FF"]]
+ def readBytes = { name -> sql "SELECT id, from_binary(payload) FROM
${name} ORDER BY id" }
+ assertEquals(bytes, readBytes(table))
Review Comment:
[P3] Record fixed results with `order_qt` and keep tables after the run. The
hardcoded `assertEquals` checks here and later in the suite have no generated
`.out` result, and `finally { cleanup() }` drops the tables even when an
assertion fails. The repository's regression standards require generated `qt`
output for determined results and drops before setup only. Convert the fixed
cases to ordered `qt` checks, generate their output with the runner, and remove
the post-test cleanup.
--
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]