Skip to content
7 changes: 6 additions & 1 deletion .github/workflows/ci-trino-adapter.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -21,8 +21,13 @@ jobs:
runs-on: ubuntu-latest

strategy:
fail-fast: false
matrix:
java: [ '25' ]
test_args:
- '-Duser.timezone=UTC'
- '-Duser.timezone=UTC -Dydb.test.force-signed-datetimes=true'
- "-Duser.timezone=Europe/Moscow -Dtest='TestYdbConnectorTest#testNativeDateCompatibility+testDatePredicateWidening'"

steps:
- uses: actions/checkout@v5
Expand All @@ -40,4 +45,4 @@ jobs:

- name: Build and test Trino Adapter
working-directory: ./ydb-trino-adapter
run: mvn $MAVEN_ARGS clean test
run: mvn $MAVEN_ARGS ${{matrix.test_args}} clean test
9 changes: 9 additions & 0 deletions ydb-trino-adapter/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,3 +49,12 @@ SELECT * FROM local.default.orders;

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

## Даты

По умолчанию Trino `date` создаётся как YDB `Date`. Параметр JDBC URL
`forceSignedDatetimes=true` переключает новые столбцы на `Date32`; существующие
столбцы обоих типов поддерживаются независимо от этого параметра. Для MERGE с
такими столбцами требуется стандартная подготовка запросов YDB JDBC — режим
`disablePrepareDataQuery=true` не поддерживается. Здесь проверена совместимость
только `Date` и `Date32`; предикаты ограничены диапазоном YDB `Date32`.
22 changes: 18 additions & 4 deletions ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Types;
import java.time.LocalDate;
import java.util.Collection;
import java.util.List;
import java.util.Locale;
Expand All @@ -68,6 +69,7 @@
import java.util.Optional;
import java.util.OptionalInt;
import java.util.OptionalLong;
import java.util.Properties;
import java.util.Set;
import java.util.TreeMap;
import java.util.function.BiFunction;
Expand All @@ -76,15 +78,19 @@
import java.util.stream.Collectors;
import java.util.stream.Stream;

import tech.ydb.jdbc.settings.YdbConfig;
import tech.ydb.jdbc.settings.YdbOperationProperties;
import tech.ydb.table.values.PrimitiveValue;

import static io.trino.plugin.jdbc.DefaultJdbcMetadata.MERGE_ROW_ID;
import static io.trino.plugin.jdbc.JdbcErrorCode.JDBC_ERROR;
import static io.trino.plugin.jdbc.PredicatePushdownController.DISABLE_PUSHDOWN;
import static io.trino.plugin.jdbc.PredicatePushdownController.FULL_PUSHDOWN;
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.dateWriteFunctionUsingLocalDate;
import static io.trino.plugin.jdbc.StandardColumnMappings.decimalColumnMapping;
import static io.trino.plugin.jdbc.StandardColumnMappings.doubleColumnMapping;
import static io.trino.plugin.jdbc.StandardColumnMappings.doubleWriteFunction;
Expand Down Expand Up @@ -127,10 +133,10 @@ 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;

private final ConnectorExpressionRewriter<ParameterizedExpression> connectorExpressionRewriter;
private final AggregateFunctionRewriter<JdbcExpression, ParameterizedExpression> aggregateFunctionRewriter;
private final ProjectFunctionRewriter<JdbcExpression, ParameterizedExpression> projectFunctionRewriter;
private final boolean forceSignedDatetimes;

@Inject
public YdbClient(
Expand All @@ -150,6 +156,14 @@ public YdbClient(
true
);

try {
forceSignedDatetimes = new YdbOperationProperties(YdbConfig.from(config.getConnectionUrl(), new Properties()))
.getForceNewDatetypes();
}
catch (SQLException e) {
throw new TrinoException(JDBC_ERROR, "Invalid YDB JDBC configuration", e);
}

this.connectorExpressionRewriter = JdbcConnectorExpressionRewriterBuilder.newBuilder()
.addStandardRules(this::quoted)
.add(new RewriteIn())
Expand Down Expand Up @@ -374,7 +388,7 @@ private static ColumnMapping dateColumnMapping() {
return ColumnMapping.longMapping(
DATE,
dateReadFunctionUsingLocalDate(),
dateWriteFunctionUsingLocalDate());
(statement, index, value) -> statement.setObject(index, PrimitiveValue.newDate32(LocalDate.ofEpochDay(value))));
}

