Skip to content
Closed
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
4 changes: 4 additions & 0 deletions ydb-trino-adapter/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,3 +49,7 @@ SELECT * FROM local.default.orders;

YDB `Text` отображается в Trino как `varchar`, а `Bytes` — как `varbinary` без
декодирования UTF-8. При создании таблиц адаптер использует типы `Text` и `Bytes`.

Trino `date` и `timestamp` создают YDB `Date32` и `Timestamp64`. Адаптер по
умолчанию включает параметр JDBC `forceSignedDatetimes`; значение в
`connection-url` имеет приоритет. Существующие таблицы не мигрируются.
3 changes: 2 additions & 1 deletion ydb-trino-adapter/ROADMAP.md
Original file line number Diff line number Diff line change
Expand Up @@ -136,7 +136,8 @@ unsupported behavior, record:
2. an authoritative documentation link or tracked upstream issue;
3. a focused negative test that proves the connector fails clearly.

The existing negative-date overrides need this treatment.
Trino `date` and `timestamp` use YDB `Date32` and `Timestamp64`; the connector
enables signed datetime support by default. Existing tables are not migrated.
YDB `Date` starts at the Unix epoch; see
[primitive types](https://ydb.tech/docs/en/yql/reference/types/primitive).
CHAR is rejected with the focused inherited contract
Expand Down
18 changes: 13 additions & 5 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,6 +23,7 @@
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 @@ -50,6 +51,7 @@
import io.trino.spi.expression.ConnectorExpression;
import io.trino.spi.type.DecimalType;
import io.trino.spi.type.Type;
import io.trino.spi.type.TimestampType;
import io.trino.spi.type.VarcharType;

import jakarta.annotation.Nullable;
Expand Down Expand Up @@ -88,6 +90,7 @@
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,7 +101,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;
Expand All @@ -121,6 +123,7 @@
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;

public class YdbClient extends BaseJdbcClient {
Expand Down Expand Up @@ -381,7 +384,12 @@ private static ColumnMapping timestampColumnMapping() {
return ColumnMapping.longMapping(
TIMESTAMP_MICROS,
timestampReadFunction(TIMESTAMP_MICROS),
timestampWriteFunction(TIMESTAMP_MICROS));
timestamp64WriteFunction());
}

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

@Override
Expand Down Expand Up @@ -420,10 +428,10 @@ public WriteMapping toWriteMapping(ConnectorSession session, Type type) {
return WriteMapping.sliceMapping("Bytes", varbinaryWriteFunction());
}
if (type == DATE) {
return WriteMapping.longMapping("Date", dateWriteFunctionUsingLocalDate());
return WriteMapping.longMapping("Date32", dateWriteFunctionUsingLocalDate());
}
if (type == TIMESTAMP_MICROS) {
return WriteMapping.longMapping("Timestamp", timestampWriteFunction(TIMESTAMP_MICROS));
if (type instanceof TimestampType timestampType && timestampType.isShort()) {
return WriteMapping.longMapping("Timestamp64", timestamp64WriteFunction());
}

throw new TrinoException(NOT_SUPPORTED, "Unsupported column type: " + type);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,7 @@ public static ConnectionFactory createConnectionFactory(
Properties connectionProperties = new Properties();
// Avoid the YDB JDBC 2.3.18 shared-context close/register race under concurrent connections.
connectionProperties.setProperty("cacheConnectionsInDriver", "false");
connectionProperties.setProperty("forceSignedDatetimes", "true");
return DriverConnectionFactory.builder(
new YdbDriver(),
config.getConnectionUrl(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,9 +30,9 @@ public static Optional<JdbcTypeHandle> toTypeHandle(Type type) {
case io.trino.spi.type.DoubleType _ ->
Optional.of(new JdbcTypeHandle(Types.DOUBLE, Optional.of("Double"), Optional.empty(), Optional.empty(), Optional.empty(), Optional.empty()));
case io.trino.spi.type.DateType _ ->
Optional.of(new JdbcTypeHandle(Types.DATE, Optional.of("Date"), Optional.empty(), Optional.empty(), Optional.empty(), Optional.empty()));
Optional.of(new JdbcTypeHandle(Types.DATE, Optional.of("Date32"), Optional.empty(), Optional.empty(), Optional.empty(), Optional.empty()));
case io.trino.spi.type.TimestampType _ ->
Optional.of(new JdbcTypeHandle(Types.TIMESTAMP, Optional.of("Timestamp"), Optional.empty(), Optional.empty(), Optional.empty(), Optional.empty()));
Optional.of(new JdbcTypeHandle(Types.TIMESTAMP, Optional.of("Timestamp64"), Optional.empty(), Optional.empty(), Optional.empty(), Optional.empty()));
case DecimalType decimalType ->
Optional.of(new JdbcTypeHandle(Types.DECIMAL, Optional.of("Decimal"), Optional.of(decimalType.getPrecision()), Optional.of(decimalType.getScale()), Optional.empty(), Optional.empty()));
case io.trino.spi.type.VarcharType _ ->
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -83,24 +83,6 @@ protected boolean hasBehavior(TestingConnectorBehavior connectorBehavior) {
};
}

@Test
@Override
public void testInsertNegativeDate() {
// YDB не поддерживает, negative daysSinceEpoch
}

@Test
@Override
public void testDateYearOfEraPredicate() {
// YDB не поддерживает, negative daysSinceEpoch
}

@Test
@Override
public void testCreateTableAsSelectNegativeDate() {
// YDB не поддерживает, negative daysSinceEpoch
}

@Test
@Override
public void testCharVarcharComparison() {
Expand Down Expand Up @@ -142,15 +124,12 @@ protected boolean isColumnNameRejected(Exception exception, String columnName, b
protected Optional<DataMappingTestSetup> filterDataMappingSmokeTestData(BaseConnectorTest.DataMappingTestSetup dataMappingTestSetup) {
if (dataMappingTestSetup.getTrinoTypeName().equals("char(3)")) {
return Optional.of(dataMappingTestSetup.asUnsupported());
} else if (dataMappingTestSetup.getTrinoTypeName().equals("date")) {
return Optional.of(new DataMappingTestSetup(
dataMappingTestSetup.getTrinoTypeName(),
"DATE '2006-06-06'",
"DATE '2026-06-06'"
));
} else if (dataMappingTestSetup.getTrinoTypeName().startsWith("time")) {
} else if (dataMappingTestSetup.getTrinoTypeName().equals("time")
|| dataMappingTestSetup.getTrinoTypeName().equals("time(6)")) {
// Нет time в YQL
return Optional.empty();
} else if (dataMappingTestSetup.getTrinoTypeName().contains("with time zone")) {
return Optional.empty();
}
return Optional.of(dataMappingTestSetup);
}
Expand Down