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]