github-actions[bot] commented on code in PR #68786:
URL: https://github.com/apache/doris/pull/68786#discussion_r4217826031


##########
be/test/exec/sink/writer/iceberg/viceberg_table_writer_test.cpp:
##########
@@ -27,14 +29,227 @@
 #include "core/data_type/data_type_nullable.h"
 #include "core/data_type/data_type_number.h"
 #include "core/data_type/data_type_struct.h"
+#include "core/data_type/data_type_timestamptz.h"
+#include "core/data_type/data_type_varbinary.h"
+#include "exec/sink/writer/iceberg/partition_transformers.h"
+#include "exec/sink/writer/iceberg/vpartition_writer_base.h"
 #include "exprs/vexpr_context.h"
 #include "exprs/vslot_ref.h"
 #include "format/table/iceberg/partition_spec_parser.h"
 #include "format/table/iceberg/schema.h"
 #include "format/table/iceberg/types.h"
+#include "runtime/runtime_state.h"
 
 namespace doris {
 
+namespace {
+class RecordingPartitionWriter : public IPartitionWriterBase {
+public:
+    Status open(RuntimeState*, RuntimeProfile*, const RowDescriptor*) override 
{
+        return Status::OK();
+    }
+    Status write(Block& block) override {
+        rows += block.rows();
+        for (size_t i = 0; i < block.rows(); ++i) {
+            EXPECT_TRUE(block.get_by_position(0).column->is_null_at(i));
+        }
+        return Status::OK();
+    }
+    Status close(const Status& status) override { return status; }
+    const std::string& file_name() const override { return name; }
+    int file_name_index() const override { return 0; }
+    size_t written_len() const override { return 0; }
+    size_t rows = 0;
+    const std::string name = "part";
+};
+} // namespace
+
+TEST(VIcebergTableWriterTest, TimestampIdentityPreservesUtcCommitAndNull) {
+    std::vector<iceberg::NestedField> fields;
+    fields.emplace_back(false, 1, "event_time", 
std::make_unique<iceberg::TimestampType>(true),
+                        std::nullopt);
+    auto schema = std::make_shared<iceberg::Schema>(std::move(fields));
+    auto spec = iceberg::PartitionSpecParser::from_json(
+            schema,
+            
R"({"spec-id":0,"fields":[{"name":"event_time","transform":"identity","source-id":1,"field-id":1000}]})");
+    auto type = std::make_shared<DataTypeTimeStampTz>(6);
+    TIcebergTableSink sink;
+    TDataSink data_sink;
+    data_sink.__set_iceberg_table_sink(sink);
+    VExprContextSPtrs exprs;
+    VIcebergTableWriter writer(data_sink, exprs, nullptr, nullptr);
+    RuntimeState state;
+    state.set_timezone("America/New_York");
+    writer._state = &state;
+    writer._schema = schema;
+    writer._iceberg_partition_columns.emplace_back(
+            spec->fields()[0], TYPE_TIMESTAMPTZ, 0, std::vector<size_t> {},
+            std::make_unique<IdentityPartitionColumnTransform>(type));
+    auto column = ColumnTimeStampTz::create();
+    TimestampTzValue first;
+    first.unchecked_set_time(2021, 11, 7, 5, 30, 0, 123456);
+    TimestampTzValue second;
+    second.unchecked_set_time(2021, 11, 7, 6, 30, 0, 123456);
+    column->get_data().assign({first, second, TimestampTzValue()});
+    auto null_map = ColumnUInt8::create();
+    null_map->get_data().assign({0, 0, 1});
+    ColumnWithTypeAndName partition(ColumnNullable::create(std::move(column), 
std::move(null_map)),
+                                    make_nullable(type), "event_time");
+    std::vector<std::string> paths;
+    for (int row = 0; row < 3; ++row) {
+        auto value = writer._get_iceberg_partition_value(TYPE_TIMESTAMPTZ, 
partition, row);
+        IcebergPartitionData data({value});
+        paths.push_back(writer._partition_to_path(data));
+        const std::vector<std::string> expected = {"2021-11-07 
05:30:00.123456+00:00",
+                                                   "2021-11-07 
06:30:00.123456+00:00", "null"};
+        EXPECT_EQ(std::vector<std::string>({expected[row]}), 
writer._partition_values(data));
+    }
+    EXPECT_NE(paths[0], paths[1]);
+    EXPECT_EQ("event_time=null", paths[2]);
+    // Static literals with an offset must route to exactly the same partition 
as dynamic rows.
+    writer._t_sink.iceberg_table_sink.__set_static_partition_values(
+            {{"event_time", "2021-11-07 01:30:00.123456-05:00"}});
+    writer._init_static_partition_values();
+    EXPECT_EQ(paths[1], writer._static_partition_path);
+    EXPECT_EQ(std::vector<std::string>({"2021-11-07 06:30:00.123456+00:00"}),
+              writer._static_partition_value_list);
+}
+
+TEST(VIcebergTableWriterTest, 
StaticNullBinaryPartitionsKeepNullPathsAndCommitValues) {
+    for (int kind = 0; kind < 3; ++kind) {
+        SCOPED_TRACE(kind);
+        std::unique_ptr<iceberg::Type> source_type;
+        if (kind == 0) {
+            source_type = std::make_unique<iceberg::BinaryType>();
+        } else if (kind == 1) {
+            source_type = std::make_unique<iceberg::FixedType>(16);
+        } else {
+            source_type = std::make_unique<iceberg::UUIDType>();
+        }
+        std::vector<iceberg::NestedField> fields;
+        fields.emplace_back(true, 1, "key", std::move(source_type), 
std::nullopt);
+        fields.emplace_back(true, 2, "id", 
std::make_unique<iceberg::IntegerType>(), std::nullopt);
+        auto schema = std::make_shared<iceberg::Schema>(std::move(fields));
+        auto spec = iceberg::PartitionSpecParser::from_json(
+                schema,
+                
R"({"spec-id":0,"fields":[{"name":"key","transform":"identity","source-id":1,"field-id":1000}]})");
+        // Encode the nullable wire contract independently of the generated 
reader's schema.
+        using namespace apache::thrift::protocol;
+        auto buffer = 
std::make_shared<apache::thrift::transport::TMemoryBuffer>();
+        TBinaryProtocol protocol(buffer);
+        protocol.writeStructBegin("TIcebergTableSink");
+        protocol.writeFieldBegin("static_partition_values", T_MAP, 15);
+        protocol.writeMapBegin(T_STRING, T_STRING, 1);
+        protocol.writeString(std::string("key"));
+        protocol.writeString(std::string("null"));
+        protocol.writeMapEnd();
+        protocol.writeFieldEnd();
+        protocol.writeFieldBegin("static_partition_null_keys", T_SET, 19);

Review Comment:
   [P1] Encode the NULL-key set with Thrift field 20. Field 19 is already 
`nan_count_field_ids` (`list<i32>`), so the generated `TIcebergTableSink` 
reader skips this `T_SET` and leaves `static_partition_null_keys` empty. The 
writer then parses the placeholder `null` as VARBINARY hex and 
`_init_static_partition_values` throws, failing this new test; it never 
verifies the intended FE-to-BE NULL marker. Use field 20 or serialize a 
generated sink.



##########
be/src/exec/sink/writer/vjdbc_table_writer.cpp:
##########
@@ -50,6 +50,8 @@ std::map<std::string, std::string> 
VJdbcTableWriter::_build_writer_params(const
     params["jdbc_driver_url"] = driver_url;
 
     params["jdbc_driver_checksum"] = 
t_jdbc_sink.jdbc_table.jdbc_driver_checksum;
+    // The unified JNI writer needs the dialect to preserve instant semantics 
in parameter binds.
+    params["table_type"] = to_string(t_jdbc_sink.table_type);

Review Comment:
   [P1] Send the JDBC dialect name to the JNI writer. `TOdbcTableType` is a 
numeric Thrift enum, so `to_string` supplies values such as `0` for MYSQL, but 
`JdbcTypeHandlerFactory.create` matches names such as `MYSQL` and otherwise 
constructs `DefaultTypeHandler`. JDBC sinks therefore skip dialect-specific 
handling: MySQL TIMESTAMPTZ writes miss UTC session setup and use the generic 
bind, which can shift instants under a non-UTC session. Reuse the enum-to-name 
mapping used by `JdbcScanner` and test the BE-generated parameter.



##########
fe/fe-connector/fe-connector-jdbc/src/main/java/org/apache/doris/connector/jdbc/JdbcQueryBuilder.java:
##########
@@ -190,6 +191,72 @@ 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(";+$", "");
+        return "SELECT " + projections + " FROM (" + inner + ") 
doris_jdbc_query";

Review Comment:
   [P2] Keep a newline before closing a wrapped TVF query. For valid input 
ending in a `--` line comment and newline, `trim()` removes that newline and 
the generated `)` lands inside the comment. ClickHouse/Trino/Presto TVFs with a 
TIMESTAMPTZ output now send invalid SQL although the raw query worked before. 
Insert a line break before the closing parenthesis and cover a terminal line 
comment.



##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanPlanProvider.java:
##########
@@ -1927,6 +1933,45 @@ private static boolean projectsVariant(
                 .anyMatch(index -> containsVariant(rowType.getTypeAt(index)));
     }
 
+    static boolean requiresLegacyOrcTimestampReader(Table table, 
Optional<List<RawFile>> rawFiles,
+            Map<Long, Boolean> schemaTimestamps) {
+        if (!rawFiles.isPresent() || rawFiles.get().stream().noneMatch(f -> 
f.path().endsWith(".orc"))
+                || !new org.apache.paimon.options.Options(table.options()).get(
+                        
org.apache.paimon.format.OrcOptions.ORC_TIMESTAMP_LTZ_LEGACY_TYPE)) {
+            return false;
+        }
+        // Legacy ORC LTZ bytes require the SDK's JVM-zone conversion, 
including old nested schemas.
+        if (containsTimestampLtz(table.rowType())) {

Review Comment:
   [P2] Scope the legacy ORC fallback to LTZ fields this scan reads. 
`planScanInternal` already computes projected fields, but this checks the whole 
table: with an unused LTZ column and `orc.timestamp-ltz.legacy-type` enabled, 
even `SELECT id` takes JNI and loses native file splitting. A query that also 
selects a Paimon metadata column then fails because the JNI arm calls 
`validateMetadataColumnReader(true, false)`, although its native scan would 
read no LTZ bytes. Check projected and filter fields, including historical file 
schemas, before forcing JNI.



-- 
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