github-actions[bot] commented on code in PR #68786:
URL: https://github.com/apache/doris/pull/68786#discussion_r4230876788
##########
fe/fe-connector/fe-connector-fluss/src/main/java/org/apache/doris/connector/fluss/FlussPartitionColumnTypes.java:
##########
@@ -72,12 +72,8 @@ static String rejection(DataType type,
FlussTypeMapping.Options options) {
return null;
case BINARY:
case BYTES:
- return options.isMapBinaryToVarbinary()
- ? "fluss names such a partition with the hex text of
the bytes, which this catalog"
- + " cannot read back as the VARBINARY column
that '"
- +
FlussCatalogProperties.ENABLE_MAPPING_VARBINARY + "=true' asks for;"
- + " turning that property off reads the
column, and the partition, as text"
- : null;
+ // Binary partition names are lossless hex; FlussPartitions
adds the VARBINARY wire prefix.
+ return null;
Review Comment:
[P2] Preserve binary partition keys in the Paimon lake half. This branch now
admits BINARY/BYTES partitions, including append-only Fluss tables with a
Paimon lake, but the Paimon sibling's `serializePartitionValue` throws for
BINARY/VARBINARY and `getPartitionInfoMap` drops the entire map. Because the
scan declares these as path partition keys, lake ranges then have no
`columnsFromPath` value and can return missing keys or lose rows under a key
predicate. Encode the lake partition bytes for the VARBINARY path, and cover a
binary-partitioned lake union with rows in both halves.
##########
fe/be-java-extensions/hadoop-hudi-scanner/src/main/java/org/apache/doris/hudi/HadoopHudiColumnValue.java:
##########
@@ -131,25 +145,9 @@ public LocalDate getDate() {
public LocalDateTime getDateTime() {
if (fieldData instanceof Timestamp) {
return ((Timestamp) fieldData).toLocalDateTime();
- } else if (fieldData instanceof TimestampWritableV2) {
- return
LocalDateTime.ofInstant(Instant.ofEpochSecond((((TimestampObjectInspector)
fieldInspector)
- .getPrimitiveJavaObject(fieldData)).toEpochSecond()),
zoneId);
- } else {
- long datetime = ((LongWritable) fieldData).get();
- long seconds;
- long nanoseconds;
- if (dorisType.getPrecision() == 3) {
- seconds = datetime / 1000;
- nanoseconds = (datetime % 1000) * 1000000;
- } else if (dorisType.getPrecision() == 6) {
- seconds = datetime / 1000000;
- nanoseconds = (datetime % 1000000) * 1000;
- } else {
- throw new RuntimeException("Hoodie timestamp only support
milliseconds and microseconds, "
- + "wrong precision = " + dorisType.getPrecision());
- }
- return LocalDateTime.ofInstant(Instant.ofEpochSecond(seconds,
nanoseconds), zoneId);
}
+ // DATETIMEV2 now denotes local-timestamp annotations: decode their
fields without a zone shift.
+ return getTimeStampTz();
Review Comment:
[P2] Preserve old-FE Hudi timestamp reads during a BE-first upgrade. The old
FE maps Avro `timestamp-millis`/`timestamp-micros` to DATETIMEV2, so its Hudi
scan still calls `getDateTime()` on the upgraded BE. For `LongWritable` and
`TimestampWritableV2`, this new delegation returns UTC fields, whereas the
previous branch converted the instant in `zoneId`; an instant at 00:00Z becomes
00:00 instead of 08:00 in an Asia/Shanghai session. Keep the session-zone
conversion for legacy DATETIMEV2 plans or version the scan contract, and cover
the mixed-version read.
##########
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 == '#'
Review Comment:
[P2] Respect MySQL's `--` comment rule when stripping the delimiter. MySQL
treats `balance--1` as subtraction because the second dash is followed by `1`,
but this branch consumes the rest of the line as a comment. A valid TIMESTAMPTZ
query TVF such as `SELECT ts, balance--1 AS n FROM t;` passes metadata
discovery; the scan wrapper then leaves `;` inside `FROM (...)` and fails.
Recognize `--` as a comment for MySQL/OceanBase only when followed by
whitespace/control, and cover this query shape.
##########
fe/be-java-extensions/max-compute-connector/src/main/java/org/apache/doris/maxcompute/MaxComputeJniWriter.java:
##########
@@ -566,8 +569,24 @@ private void fillArrowVectorStreaming(VectorSchemaRoot
root, int colIdx, OdpsTyp
vec.setValueCount(numRows);
break;
}
- case DATETIME:
case TIMESTAMP: {
+ // TIMESTAMPTZ's JNI carrier is UTC and must retain
microseconds on write.
+ org.apache.arrow.vector.TimeStampVector vec =
+ (org.apache.arrow.vector.TimeStampVector)
root.getVector(colIdx);
+ vec.allocateNew(numRows);
+ for (int i = 0; i < numRows; i++) {
+ if (vc.isNullAt(rowOffset + i)) {
+ vec.setNull(i);
+ } else {
+ LocalDateTime utc = vc.getTimeStampTz(rowOffset + i);
Review Comment:
[P2] Reject or preserve legacy MaxCompute TIMESTAMP write plans during
BE-first upgrades. The immediate pre-PR FE sends TIMESTAMP as DATETIMEV2 and
already supplies `txn_id`, so it passes the constructor guard. This branch
interprets its local 08:00 fields as 08:00Z instead of the previous 00:00Z on a
+08 BE, while that FE's write session also uses MILLI rather than the new MICRO
unit. Depending on SDK enforcement, the write fails or stores the wrong
instant. Select conversion and units from a compatible plan contract, or reject
it before committing; cover the old-FE/new-BE case.
--
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]