github-actions[bot] commented on code in PR #68786:
URL: https://github.com/apache/doris/pull/68786#discussion_r4226739807
##########
fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergPartitionUtils.java:
##########
@@ -478,10 +516,12 @@ private static String
normalizeFloatingPointPartitionValue(String valueStr) {
}
private static long parseTimestampToMicros(String valueStr, TimestampType
timestampType, ZoneId sessionZone) {
- LocalDateTime ldt = LocalDateTime.parse(valueStr,
TIMESTAMP_PARTITION_FORMAT);
- // timestamptz (shouldAdjustToUTC): interpret the wall-clock string in
the session zone; plain timestamp:
- // interpret it in UTC. Mirrors legacy parseTimestampToMicros
(DateUtils.getTimeZone vs ZoneId.of("UTC")).
- ZoneId zone = timestampType.shouldAdjustToUTC() ? sessionZone :
ZoneOffset.UTC;
+ java.time.temporal.TemporalAccessor parsed =
TIMESTAMP_PARTITION_FORMAT.parse(valueStr.replace('T', ' '));
+ LocalDateTime ldt = LocalDateTime.from(parsed);
+ ZoneId explicitZone =
parsed.query(java.time.temporal.TemporalQueries.zone());
+ // Dynamic commits carry an offset; only unqualified static literals
use the session zone.
+ ZoneId zone = timestampType.shouldAdjustToUTC()
+ ? (explicitZone != null ? explicitZone : sessionZone) :
ZoneOffset.UTC;
Review Comment:
[P2] Resolve skipped local times identically for overwrite filtering and
writing. With session zone America/New_York, static TIMESTAMPTZ partition
literal '2021-03-14 02:30:00' falls in the spring gap. This FE atZone call
advances it to 03:30 EDT (07:30 UTC) for overwriteByRowFilter, while the BE
static writer's cctz::convert resolves the same literal to the transition
instant, 03:00 EDT (07:00 UTC). The overwrite filter therefore targets a
different partition from the file being added; an overwrite can fail or leave
the intended old partition in place. Resolve once and pass the UTC value to
both sides, or align gap handling, with a spring-gap overwrite test.
##########
regression-test/suites/external_table_p0/iceberg/test_iceberg_initial_defaults.groovy:
##########
@@ -528,6 +528,9 @@ suite("test_iceberg_initial_defaults",
"p0,external,nonConcurrent") {
sql """switch ${legacyCatalog}"""
sql """use ${namespace}"""
+ // Disabling the retired toggle must preserve the zoned type and
its UTC instant.
+ assertEquals("timestamptz(6)",
Review Comment:
[P3] Capture this fixed DESC result in the generated output. The new
assertEquals expects `timestamptz(6)` from `DESC ${tableName}`, but this
suite's .out has no corresponding schema result. Repository regression-test
rules require deterministic results to use qt/order_qt with runner-generated
output. Please add a named qt check and regenerate this suite's .out through
the runner. This initial-defaults suite is separate from the suites in the
existing regression-style thread and the static-overwrite comment.
##########
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))) {
Review Comment:
[P2] Preserve persisted view output types during catalog mapping migration.
A Doris view created over an external BINARY column while the old marker was
false persists STRING in View.fullSchema, but this migration forces the catalog
to map the source as VARBINARY. Each SELECT reparses the stored view SQL, and
LogicalView.computeOutput keeps the new child slot dataType while SHOW COLUMNS
still reads the old schema. Thus the same view advertises STRING and returns
VARBINARY (and DATETIMEV2/TIMESTAMPTZ has the same problem), which can change
or break dependent queries. Revalidate/migrate affected views or cast expanded
outputs to their stored schema, and cover a view created before the marker flip.
##########
regression-test/suites/external_table_p0/paimon/paimon_timestamp_types.groovy:
##########
@@ -150,28 +152,53 @@ suite("paimon_timestamp_types", "p0,external") {
test_scale("2024-01-02 10:04:05.123456")
// test_ltz_ntz("test_timestamp_ntz_ltz_orc")
// test_ltz_ntz("test_timestamp_ntz_ltz_parquet")
- test_ltz_ntz_simple("test_timestamp_ntz_ltz_simple_orc", "2024-01-02
10:12:34.123456")
+ test_ltz_ntz_simple("test_timestamp_ntz_ltz_simple_orc", "2024-01-02
10:12:34.123456+08:00")
// test_ltz_ntz_simple("test_timestamp_ntz_ltz_simple_parquet")
// Native ORC rounds scales 7-9 to microseconds; JNI truncates. Native
Parquet uses
// Paimon history semantics for high-precision INT96 instead of the
session timezone.
sql """set force_jni_scanner=false"""
+ // Precision-aware native predicates use V2; V1 compares raw
nanoseconds before truncation.
+ // SELECT * still falls back to JNI because it includes legacy ORC LTZ
fields.
+ sql """set enable_file_scanner_v2=true"""
test_scale("2024-01-02 10:04:05.123457")
+ // Keep testing legacy LTZ fallback with V2 disabled independently of
the NTZ predicate tests.
+ sql """set enable_file_scanner_v2=false"""
// test_ltz_ntz("test_timestamp_ntz_ltz_orc")
// test_ltz_ntz("test_timestamp_ntz_ltz_parquet")
- test_ltz_ntz_simple("test_timestamp_ntz_ltz_simple_orc", "2024-01-02
02:12:34.123456")
- test_ltz_ntz_simple("test_timestamp_ntz_ltz_simple_parquet",
"2024-01-02 10:12:34.123456")
+ test_ltz_ntz_simple("test_timestamp_ntz_ltz_simple_orc", "2024-01-02
10:12:34.123456+08:00")
+ test_ltz_ntz_simple("test_timestamp_ntz_ltz_simple_parquet",
"2024-01-02 10:12:34.123456+08:00")
+
+ // Legacy ORC LTZ needs Paimon's compatibility conversion. Reader
selection must
+ // preserve LTZ instants and NTZ civil fields, including inside
containers.
+ for (def zone : ["UTC", "Asia/Shanghai", "America/New_York"]) {
+ sql "set time_zone = '${zone}'"
+ for (def scannerV2 : [false, true]) {
+ sql "set enable_file_scanner_v2 = ${scannerV2}"
+ sql "set force_jni_scanner = true"
+ def scalarSql = """select cast(ts6 as string), cast(ts16 as
string),
+ unix_timestamp(ts16) from ts_scale_orc"""
+ def nestedSql = """select cast(cmap1 as string), cast(cmap2 as
string),
+ cast(carray1 as string), cast(carray2 as string),
cast(crow as string)
+ from test_timestamp_ntz_ltz_simple_orc"""
+ def expectedScalar = sql(scalarSql)
+ def expectedNested = sql(nestedSql)
+ assertEquals("2024-01-02 10:04:05.123456",
expectedScalar[0][0])
Review Comment:
[P3] Record the fixed Paimon scalar expectations in generated output. These
two new assertEquals checks hardcode the wall time and epoch returned by
scalarSql for every zone and reader setting, but the suite's .out has no result
for those new queries. Repository regression rules require fixed expected query
results in named qt/order_qt checks with runner-generated output. Please add
generated entries for the fixed values; the later JNI-versus-native equality
assertions can stay as dynamic comparisons. This Paimon suite is distinct from
the existing regression-style thread.
##########
be/src/exec/sink/writer/iceberg/viceberg_table_writer.cpp:
##########
@@ -251,7 +254,74 @@ void VIcebergTableWriter::_init_static_partition_values() {
auto it = static_values_map.find(col_name);
if (it != static_values_map.end()) {
_partition_column_static_values[i] = it->second;
+ _partition_column_static_path_values[i] = it->second;
_partition_column_is_static[i] = 1;
+ auto type =
+
_iceberg_partition_columns[i].partition_column_transform().get_result_type();
+ if (iceberg_sink.static_partition_null_keys.contains(col_name)) {
+ // SQL NULL is not a UUID/hex string, and must not share the
empty-byte path.
+ _partition_column_static_values[i] =
+ _iceberg_partition_columns[i]
+ .partition_column_transform()
+ .get_partition_value(type, std::any {});
+ _partition_column_static_path_values[i] =
+ _partition_value_to_human_string(i, std::any {});
+ continue;
+ }
+ if (type->get_primitive_type() == TYPE_VARBINARY) {
+ // Static and dynamic partitions share typed hex commit
values, but Iceberg paths
+ // render binary as base64 (UUID as canonical text). Decode
before rendering.
+ // Nullable SerDe turns malformed bytes into NULL. Non-null
wire values must fail
+ // validation instead of being silently rendered as an empty
binary partition.
+ auto non_null_type = remove_nullable(type);
+ auto column = non_null_type->create_column();
+ auto& encoded = _partition_column_static_values[i];
+ if
(_schema->find_type(_iceberg_partition_columns[i].field().source_id())
+ ->type_id() == iceberg::TypeID::UUID &&
+ !encoded.starts_with("0x")) {
+ // Preserve the canonical UUID string input accepted
before VARBINARY mapping.
+ boost::uuids::uuid uuid;
+ try {
+ uuid = boost::uuids::string_generator()(encoded);
+ } catch (const std::runtime_error&) {
+ throw Exception(ErrorCode::INVALID_ARGUMENT,
+ "Invalid UUID partition value");
+ }
+ encoded = _iceberg_partition_columns[i]
+ .partition_column_transform()
+ .get_partition_value(type,
+
std::string(uuid.begin(), uuid.end()));
+ }
+ StringRef value(encoded.data(), encoded.size());
+ DataTypeSerDe::FormatOptions options;
+ auto status = non_null_type->get_serde()->from_string(value,
*column, options);
Review Comment:
[P1] Keep static binary partition bytes consistent with the materialized
row. `PARTITION(part_key='0xDEAD')` passes the literal and identity-field
checks, and BindSink casts that StringLiteral to VARBINARY as the six ASCII
bytes `30 78 44 45 41 44` in the data row. This new `from_string` call instead
hex-decodes the same text to `DE AD` for the Iceberg path and committed
partition metadata. Full-static and hybrid writes can therefore commit a file
whose BINARY rows disagree with its partition value; Iceberg readers that prune
by partition can silently miss those rows. The `X'DEAD'` regression uses a
typed literal and does not cover this path. Normalize the value once for both
row and partition, or reject quoted-hex static specs before writing, and cover
this case.
##########
fe/be-java-extensions/jdbc-scanner/src/main/java/org/apache/doris/jdbc/MySQLTypeHandler.java:
##########
@@ -334,4 +343,31 @@ private java.lang.reflect.Type
getListTypeForArray(ColumnType type) {
throw new IllegalArgumentException("Unsupported array child
type: " + type.getType());
}
}
+
+ @Override
+ public void setTimestampTz(java.sql.PreparedStatement statement, int
parameterIndex, LocalDateTime value)
+ throws SQLException {
+ if (usesMySqlTimestampProtocol()) {
+ // Connector/J 5.x can ignore Calendar during timestamp
conversion. Bind the UTC
+ // session's civil fields directly so neither driver nor JVM
applies another offset.
+ statement.setString(parameterIndex, value.toString().replace('T',
' '));
+ } else {
+ statement.setObject(parameterIndex,
java.sql.Timestamp.from(value.toInstant(ZoneOffset.UTC)));
+ }
+ }
+
+ private boolean usesMySqlTimestampProtocol() {
+ return "MYSQL".equals(tableType) || "OCEANBASE".equals(tableType);
+ }
+
+ @Override
+ public void initializeWriteConnection(Connection connection) throws
SQLException {
+ if (usesMySqlTimestampProtocol()) {
+ // Reset every pool checkout: TIMESTAMP travels as session-local
fields on the wire.
+ try (Statement statement = connection.createStatement()) {
+ statement.execute("SET SESSION time_zone = '+00:00'");
Review Comment:
[P2] Align Connector/J's cached timezone with the forced UTC session. With
Connector/J 5.1.49 and `useTimezone=true`, the driver caches the original
server timezone when the connection opens. This new SET changes the SQL session
to UTC, but the TIMESTAMPTZ read uses `getObject(..., LocalDateTime.class)`,
which delegates to `getTimestamp`; its legacy path still applies the cached old
zone. A connection initially in +08 with a UTC BE JVM can therefore return a
MySQL TIMESTAMP instant eight hours off after this reset. Read the UTC session
fields without the driver's cached-zone conversion or initialize both zones
consistently, and cover this supported driver configuration in a scan
regression.
##########
regression-test/suites/external_table_p0/iceberg/write/test_iceberg_static_partition_overwrite.groovy:
##########
@@ -48,6 +49,44 @@ suite("test_iceberg_static_partition_overwrite",
"p0,external") {
sql """ create database ${db1} """
sql """ use ${db1} """
+ // Binary partition values must retain their bytes across dynamic,
full-static and hybrid writes.
+ String binaryTable = "binary_partition_overwrite"
+ spark_iceberg """
+ CREATE TABLE demo.${db1}.${binaryTable} (id INT, part_key BINARY,
region STRING)
+ USING iceberg PARTITIONED BY (part_key, region)
+ TBLPROPERTIES ('write.format.default'='parquet')
+ """
+ sql """ INSERT INTO ${binaryTable} VALUES
+ (1, X'DEAD', 'a'), (2, X'DEAD', 'b'), (3, X'00FF', 'a') """
+ sql """ INSERT OVERWRITE TABLE ${binaryTable}
+ PARTITION (part_key=X'DEAD', region='a') SELECT 10 """
+ assertEquals([[2, "DEAD", "b"], [3, "00FF", "a"], [10, "DEAD", "a"]],
Review Comment:
[P3] Record these fixed overwrite results in the generated output file. The
four new assertEquals checks at lines 63, 68, 79, and 84 compare deterministic
ordered SELECT rows, but this suite's .out file has no corresponding entries.
Repository regression-test rules require qt/order_qt with runner-generated .out
for determined results. Please convert these checks and regenerate the expected
output through the test runner. This is a separate suite from the existing
regression-style thread.
--
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]