github-actions[bot] commented on code in PR #68786:
URL: https://github.com/apache/doris/pull/68786#discussion_r4226942343
##########
fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergPartitionUtils.java:
##########
@@ -840,6 +881,14 @@ private static IcebergRawPartition generateRawPartition(
int ordinal = partitionFieldOrdinals.get(partitionField.fieldId());
Object o = partitionData.get(ordinal, fieldClass);
String fieldValue = o == null ? null : o.toString();
+ Type fieldType =
partitionSpec.partitionType().fields().get(i).type();
+ // ByteBuffer.toString() describes its bounds, not the partition
bytes. Use the
+ // transformed type so bucket(binary) remains an integer while
identity/truncate keep bytes.
+ if (fieldType.typeId() == Type.TypeID.BINARY || fieldType.typeId()
== Type.TypeID.FIXED) {
Review Comment:
[P2] Preserve partition counts for bucketed binary columns. A `bucket(16,
binary_key)` partition has an integer field value, which this branch leaves as
text such as `3`, but the FE list pairs it with the source column's newly
forced VARBINARY type. `LiteralExprUtils.createLiteral` requires `0x` hex, so
every non-null bucket item is skipped. The latest Iceberg pin then exposes
`partition=0/0` for a nonempty table, and COUNT pushdown retains this false
count for `partition_num` SQL block rules. Iceberg's scan provider can still
read the rows, but partition metadata and rule enforcement are wrong. Carry the
transformed type for these items and cover bucket(binary) listing/counts; the
identity-binary fix does not cover an INT transform result.
##########
fe/fe-connector/fe-connector-jdbc/src/main/java/org/apache/doris/connector/jdbc/JdbcQueryBuilder.java:
##########
@@ -190,6 +191,86 @@ 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.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)) {
+ 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 = query.trim().replaceAll(";+$", "");
+ // A trailing SQL line comment must end before the wrapper closes its
derived table.
+ return "SELECT " + projections + " FROM (" + inner + "\n)
doris_jdbc_query";
Review Comment:
[P2] Keep Trino WITH SESSION at statement scope for TVF scans. A Trino query
TVF can use `WITH SESSION query_max_execution_time='2h' SELECT ts FROM t` with
a zoned timestamp column. Metadata discovery prepares that original statement
successfully, but this new projection wraps it as `SELECT ... FROM (WITH
SESSION ... SELECT ...)`, where Trino only permits `WITH SESSION` at statement
scope. The TVF therefore passes schema discovery and fails every scan with a
syntax error. Keep the session clause outside the derived table, or project the
timestamp without nesting the statement; add a TVF case with a TIMESTAMPTZ
column.
##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanPlanProvider.java:
##########
@@ -960,7 +966,9 @@ && hasTrustworthyLimitAccounting(table)
if (shouldUseNativeReader(paimonHandle.isForceJni(),
isForceJniScannerEnabled(session), hasVariantProjection,
- physicalVariantSchemaIds, optRawFiles)) {
+ physicalVariantSchemaIds, optRawFiles)
+ && !requiresLegacyOrcTimestampReader(
Review Comment:
[P2] Keep legacy LTZ and synthesized metadata readable together. On an ORC
table with `orc.timestamp-ltz.legacy-type=true`, `SELECT event_time,
__paimon_file_path FROM t` passes the initial FileScannerV2 metadata check.
Reading `event_time` makes this new condition force the split to JNI, where
`validateMetadataColumnReader(true, false)` unconditionally rejects
`__paimon_file_path` (and likewise `__paimon_row_index`). Before this routing
change the native split accepted the mixed projection. Provide a reader path
that can return both values, and test them in one query; the new test checks
each projection separately.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/CatalogMgr.java:
##########
@@ -475,6 +485,47 @@ 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 {
+ List<Runnable> cleanups = Lists.newArrayList();
+ 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.
+ cleanups.add(applyAlterCatalogProps(log, null, true, true,
false));
Review Comment:
[P2] Reconcile existing MTMVs when changing external type mappings. An MTMV
created before this migration can project an external BINARY or zoned timestamp
column while joining an OLAP table; it persists STRING or DATETIMEV2 in its
physical schema and keeps the original SELECT text. This loop flips the catalog
mapping without updating that MV. If a later ALTER of the OLAP table sets
SCHEMA_CHANGE, `MTMVTask` reanalyzes the unchanged SELECT;
`MTMVPlanUtil.validateColumns` rejects its now-VARBINARY output, while
`checkColumnIfChange` rejects old DATETIMEV2 versus new TIMESTAMPTZ. The MV can
no longer refresh through its normal schema-change path even if that ALTER left
the projected external column intact. Reconcile or explicitly cast affected
definitions and cover an upgraded MTMV refresh.
##########
fe/be-java-extensions/jdbc-scanner/src/main/java/org/apache/doris/jdbc/MySQLTypeHandler.java:
##########
@@ -192,6 +200,7 @@ public ColumnValueConverter getOutputConverter(ColumnType
columnType, String rep
@Override
public PreparedStatement initializeStatement(Connection conn, String sql,
int fetchSize) throws
SQLException {
+ initializeWriteConnection(conn);
Review Comment:
[P2] Preserve MySQL TIMESTAMP results during BE-first upgrades. Doris
upgrades BEs before FEs. While the old FE still maps a default MySQL TIMESTAMP
column to DATETIMEV2, this new read-path call resets every MySQL scan session
to UTC. The DATETIMEV2 decoder still reads session-local `LocalDateTime`
fields, so a TIMESTAMP storing 00:00Z on a +08 server changes from 08:00 on an
old BE to 00:00 on a new BE for the same query. Results can vary by BE during
rollout. Gate the UTC reset on a TIMESTAMPTZ scan contract, or preserve legacy
DATETIMEV2 interpretation, and cover old-FE/new-BE scans.
--
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]