private static ColumnMapping timestampColumnMapping() {
Expand Down Expand Up @@ -420,7 +434,7 @@ public WriteMapping toWriteMapping(ConnectorSession session, Type type) {
return WriteMapping.sliceMapping("Bytes", varbinaryWriteFunction());
}
if (type == DATE) {
return WriteMapping.longMapping("Date", dateWriteFunctionUsingLocalDate());
return WriteMapping.longMapping(forceSignedDatetimes ? "Date32" : "Date", dateWriteFunctionUsingLocalDate());
}
if (type == TIMESTAMP_MICROS) {
return WriteMapping.longMapping("Timestamp", timestampWriteFunction(TIMESTAMP_MICROS));
Expand Down
139 changes: 125 additions & 14 deletions ydb-trino-adapter/src/test/java/tech/ydb/trino/TestYdbConnectorTest.java
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
package tech.ydb.trino;

import com.google.common.collect.ImmutableMap;
import io.trino.spi.type.Type;
import io.trino.spi.type.VarcharType;
import io.trino.testing.BaseConnectorTest;
Expand All @@ -11,24 +12,41 @@
import org.assertj.core.api.Assertions;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
import tech.ydb.table.values.PrimitiveValue;
import tech.ydb.test.junit5.YdbHelperExtension;

import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.Statement;
import java.time.LocalDate;
import java.util.Optional;
import java.util.OptionalInt;

import static io.trino.testing.TestingNames.randomNameSuffix;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.junit.jupiter.api.Assertions.assertAll;

public class TestYdbConnectorTest extends BaseConnectorTest {
private static final String ALTERNATE_DATE_CATALOG = "alternate_dates";
private static final boolean FORCE_SIGNED_DATETIMES = Boolean.getBoolean("ydb.test.force-signed-datetimes");

@RegisterExtension
static final YdbHelperExtension ydb = new YdbHelperExtension();

@Override
protected QueryRunner createQueryRunner() throws Exception {
return YdbQueryRunner.builder(ydb)
String jdbcUrl = YdbQueryRunner.buildJdbcUrl(ydb);
QueryRunner queryRunner = YdbQueryRunner.builder(ydb)
.addConnectorProperty("connection-url", jdbcUrl + "&forceSignedDatetimes=" + FORCE_SIGNED_DATETIMES)
.setInitialTables(REQUIRED_TPCH_TABLES)
.build();
queryRunner.createCatalog(ALTERNATE_DATE_CATALOG, "ydb", ImmutableMap.of(
"connection-url", jdbcUrl + "&forceSignedDatetimes=" + !FORCE_SIGNED_DATETIMES,
"insert.non-transactional-insert.enabled", "true"));
return queryRunner;
}

@Test
Expand Down Expand Up @@ -79,26 +97,19 @@ protected boolean hasBehavior(TestingConnectorBehavior connectorBehavior) {
SUPPORTS_DROP_DEFAULT_COLUMN_VALUE,
SUPPORTS_ADD_COLUMN_NOT_NULL_CONSTRAINT -> false;
case SUPPORTS_TOPN_PUSHDOWN_WITH_VARCHAR -> true;
case SUPPORTS_NEGATIVE_DATE -> FORCE_SIGNED_DATETIMES;
default -> super.hasBehavior(connectorBehavior);
};
}

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

@Test
@Override
public void testDateYearOfEraPredicate() {
// YDB не поддерживает, negative daysSinceEpoch
protected String errorMessageForInsertNegativeDate(String date) {
return ".*negative daysSinceEpoch.*";
}

@Test
@Override
public void testCreateTableAsSelectNegativeDate() {
// YDB не поддерживает, negative daysSinceEpoch
protected String errorMessageForCreateTableAsSelectNegativeDate(String date) {
return ".*negative daysSinceEpoch.*";
}

@Test
Expand Down Expand Up @@ -142,7 +153,7 @@ 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")) {
} else if (dataMappingTestSetup.getTrinoTypeName().equals("date") && !FORCE_SIGNED_DATETIMES) {
return Optional.of(new DataMappingTestSetup(
dataMappingTestSetup.getTrinoTypeName(),
"DATE '2006-06-06'",
Expand All @@ -163,6 +174,106 @@ public void testVarbinaryCreateTableAndInsert() {
}
}

@Test
public void testDatePredicateWidening() throws Exception {
assertAll(
() -> verifyDatePredicateWidening("local"),
() -> verifyDatePredicateWidening(ALTERNATE_DATE_CATALOG));
}

private void verifyDatePredicateWidening(String catalog) throws Exception {
try (TestTable table = new TestTable(
new JdbcSqlExecutor(YdbQueryRunner.buildJdbcUrl(ydb)),
"date_predicate_",
"(legacy_key Date NOT NULL, signed_key Date32 NOT NULL, PRIMARY KEY (legacy_key, signed_key))")) {
try (Connection connection = DriverManager.getConnection(YdbQueryRunner.buildJdbcUrl(ydb));
PreparedStatement statement = connection.prepareStatement(
"INSERT INTO `" + table.getName() + "` (legacy_key, signed_key) VALUES (?, ?)")) {
statement.setObject(1, PrimitiveValue.newDate(LocalDate.of(2020, 1, 1)));
statement.setObject(2, PrimitiveValue.newDate32(LocalDate.of(-1, 1, 1)));
statement.executeUpdate();
}
try (Connection connection = DriverManager.getConnection(YdbQueryRunner.buildJdbcUrl(ydb));
Statement statement = connection.createStatement();
ResultSet rows = statement.executeQuery("SELECT * FROM `" + table.getName() + "`")) {
assertThat(rows.next()).isTrue();
assertThat(rows.getObject("legacy_key", LocalDate.class)).isEqualTo(LocalDate.of(2020, 1, 1));
assertThat(rows.getObject("signed_key", LocalDate.class)).isEqualTo(LocalDate.of(-1, 1, 1));
assertThat(rows.next()).isFalse();
}
String name = catalog + ".default." + table.getName();
assertQuery("SELECT count(*) FROM " + name, "VALUES CAST(1 AS BIGINT)");
assertAll(
() -> assertQueryReturnsEmptyResult("SELECT * FROM " + name + " WHERE legacy_key = DATE '-0001-01-01'"),
() -> assertQuery("SELECT count(*) FROM " + name + " WHERE legacy_key > DATE '-0001-01-01'", "VALUES CAST(1 AS BIGINT)"),
() -> assertQueryReturnsEmptyResult("SELECT * FROM " + name + " WHERE legacy_key = DATE '2106-01-01'"),
() -> assertQuery("SELECT count(*) FROM " + name + " WHERE legacy_key < DATE '2106-01-01'", "VALUES CAST(1 AS BIGINT)"));
}
}

@Test
public void testNativeDateCompatibility() throws Exception {
assertAll(
() -> verifyNativeDateCompatibility("local"),
() -> verifyNativeDateCompatibility(ALTERNATE_DATE_CATALOG));
}

private void verifyNativeDateCompatibility(String catalog) throws Exception {
String ddlTable = "trino_date_" + randomNameSuffix();
boolean signed = catalog.equals("local") ? FORCE_SIGNED_DATETIMES : !FORCE_SIGNED_DATETIMES;
try (TestTable table = new TestTable(
new JdbcSqlExecutor(YdbQueryRunner.buildJdbcUrl(ydb)),
"native_dates_",
"(legacy_key Date NOT NULL, signed_key Date32 NOT NULL, legacy_value Date, signed_value Date32, " +
"PRIMARY KEY (legacy_key, signed_key))")) {
String name = catalog + ".default." + table.getName();
assertUpdate("INSERT INTO " + name + " VALUES " +
"(DATE '2020-01-01', DATE '-0001-01-01', DATE '2000-01-01', NULL), " +
"(DATE '2020-01-02', DATE '-0001-01-02', NULL, DATE '-0002-01-01'), " +
"(DATE '2020-01-01', DATE '-0001-01-04', DATE '2004-01-01', NULL)", 3);
assertQuery("SELECT legacy_value, signed_value FROM " + name +
" WHERE legacy_key = DATE '2020-01-01' AND signed_key = DATE '-0001-01-01'",
"VALUES (DATE '2000-01-01', CAST(NULL AS DATE))");
assertUpdate("UPDATE " + name + " SET legacy_value = DATE '2001-01-01', signed_value = DATE '-0003-01-01'" +
" WHERE legacy_key = DATE '2020-01-01' AND signed_key = DATE '-0001-01-01'", 1);
assertUpdate("""
MERGE INTO %s t
USING (VALUES
(DATE '2020-01-01', DATE '-0001-01-01', CAST(NULL AS DATE), DATE '-0004-01-01', 'update'),
(DATE '2020-01-02', DATE '-0001-01-02', CAST(NULL AS DATE), CAST(NULL AS DATE), 'delete'),
(DATE '2020-01-03', DATE '-0001-01-03', DATE '2003-01-01', CAST(NULL AS DATE), 'insert')
) s (legacy_key, signed_key, legacy_value, signed_value, operation)
ON (t.legacy_key = s.legacy_key AND t.signed_key = s.signed_key)
WHEN MATCHED AND s.operation = 'delete' THEN DELETE
WHEN MATCHED THEN UPDATE SET legacy_value = s.legacy_value, signed_value = s.signed_value
WHEN NOT MATCHED THEN INSERT VALUES (s.legacy_key, s.signed_key, s.legacy_value, s.signed_value)
""".formatted(name), 3);
assertQuery("SELECT legacy_key, signed_key, legacy_value, signed_value FROM " + name + " ORDER BY legacy_key",
"VALUES " +
"(DATE '2020-01-01', DATE '-0001-01-01', NULL, DATE '-0004-01-01'), " +
"(DATE '2020-01-01', DATE '-0001-01-04', DATE '2004-01-01', NULL), " +
"(DATE '2020-01-03', DATE '-0001-01-03', DATE '2003-01-01', NULL)");
try (Connection connection = DriverManager.getConnection(YdbQueryRunner.buildJdbcUrl(ydb));
Statement statement = connection.createStatement();
ResultSet rows = statement.executeQuery("SELECT * FROM `" + table.getName() + "` ORDER BY legacy_key, signed_key")) {
assertThat(rows.next()).isTrue();
assertThat(rows.getObject("legacy_value", LocalDate.class)).isNull();
assertThat(rows.getObject("signed_value", LocalDate.class)).isEqualTo(LocalDate.of(-4, 1, 1));
}
}
try {
assertUpdate("CREATE TABLE " + catalog + ".default." + ddlTable + " (value DATE)");
try (Connection connection = DriverManager.getConnection(YdbQueryRunner.buildJdbcUrl(ydb));
ResultSet columns = connection.getMetaData().getColumns(null, null, ddlTable, "value")) {
assertThat(columns.next()).isTrue();
assertThat(columns.getString("TYPE_NAME")).isEqualTo(signed ? "Date32" : "Date");
}
}
finally {
assertUpdate("DROP TABLE IF EXISTS " + catalog + ".default." + ddlTable);
}
}

@Override
protected Optional<DataMappingTestSetup> filterCaseSensitiveDataMappingTestData(DataMappingTestSetup dataMappingTestSetup) {
if (dataMappingTestSetup.getTrinoTypeName().equals("char(1)")) {
Expand Down
Loading