Gabriel39 commented on code in PR #68786:
URL: https://github.com/apache/doris/pull/68786#discussion_r4218487765


##########
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:
   Fixed the wire fixture to encode `static_partition_null_keys` as field 20, 
with explicit assertions that the generated reader decoded the NULL-key set. 
Field 19 remains reserved for NaN counts.
   
   Addressed in f2c7b1d41ad2ba9d5ab12e2ee3b6859123ba5563.



##########
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:
   The unqualified call resolves to Thrift's generated `doris::to_string(const 
TOdbcTableType::type&)`, declared in `Types_types.h`. That overload uses 
`_TOdbcTableType_VALUES_TO_NAMES` and returns names such as `MYSQL`, `TRINO`, 
and `PRESTO`; only unknown enum values fall back to numbers. Added 
`VJdbcTableWriterTest.ParametersCarryDialectNamesForJavaDispatch` to cover the 
actual BE-generated parameter map for all 12 dialects. The production mapping 
is already correct and is unchanged. The new BE test passes syntax compilation; 
full execution is pending CI because the local BE build fails in the unchanged 
wide-integer implementation.
   
   Addressed in f2c7b1d41ad2ba9d5ab12e2ee3b6859123ba5563.



##########
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:
   Fixed by placing a newline before the derived table's closing parenthesis. 
The regression test reproduced the failure with a terminal `--` comment and now 
passes for ClickHouse, Trino, and Presto. Both JDBC test classes pass (35 
tests).
   
   Addressed in f2c7b1d41ad2ba9d5ab12e2ee3b6859123ba5563.



##########
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:
   Scoped the legacy ORC fallback to projected and filter-referenced fields, 
using stable field IDs when inspecting historical schemas. Added a real ORC 
scan-planning test covering native metadata-column scans with an unread LTZ 
column, JNI for projected/filter-referenced LTZ, and native scans after the LTZ 
column is dropped. The new test reproduced the metadata-column error before the 
fix; both Paimon test classes now pass (105 tests).
   
   Addressed in f2c7b1d41ad2ba9d5ab12e2ee3b6859123ba5563.



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