From 8e26d921fb73e3dab9567060e396651836693089 Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Mon, 14 Sep 2026 16:24:31 +0300 Subject: [PATCH 1/4] Support Trino timestamp(3) writes --- ydb-trino-adapter/README.md | 3 ++- ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java | 3 ++- .../src/test/java/tech/ydb/trino/TestYdbConnectorTest.java | 4 +++- 3 files changed, 7 insertions(+), 3 deletions(-) diff --git a/ydb-trino-adapter/README.md b/ydb-trino-adapter/README.md index 3a3f6816..ba77df2b 100644 --- a/ydb-trino-adapter/README.md +++ b/ydb-trino-adapter/README.md @@ -54,6 +54,7 @@ YDB `Text` отображается в Trino как `varchar`, а `Bytes` — к Новые столбцы Trino `date` создаются как YDB `Date32`; адаптер также включает JDBC-параметр `forceSignedDatetimes=true`. Существующие столбцы YDB `Date` и `Date32` читаются как Trino `date`. -Новые столбцы Trino `timestamp(6)` создаются как YDB `Timestamp64`. +Новые столбцы Trino `timestamp(3)` и `timestamp(6)` создаются как YDB `Timestamp64`. +При чтении физический `Timestamp64` описывается как Trino `timestamp(6)`. Это не исправляет диапазонные фильтры и `UPDATE` существующих столбцов YDB `Date`. Не задавайте `forceSignedDatetimes=false` в JDBC URL: параметры URL имеют приоритет над настройками адаптера. diff --git a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java index 3a4822e0..44bd4c90 100644 --- a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java +++ b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java @@ -116,6 +116,7 @@ import static io.trino.spi.type.IntegerType.INTEGER; import static io.trino.spi.type.RealType.REAL; import static io.trino.spi.type.SmallintType.SMALLINT; +import static io.trino.spi.type.TimestampType.TIMESTAMP_MILLIS; import static io.trino.spi.type.TimestampType.TIMESTAMP_MICROS; import static io.trino.spi.type.TinyintType.TINYINT; import static io.trino.spi.type.VarbinaryType.VARBINARY; @@ -435,7 +436,7 @@ public WriteMapping toWriteMapping(ConnectorSession session, Type type) { if (type == DATE) { return WriteMapping.longMapping("Date32", dateWriteFunctionUsingLocalDate()); } - if (type == TIMESTAMP_MICROS) { + if (type == TIMESTAMP_MILLIS || type == TIMESTAMP_MICROS) { return WriteMapping.longMapping("Timestamp64", timestamp64WriteFunction()); } diff --git a/ydb-trino-adapter/src/test/java/tech/ydb/trino/TestYdbConnectorTest.java b/ydb-trino-adapter/src/test/java/tech/ydb/trino/TestYdbConnectorTest.java index cec9b9d5..9635310b 100644 --- a/ydb-trino-adapter/src/test/java/tech/ydb/trino/TestYdbConnectorTest.java +++ b/ydb-trino-adapter/src/test/java/tech/ydb/trino/TestYdbConnectorTest.java @@ -125,7 +125,9 @@ protected Optional filterDataMappingSmokeTestData(BaseConn String trinoTypeName = dataMappingTestSetup.getTrinoTypeName(); if (trinoTypeName.equals("char(3)")) { return Optional.of(dataMappingTestSetup.asUnsupported()); - } else if (trinoTypeName.startsWith("time") && !trinoTypeName.equals("timestamp(6)")) { + } else if (trinoTypeName.startsWith("time") + && !trinoTypeName.equals("timestamp") + && !trinoTypeName.equals("timestamp(6)")) { // Нет time в YQL return Optional.empty(); } From 9b4eb1f55b3d170ea37a68b1e39f066eee7e607d Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Mon, 14 Sep 2026 16:54:17 +0300 Subject: [PATCH 2/4] Bind YDB temporal values with native JDBC types --- ydb-trino-adapter/README.md | 12 ++-- .../main/java/tech/ydb/trino/YdbClient.java | 60 +++++++++++++++---- .../java/tech/ydb/trino/YdbClientModule.java | 5 -- 3 files changed, 55 insertions(+), 22 deletions(-) diff --git a/ydb-trino-adapter/README.md b/ydb-trino-adapter/README.md index ba77df2b..688b22d7 100644 --- a/ydb-trino-adapter/README.md +++ b/ydb-trino-adapter/README.md @@ -52,9 +52,13 @@ YDB `Text` отображается в Trino как `varchar`, а `Bytes` — к ## Даты -Новые столбцы Trino `date` создаются как YDB `Date32`; адаптер также включает JDBC-параметр `forceSignedDatetimes=true`. +Новые столбцы Trino `date` создаются как YDB `Date32`. Существующие столбцы YDB `Date` и `Date32` читаются как Trino `date`. Новые столбцы Trino `timestamp(3)` и `timestamp(6)` создаются как YDB `Timestamp64`. -При чтении физический `Timestamp64` описывается как Trino `timestamp(6)`. -Это не исправляет диапазонные фильтры и `UPDATE` существующих столбцов YDB `Date`. -Не задавайте `forceSignedDatetimes=false` в JDBC URL: параметры URL имеют приоритет над настройками адаптера. +Существующие YDB `Datetime`, `Datetime64`, `Timestamp` и `Timestamp64` читаются как Trino `timestamp(6)`; +исходная точность `timestamp(3)` в метаданных не сохраняется. +При записи адаптер передаёт YDB JDBC точный vendor type по `TYPE_NAME` и больше не задаёт `forceSignedDatetimes`: +`Date`/`Date32` получают дни от эпохи, `Datetime`/`Datetime64` — секунды UTC, а `Timestamp`/`Timestamp64` — `Instant` с микросекундами. +В `UPDATE`/`DELETE`, которые стандартный merge sink выполняет внутри `MERGE`, native `TYPE_NAME` не передаётся; +fallback `Instant` для `Datetime`/`Datetime64` может зависеть от часового пояса JVM и остаётся отдельным риском. +Диапазонные предикаты вне диапазона legacy YDB `Date` остаются отдельным известным ограничением. diff --git a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java index 44bd4c90..ab9126a5 100644 --- a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java +++ b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java @@ -84,7 +84,6 @@ import static io.trino.plugin.jdbc.StandardColumnMappings.bigintColumnMapping; import static io.trino.plugin.jdbc.StandardColumnMappings.bigintWriteFunction; import static io.trino.plugin.jdbc.StandardColumnMappings.booleanColumnMapping; -import static io.trino.plugin.jdbc.StandardColumnMappings.dateWriteFunctionUsingLocalDate; import static io.trino.plugin.jdbc.StandardColumnMappings.dateReadFunctionUsingLocalDate; import static io.trino.plugin.jdbc.StandardColumnMappings.decimalColumnMapping; import static io.trino.plugin.jdbc.StandardColumnMappings.doubleColumnMapping; @@ -100,7 +99,6 @@ import static io.trino.plugin.jdbc.StandardColumnMappings.smallintWriteFunction; import static io.trino.plugin.jdbc.StandardColumnMappings.timestampColumnMapping; import static io.trino.plugin.jdbc.StandardColumnMappings.timestampReadFunction; -import static io.trino.plugin.jdbc.StandardColumnMappings.timestampWriteFunction; import static io.trino.plugin.jdbc.StandardColumnMappings.tinyintWriteFunction; import static io.trino.plugin.jdbc.StandardColumnMappings.varbinaryColumnMapping; import static io.trino.plugin.jdbc.StandardColumnMappings.varbinaryWriteFunction; @@ -126,11 +124,19 @@ import static java.lang.String.format; import static java.time.ZoneOffset.UTC; import static java.util.stream.Collectors.joining; +import static tech.ydb.jdbc.YdbConst.SQL_KIND_PRIMITIVE; public class YdbClient extends BaseJdbcClient { static final String DEFAULT_SCHEMA = "default"; private static final int YDB_DEFAULT_DECIMAL_PRECISION = 22; private static final int YDB_DEFAULT_DECIMAL_SCALE = 9; + // https://github.com/ydb-platform/ydb-jdbc-driver/blob/v2.4.1/jdbc/src/main/java/tech/ydb/jdbc/common/YdbTypes.java#L64-L77 + private static final int YDB_DATE_SQL_TYPE = SQL_KIND_PRIMITIVE + 16; + private static final int YDB_DATETIME_SQL_TYPE = SQL_KIND_PRIMITIVE + 17; + private static final int YDB_TIMESTAMP_SQL_TYPE = SQL_KIND_PRIMITIVE + 18; + private static final int YDB_DATE32_SQL_TYPE = SQL_KIND_PRIMITIVE + 25; + private static final int YDB_DATETIME64_SQL_TYPE = SQL_KIND_PRIMITIVE + 26; + private static final int YDB_TIMESTAMP64_SQL_TYPE = SQL_KIND_PRIMITIVE + 27; private final ConnectorExpressionRewriter connectorExpressionRewriter; private final AggregateFunctionRewriter aggregateFunctionRewriter; @@ -342,7 +348,7 @@ public Optional toColumnMapping( : typeHandle.columnSize().orElse(VarcharType.MAX_LENGTH); yield Optional.of(varcharColumnMapping(length)); } - case Types.DATE -> Optional.of(dateColumnMapping()); + case Types.DATE -> Optional.of(dateColumnMapping(typeHandle)); case Types.TIMESTAMP -> Optional.of(timestampColumnMapping(typeHandle)); default -> Optional.empty(); }; @@ -374,28 +380,56 @@ private static ColumnMapping varcharColumnMapping(int varcharLength) { FULL_PUSHDOWN); } - private static ColumnMapping dateColumnMapping() { + private static ColumnMapping dateColumnMapping(JdbcTypeHandle typeHandle) { + String typeName = typeHandle.jdbcTypeName().orElse(""); + LongWriteFunction writeFunction = switch (typeName.toLowerCase(Locale.ROOT)) { + case "date" -> dateWriteFunction(YDB_DATE_SQL_TYPE); + case "date32" -> dateWriteFunction(YDB_DATE32_SQL_TYPE); + default -> throw new TrinoException(NOT_SUPPORTED, "Unsupported YDB date type: " + typeName); + }; return ColumnMapping.longMapping( DATE, dateReadFunctionUsingLocalDate(), - dateWriteFunctionUsingLocalDate()); + writeFunction); } private static ColumnMapping timestampColumnMapping(JdbcTypeHandle typeHandle) { String typeName = typeHandle.jdbcTypeName().orElse(""); - LongWriteFunction writeFunction = typeName.equalsIgnoreCase("Timestamp64") - ? timestamp64WriteFunction() - : timestampWriteFunction(TIMESTAMP_MICROS); + LongWriteFunction writeFunction = switch (typeName.toLowerCase(Locale.ROOT)) { + case "datetime" -> datetimeWriteFunction(YDB_DATETIME_SQL_TYPE); + case "datetime64" -> datetimeWriteFunction(YDB_DATETIME64_SQL_TYPE); + case "timestamp" -> timestampWriteFunction(YDB_TIMESTAMP_SQL_TYPE); + case "timestamp64" -> timestampWriteFunction(YDB_TIMESTAMP64_SQL_TYPE); + default -> throw new TrinoException(NOT_SUPPORTED, "Unsupported YDB timestamp type: " + typeName); + }; return ColumnMapping.longMapping( TIMESTAMP_MICROS, timestampReadFunction(TIMESTAMP_MICROS), writeFunction); } - private static LongWriteFunction timestamp64WriteFunction() { + private static LongWriteFunction dateWriteFunction(int sqlType) { + return LongWriteFunction.of( + sqlType, + (statement, index, value) -> statement.setObject(index, value, sqlType)); + } + + private static LongWriteFunction datetimeWriteFunction(int sqlType) { + return LongWriteFunction.of( + sqlType, + (statement, index, value) -> statement.setObject( + index, + fromTrinoTimestamp(value).toEpochSecond(UTC), + sqlType)); + } + + private static LongWriteFunction timestampWriteFunction(int sqlType) { return LongWriteFunction.of( - Types.TIMESTAMP, - (statement, index, value) -> statement.setObject(index, fromTrinoTimestamp(value).toInstant(UTC))); + sqlType, + (statement, index, value) -> statement.setObject( + index, + fromTrinoTimestamp(value).toInstant(UTC), + sqlType)); } @Override @@ -434,10 +468,10 @@ public WriteMapping toWriteMapping(ConnectorSession session, Type type) { return WriteMapping.sliceMapping("Bytes", varbinaryWriteFunction()); } if (type == DATE) { - return WriteMapping.longMapping("Date32", dateWriteFunctionUsingLocalDate()); + return WriteMapping.longMapping("Date32", dateWriteFunction(YDB_DATE32_SQL_TYPE)); } if (type == TIMESTAMP_MILLIS || type == TIMESTAMP_MICROS) { - return WriteMapping.longMapping("Timestamp64", timestamp64WriteFunction()); + return WriteMapping.longMapping("Timestamp64", timestampWriteFunction(YDB_TIMESTAMP64_SQL_TYPE)); } throw new TrinoException(NOT_SUPPORTED, "Unsupported column type: " + type); diff --git a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClientModule.java b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClientModule.java index 01cba8e7..53b2da8c 100644 --- a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClientModule.java +++ b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClientModule.java @@ -17,8 +17,6 @@ import io.trino.plugin.jdbc.logging.RemoteQueryModifier; import tech.ydb.jdbc.YdbDriver; -import java.util.Properties; - import static com.google.inject.multibindings.OptionalBinder.newOptionalBinder; public class YdbClientModule implements Module { @@ -51,13 +49,10 @@ public JdbcClient provideJdbcClient( public static ConnectionFactory createConnectionFactory( BaseJdbcConfig config, CredentialProvider credentialProvider) { - Properties connectionProperties = new Properties(); - connectionProperties.setProperty("forceSignedDatetimes", "true"); return DriverConnectionFactory.builder( new YdbDriver(), config.getConnectionUrl(), credentialProvider) - .setConnectionProperties(connectionProperties) .build(); } } From 45db5f8ed8a55c690723d8358ace2327886d69a0 Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Mon, 14 Sep 2026 17:49:19 +0300 Subject: [PATCH 3/4] List unsupported Trino time mappings explicitly --- .../src/test/java/tech/ydb/trino/TestYdbConnectorTest.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/ydb-trino-adapter/src/test/java/tech/ydb/trino/TestYdbConnectorTest.java b/ydb-trino-adapter/src/test/java/tech/ydb/trino/TestYdbConnectorTest.java index 9635310b..fdde9656 100644 --- a/ydb-trino-adapter/src/test/java/tech/ydb/trino/TestYdbConnectorTest.java +++ b/ydb-trino-adapter/src/test/java/tech/ydb/trino/TestYdbConnectorTest.java @@ -125,9 +125,10 @@ protected Optional filterDataMappingSmokeTestData(BaseConn String trinoTypeName = dataMappingTestSetup.getTrinoTypeName(); if (trinoTypeName.equals("char(3)")) { return Optional.of(dataMappingTestSetup.asUnsupported()); - } else if (trinoTypeName.startsWith("time") - && !trinoTypeName.equals("timestamp") - && !trinoTypeName.equals("timestamp(6)")) { + } else if (trinoTypeName.equals("time") + || trinoTypeName.equals("time(6)") + || trinoTypeName.equals("timestamp(3) with time zone") + || trinoTypeName.equals("timestamp(6) with time zone")) { // Нет time в YQL return Optional.empty(); } From f0639955bd9fb0ec403963a78717f42432c97601 Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Mon, 14 Sep 2026 18:05:28 +0300 Subject: [PATCH 4/4] Extract YDB temporal column mappings --- .../main/java/tech/ydb/trino/YdbClient.java | 72 ++--------------- .../tech/ydb/trino/YdbColumnMappings.java | 81 +++++++++++++++++++ 2 files changed, 87 insertions(+), 66 deletions(-) create mode 100644 ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbColumnMappings.java diff --git a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java index ab9126a5..546f003f 100644 --- a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java +++ b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java @@ -23,7 +23,6 @@ import io.trino.plugin.jdbc.JdbcSortItem; import io.trino.plugin.jdbc.JdbcTableHandle; import io.trino.plugin.jdbc.JdbcTypeHandle; -import io.trino.plugin.jdbc.LongWriteFunction; import io.trino.plugin.jdbc.PreparedQuery; import io.trino.plugin.jdbc.QueryBuilder; import io.trino.plugin.jdbc.RemoteTableName; @@ -84,11 +83,9 @@ import static io.trino.plugin.jdbc.StandardColumnMappings.bigintColumnMapping; import static io.trino.plugin.jdbc.StandardColumnMappings.bigintWriteFunction; import static io.trino.plugin.jdbc.StandardColumnMappings.booleanColumnMapping; -import static io.trino.plugin.jdbc.StandardColumnMappings.dateReadFunctionUsingLocalDate; import static io.trino.plugin.jdbc.StandardColumnMappings.decimalColumnMapping; import static io.trino.plugin.jdbc.StandardColumnMappings.doubleColumnMapping; import static io.trino.plugin.jdbc.StandardColumnMappings.doubleWriteFunction; -import static io.trino.plugin.jdbc.StandardColumnMappings.fromTrinoTimestamp; import static io.trino.plugin.jdbc.StandardColumnMappings.integerColumnMapping; import static io.trino.plugin.jdbc.StandardColumnMappings.integerWriteFunction; import static io.trino.plugin.jdbc.StandardColumnMappings.longDecimalWriteFunction; @@ -97,8 +94,6 @@ import static io.trino.plugin.jdbc.StandardColumnMappings.shortDecimalWriteFunction; import static io.trino.plugin.jdbc.StandardColumnMappings.smallintColumnMapping; import static io.trino.plugin.jdbc.StandardColumnMappings.smallintWriteFunction; -import static io.trino.plugin.jdbc.StandardColumnMappings.timestampColumnMapping; -import static io.trino.plugin.jdbc.StandardColumnMappings.timestampReadFunction; import static io.trino.plugin.jdbc.StandardColumnMappings.tinyintWriteFunction; import static io.trino.plugin.jdbc.StandardColumnMappings.varbinaryColumnMapping; import static io.trino.plugin.jdbc.StandardColumnMappings.varbinaryWriteFunction; @@ -122,21 +117,18 @@ import static io.trino.spi.type.VarcharType.createVarcharType; import static java.lang.Math.max; import static java.lang.String.format; -import static java.time.ZoneOffset.UTC; import static java.util.stream.Collectors.joining; -import static tech.ydb.jdbc.YdbConst.SQL_KIND_PRIMITIVE; +import static tech.ydb.trino.YdbColumnMappings.YDB_DATE32_SQL_TYPE; +import static tech.ydb.trino.YdbColumnMappings.YDB_TIMESTAMP64_SQL_TYPE; +import static tech.ydb.trino.YdbColumnMappings.dateColumnMapping; +import static tech.ydb.trino.YdbColumnMappings.dateWriteFunction; +import static tech.ydb.trino.YdbColumnMappings.timestampColumnMapping; +import static tech.ydb.trino.YdbColumnMappings.timestampWriteFunction; public class YdbClient extends BaseJdbcClient { static final String DEFAULT_SCHEMA = "default"; private static final int YDB_DEFAULT_DECIMAL_PRECISION = 22; private static final int YDB_DEFAULT_DECIMAL_SCALE = 9; - // https://github.com/ydb-platform/ydb-jdbc-driver/blob/v2.4.1/jdbc/src/main/java/tech/ydb/jdbc/common/YdbTypes.java#L64-L77 - private static final int YDB_DATE_SQL_TYPE = SQL_KIND_PRIMITIVE + 16; - private static final int YDB_DATETIME_SQL_TYPE = SQL_KIND_PRIMITIVE + 17; - private static final int YDB_TIMESTAMP_SQL_TYPE = SQL_KIND_PRIMITIVE + 18; - private static final int YDB_DATE32_SQL_TYPE = SQL_KIND_PRIMITIVE + 25; - private static final int YDB_DATETIME64_SQL_TYPE = SQL_KIND_PRIMITIVE + 26; - private static final int YDB_TIMESTAMP64_SQL_TYPE = SQL_KIND_PRIMITIVE + 27; private final ConnectorExpressionRewriter connectorExpressionRewriter; private final AggregateFunctionRewriter aggregateFunctionRewriter; @@ -380,58 +372,6 @@ private static ColumnMapping varcharColumnMapping(int varcharLength) { FULL_PUSHDOWN); } - private static ColumnMapping dateColumnMapping(JdbcTypeHandle typeHandle) { - String typeName = typeHandle.jdbcTypeName().orElse(""); - LongWriteFunction writeFunction = switch (typeName.toLowerCase(Locale.ROOT)) { - case "date" -> dateWriteFunction(YDB_DATE_SQL_TYPE); - case "date32" -> dateWriteFunction(YDB_DATE32_SQL_TYPE); - default -> throw new TrinoException(NOT_SUPPORTED, "Unsupported YDB date type: " + typeName); - }; - return ColumnMapping.longMapping( - DATE, - dateReadFunctionUsingLocalDate(), - writeFunction); - } - - private static ColumnMapping timestampColumnMapping(JdbcTypeHandle typeHandle) { - String typeName = typeHandle.jdbcTypeName().orElse(""); - LongWriteFunction writeFunction = switch (typeName.toLowerCase(Locale.ROOT)) { - case "datetime" -> datetimeWriteFunction(YDB_DATETIME_SQL_TYPE); - case "datetime64" -> datetimeWriteFunction(YDB_DATETIME64_SQL_TYPE); - case "timestamp" -> timestampWriteFunction(YDB_TIMESTAMP_SQL_TYPE); - case "timestamp64" -> timestampWriteFunction(YDB_TIMESTAMP64_SQL_TYPE); - default -> throw new TrinoException(NOT_SUPPORTED, "Unsupported YDB timestamp type: " + typeName); - }; - return ColumnMapping.longMapping( - TIMESTAMP_MICROS, - timestampReadFunction(TIMESTAMP_MICROS), - writeFunction); - } - - private static LongWriteFunction dateWriteFunction(int sqlType) { - return LongWriteFunction.of( - sqlType, - (statement, index, value) -> statement.setObject(index, value, sqlType)); - } - - private static LongWriteFunction datetimeWriteFunction(int sqlType) { - return LongWriteFunction.of( - sqlType, - (statement, index, value) -> statement.setObject( - index, - fromTrinoTimestamp(value).toEpochSecond(UTC), - sqlType)); - } - - private static LongWriteFunction timestampWriteFunction(int sqlType) { - return LongWriteFunction.of( - sqlType, - (statement, index, value) -> statement.setObject( - index, - fromTrinoTimestamp(value).toInstant(UTC), - sqlType)); - } - @Override public WriteMapping toWriteMapping(ConnectorSession session, Type type) { if (type == BOOLEAN) { diff --git a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbColumnMappings.java b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbColumnMappings.java new file mode 100644 index 00000000..afc4a6cb --- /dev/null +++ b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbColumnMappings.java @@ -0,0 +1,81 @@ +package tech.ydb.trino; + +import io.trino.plugin.jdbc.ColumnMapping; +import io.trino.plugin.jdbc.JdbcTypeHandle; +import io.trino.plugin.jdbc.LongWriteFunction; +import io.trino.spi.TrinoException; + +import java.util.Locale; + +import static io.trino.plugin.jdbc.StandardColumnMappings.dateReadFunctionUsingLocalDate; +import static io.trino.plugin.jdbc.StandardColumnMappings.fromTrinoTimestamp; +import static io.trino.plugin.jdbc.StandardColumnMappings.timestampReadFunction; +import static io.trino.spi.StandardErrorCode.NOT_SUPPORTED; +import static io.trino.spi.type.DateType.DATE; +import static io.trino.spi.type.TimestampType.TIMESTAMP_MICROS; +import static java.time.ZoneOffset.UTC; +import static tech.ydb.jdbc.YdbConst.SQL_KIND_PRIMITIVE; + +final class YdbColumnMappings { + // https://github.com/ydb-platform/ydb-jdbc-driver/blob/v2.4.1/jdbc/src/main/java/tech/ydb/jdbc/common/YdbTypes.java#L64-L77 + private static final int YDB_DATE_SQL_TYPE = SQL_KIND_PRIMITIVE + 16; + private static final int YDB_DATETIME_SQL_TYPE = SQL_KIND_PRIMITIVE + 17; + private static final int YDB_TIMESTAMP_SQL_TYPE = SQL_KIND_PRIMITIVE + 18; + static final int YDB_DATE32_SQL_TYPE = SQL_KIND_PRIMITIVE + 25; + private static final int YDB_DATETIME64_SQL_TYPE = SQL_KIND_PRIMITIVE + 26; + static final int YDB_TIMESTAMP64_SQL_TYPE = SQL_KIND_PRIMITIVE + 27; + + private YdbColumnMappings() {} + + static ColumnMapping dateColumnMapping(JdbcTypeHandle typeHandle) { + String typeName = typeHandle.jdbcTypeName().orElse(""); + LongWriteFunction writeFunction = switch (typeName.toLowerCase(Locale.ROOT)) { + case "date" -> dateWriteFunction(YDB_DATE_SQL_TYPE); + case "date32" -> dateWriteFunction(YDB_DATE32_SQL_TYPE); + default -> throw new TrinoException(NOT_SUPPORTED, "Unsupported YDB date type: " + typeName); + }; + return ColumnMapping.longMapping( + DATE, + dateReadFunctionUsingLocalDate(), + writeFunction); + } + + static ColumnMapping timestampColumnMapping(JdbcTypeHandle typeHandle) { + String typeName = typeHandle.jdbcTypeName().orElse(""); + LongWriteFunction writeFunction = switch (typeName.toLowerCase(Locale.ROOT)) { + case "datetime" -> datetimeWriteFunction(YDB_DATETIME_SQL_TYPE); + case "datetime64" -> datetimeWriteFunction(YDB_DATETIME64_SQL_TYPE); + case "timestamp" -> timestampWriteFunction(YDB_TIMESTAMP_SQL_TYPE); + case "timestamp64" -> timestampWriteFunction(YDB_TIMESTAMP64_SQL_TYPE); + default -> throw new TrinoException(NOT_SUPPORTED, "Unsupported YDB timestamp type: " + typeName); + }; + return ColumnMapping.longMapping( + TIMESTAMP_MICROS, + timestampReadFunction(TIMESTAMP_MICROS), + writeFunction); + } + + static LongWriteFunction dateWriteFunction(int sqlType) { + return LongWriteFunction.of( + sqlType, + (statement, index, value) -> statement.setObject(index, value, sqlType)); + } + + private static LongWriteFunction datetimeWriteFunction(int sqlType) { + return LongWriteFunction.of( + sqlType, + (statement, index, value) -> statement.setObject( + index, + fromTrinoTimestamp(value).toEpochSecond(UTC), + sqlType)); + } + + static LongWriteFunction timestampWriteFunction(int sqlType) { + return LongWriteFunction.of( + sqlType, + (statement, index, value) -> statement.setObject( + index, + fromTrinoTimestamp(value).toInstant(UTC), + sqlType)); + } +}