diff --git a/.github/workflows/ci-trino-adapter.yaml b/.github/workflows/ci-trino-adapter.yaml index 1b251633..ae4bca04 100644 --- a/.github/workflows/ci-trino-adapter.yaml +++ b/.github/workflows/ci-trino-adapter.yaml @@ -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 @@ -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 diff --git a/ydb-trino-adapter/README.md b/ydb-trino-adapter/README.md index 40b27603..d91c8af8 100644 --- a/ydb-trino-adapter/README.md +++ b/ydb-trino-adapter/README.md @@ -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`. 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 466a277c..f7a68f7e 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 @@ -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; @@ -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; @@ -76,6 +78,10 @@ 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; @@ -83,8 +89,8 @@ 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; @@ -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 connectorExpressionRewriter; private final AggregateFunctionRewriter aggregateFunctionRewriter; private final ProjectFunctionRewriter projectFunctionRewriter; + private final boolean forceSignedDatetimes; @Inject public YdbClient( @@ -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()) @@ -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() { @@ -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)); 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 11f8d34f..4967376e 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 @@ -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; @@ -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 @@ -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 @@ -142,7 +153,7 @@ protected boolean isColumnNameRejected(Exception exception, String columnName, b protected Optional 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'", @@ -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 filterCaseSensitiveDataMappingTestData(DataMappingTestSetup dataMappingTestSetup) { if (dataMappingTestSetup.getTrinoTypeName().equals("char(1)")) {