Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 9 additions & 4 deletions ydb-trino-adapter/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -52,8 +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(6)` создаются как YDB `Timestamp64`.
Это не исправляет диапазонные фильтры и `UPDATE` существующих столбцов YDB `Date`.
Не задавайте `forceSignedDatetimes=false` в JDBC URL: параметры URL имеют приоритет над настройками адаптера.
Новые столбцы Trino `timestamp(3)` и `timestamp(6)` создаются как YDB `Timestamp64`.
Существующие 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` остаются отдельным известным ограничением.
47 changes: 11 additions & 36 deletions ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -84,12 +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.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;
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;
Expand All @@ -98,9 +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.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;
Expand All @@ -116,15 +109,21 @@
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;
import static io.trino.spi.type.VarcharType.createUnboundedVarcharType;
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.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";
Expand Down Expand Up @@ -341,7 +340,7 @@ public Optional<ColumnMapping> 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();
};
Expand Down Expand Up @@ -373,30 +372,6 @@ private static ColumnMapping varcharColumnMapping(int varcharLength) {
FULL_PUSHDOWN);
}

private static ColumnMapping dateColumnMapping() {
return ColumnMapping.longMapping(
DATE,
dateReadFunctionUsingLocalDate(),
dateWriteFunctionUsingLocalDate());
}

private static ColumnMapping timestampColumnMapping(JdbcTypeHandle typeHandle) {
String typeName = typeHandle.jdbcTypeName().orElse("");
LongWriteFunction writeFunction = typeName.equalsIgnoreCase("Timestamp64")
? timestamp64WriteFunction()
: timestampWriteFunction(TIMESTAMP_MICROS);
return ColumnMapping.longMapping(
TIMESTAMP_MICROS,
timestampReadFunction(TIMESTAMP_MICROS),
writeFunction);
}

private static LongWriteFunction timestamp64WriteFunction() {
return LongWriteFunction.of(
Types.TIMESTAMP,
(statement, index, value) -> statement.setObject(index, fromTrinoTimestamp(value).toInstant(UTC)));
}

@Override
public WriteMapping toWriteMapping(ConnectorSession session, Type type) {
if (type == BOOLEAN) {
Expand Down Expand Up @@ -433,10 +408,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_MICROS) {
return WriteMapping.longMapping("Timestamp64", timestamp64WriteFunction());
if (type == TIMESTAMP_MILLIS || type == TIMESTAMP_MICROS) {
return WriteMapping.longMapping("Timestamp64", timestampWriteFunction(YDB_TIMESTAMP64_SQL_TYPE));
}

throw new TrinoException(NOT_SUPPORTED, "Unsupported column type: " + type);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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();
}
}
Original file line number Diff line number Diff line change
@@ -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));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -125,7 +125,10 @@ protected Optional<DataMappingTestSetup> 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.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();
}
Expand Down