diff --git a/ydb-trino-adapter/ROADMAP.md b/ydb-trino-adapter/ROADMAP.md index 89fe9311..f8b8dbe8 100644 --- a/ydb-trino-adapter/ROADMAP.md +++ b/ydb-trino-adapter/ROADMAP.md @@ -74,37 +74,31 @@ within the CI budget. An isolated local Colima run previously exceeded 11 minutes, so MERGE scalability remains a production concern rather than a CI failure. -## P0 — harden MERGE scalability and atomicity - -1. Replace row-by-row prepared-statement execution in `YdbMergeSink` with a - set-based YQL path. Candidate primitives are - [`UPDATE ... ON`](https://ydb.tech/docs/en/yql/reference/syntax/update) and - [`AS_TABLE`](https://ydb.tech/docs/en/yql/reference/syntax/select/from_as_table), - using bounded batches of typed list-of-struct parameters. -2. Preserve one logical MERGE transaction. Until a staging/finalize design is - implemented, keep one writer task and reject Trino query/task retries for a - direct-to-target merge. +## P0 — harden standard MERGE execution + +1. Keep MERGE on Trino's `JdbcMergeSink`; YDB-specific code should only build + the handle and supply primary-key metadata. +2. The standard sink may commit INSERT, DELETE, and UPDATE batches through + separate JDBC connections. Keep one writer task and reject Trino query/task + retries, but do not claim that this makes MERGE atomic across operation sinks. 3. YQL `UPDATE` cannot change a primary-key value. Implement physical-key changes as atomic delete+insert row changes, or reject that statement with a documented `NOT_SUPPORTED` error. See the [YQL UPDATE contract](https://ydb.tech/docs/en/yql/reference/syntax/update). 4. Add focused tests for composite primary keys, a non-unique first visible - column, physical-key updates, rollback/close, and fresh-state retries. + column, physical-key updates, partial sink failures, and abort cleanup. **Exit criterion:** retain the current green inherited `testMerge*` suite and -add a bounded-memory benchmark that demonstrates acceptable production-scale -runtime for the set-based implementation. +cover partial failures between operation sinks. If atomic MERGE becomes a +requirement, implement staging plus one finalize transaction instead of +wrapping the standard sink in replay logic. ## P1 — retry and transaction hardening -- Retry only statuses classified as unconditional by the pinned YDB SDK. - `TIMEOUT`, `UNDETERMINED`, transport failures, and other conditional statuses - are unsafe for non-idempotent writes unless an operation-id/staging design - proves replay safety. See +- Do not wrap `JdbcMergeSink` in connector-owned replay after any driver or YDB + failure. A future retry design requires staging or an operation ID that proves + replay safety. See [YDB SDK error handling](https://ydb.tech/docs/en/reference/ydb-sdk/error_handling). -- Keep the JDBC `SessionPool.acquire` scheduler-rejection workaround limited to - that provably pre-execution stack. Track it against the YDB JDBC driver and - remove the connector workaround after upgrading to a fixed driver. - YDB JDBC 2.3.18 connection-context caching has a close/register race under concurrent connections. The connector therefore defaults `cacheConnectionsInDriver` to `false`; an explicit JDBC URL option can @@ -113,10 +107,7 @@ runtime for the set-based implementation. with a verified cache-lifecycle fix. - Do not replay buffered INSERT pages after `JdbcPageSink` may already have committed an internal batch. -- Add unit tests for status classification, interrupted backoff, rollback - failure suppression, connection cleanup, and a failure after commit. -- Define memory/backpressure limits for buffered merge pages; memory usage must - not remain unreported. +- Add tests for partial sink failures and abort cleanup. Concurrent `ALTER TABLE ... ADD COLUMN` statements on one table can be rejected by YDB with `OVERLOADED` (400060) and the specific issue `path is under 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 ba7105e2..466a277c 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 @@ -672,7 +672,7 @@ public void finishMerge( JdbcMergeTableHandle tableHandle, Set pageSinkIds ) { - // Each YdbMergeSink owns and commits its transaction before reporting success. + // JdbcMergeSink finishes its operation-specific sinks before reporting success. } @Override 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 2591c0a6..f1e07ecd 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 @@ -30,7 +30,6 @@ public void configure(Binder binder) { .to(YdbMetadataFactory.class) .in(Scopes.SINGLETON); - binder.bind(YdbPageSinkProvider.class).in(Scopes.SINGLETON); binder.bind(YdbConnector.class).in(Scopes.SINGLETON); } diff --git a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbConnector.java b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbConnector.java index c0452da2..851deed6 100644 --- a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbConnector.java +++ b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbConnector.java @@ -23,7 +23,7 @@ public YdbConnector( LifeCycleManager lifeCycleManager, ConnectorSplitManager jdbcSplitManager, ConnectorPageSourceProvider jdbcPageSourceProvider, - YdbPageSinkProvider jdbcPageSinkProvider, + ConnectorPageSinkProvider jdbcPageSinkProvider, Optional accessControl, Set procedures, Set connectorTableFunctions, diff --git a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbMergeSink.java b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbMergeSink.java deleted file mode 100644 index d1f29ddd..00000000 --- a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbMergeSink.java +++ /dev/null @@ -1,476 +0,0 @@ -package tech.ydb.trino; - -import com.google.common.collect.ImmutableList; -import io.airlift.slice.Slice; -import io.airlift.slice.Slices; -import io.trino.plugin.jdbc.BooleanWriteFunction; -import io.trino.plugin.jdbc.DoubleWriteFunction; -import io.trino.plugin.jdbc.JdbcClient; -import io.trino.plugin.jdbc.JdbcColumnHandle; -import io.trino.plugin.jdbc.JdbcMergeTableHandle; -import io.trino.plugin.jdbc.JdbcOutputTableHandle; -import io.trino.plugin.jdbc.LongWriteFunction; -import io.trino.plugin.jdbc.ObjectWriteFunction; -import io.trino.plugin.jdbc.QueryBuilder; -import io.trino.plugin.jdbc.SliceWriteFunction; -import io.trino.plugin.jdbc.WriteFunction; -import io.trino.plugin.jdbc.logging.RemoteQueryModifier; -import io.trino.spi.Page; -import io.trino.spi.TrinoException; -import io.trino.spi.block.Block; -import io.trino.spi.block.RowBlock; -import io.trino.spi.connector.ColumnHandle; -import io.trino.spi.connector.ConnectorMergeSink; -import io.trino.spi.connector.ConnectorMergeTableHandle; -import io.trino.spi.connector.ConnectorPageSinkId; -import io.trino.spi.connector.ConnectorSession; -import io.trino.spi.connector.ConnectorTransactionHandle; -import io.trino.spi.type.Type; - -import java.sql.Connection; -import java.sql.PreparedStatement; -import java.sql.SQLException; -import java.util.ArrayList; -import java.util.Collection; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.Set; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.RejectedExecutionException; -import java.util.stream.IntStream; - -import static com.google.common.collect.ImmutableSet.toImmutableSet; -import static io.trino.plugin.jdbc.JdbcErrorCode.JDBC_ERROR; -import static io.trino.spi.type.IntegerType.INTEGER; -import static io.trino.spi.type.TinyintType.TINYINT; -import static java.util.concurrent.CompletableFuture.completedFuture; - -/* - Code partly borrowed from JdbcMergeSink. Key differences are: - - We do not write into the target table on each storeMergedRows call; instead, we store the necessary pages - and write everything at once when finish() is called. Without it, I had serious atomicity problems and flaky tests. - - We retry in case YDB returns a retryable exception on an update attempt. - - We do not use temporary tables unlike base Trino implementation, but write straight to YDB target. - */ -public class YdbMergeSink implements ConnectorMergeSink { - private final ConnectorSession session; - private final JdbcMergeTableHandle mergeHandle; - private final JdbcClient jdbcClient; - private final ConnectorPageSinkId pageSinkId; - private final RemoteQueryModifier remoteQueryModifier; - - private final List bufferedPages = new ArrayList<>(); - private boolean finished = false; - - public YdbMergeSink( - @SuppressWarnings("unused") ConnectorTransactionHandle transactionHandle, - ConnectorSession session, - ConnectorMergeTableHandle mergeTableHandle, - JdbcClient jdbcClient, - ConnectorPageSinkId pageSinkId, - RemoteQueryModifier remoteQueryModifier, - @SuppressWarnings("unused") QueryBuilder queryBuilder - ) { - this.session = session; - this.mergeHandle = (JdbcMergeTableHandle) mergeTableHandle; - this.jdbcClient = jdbcClient; - this.pageSinkId = pageSinkId; - this.remoteQueryModifier = remoteQueryModifier; - } - - @Override - public void storeMergedRows(Page page) { - if (finished) { - throw new IllegalStateException(); - } - bufferedPages.add(page); - } - - @Override - public CompletableFuture> finish() { - finished = true; - - int maxAttempts = 10; - Exception lastException = null; - - for (int attempt = 0; attempt < maxAttempts; attempt++) { - Connection connection = null; - boolean committed = false; - try { - connection = openConnection(); - executeMergeInTransaction(connection); - connection.commit(); - committed = true; - connection.close(); - connection = null; - - Slice value = Slices.allocate(Long.BYTES); - value.setLong(0, pageSinkId.getId()); - return completedFuture(ImmutableList.of(value)); - } - catch (Exception e) { - if (connection != null) { - if (!committed) { - try { - connection.rollback(); - } - catch (SQLException rollbackError) { - e.addSuppressed(rollbackError); - } - } - try { - connection.close(); - } - catch (SQLException closeError) { - e.addSuppressed(closeError); - } - } - - if (committed) { - throw new TrinoException(JDBC_ERROR, "YDB MERGE committed, but closing its connection failed", e); - } - - if (!YdbRetryUtils.isRetryable(e) && !isRejectedDriverSessionAcquisition(e)) { - throw new TrinoException(JDBC_ERROR, e); - } - - lastException = e; - - long delay = YdbRetryUtils.calculateBackoff(attempt); - - try { - Thread.sleep(delay); - } - catch (InterruptedException ie) { - Thread.currentThread().interrupt(); - throw new TrinoException(JDBC_ERROR, ie); - } - } - } - - throw new TrinoException(JDBC_ERROR, lastException); - } - - private Connection openConnection() throws SQLException { - JdbcOutputTableHandle outputHandle = mergeHandle.getOutputTableHandle(); - Connection connection = jdbcClient.getConnection(session, outputHandle); - try { - connection.setAutoCommit(false); - return connection; - } - catch (SQLException e) { - try { - connection.close(); - } - catch (SQLException closeError) { - e.addSuppressed(closeError); - } - throw e; - } - } - - private void executeMergeInTransaction(Connection connection) throws SQLException { - JdbcOutputTableHandle outputHandle = mergeHandle.getOutputTableHandle(); - List primaryKeys = mergeHandle.getPrimaryKeys(); - - int columnCount = outputHandle.getColumnNames().size(); - List columns = mergeHandle.getDataColumns(); - - List insertPages = new ArrayList<>(); - List deletePages = new ArrayList<>(); - Map> updatePagesByCase = new HashMap<>(); - - for (Page page : bufferedPages) { - separateMergeOperations(page, columnCount, insertPages, deletePages, updatePagesByCase, columns); - } - - if (!deletePages.isEmpty()) { - executeDeleteOperations(connection, deletePages, primaryKeys); - } - - if (!updatePagesByCase.isEmpty()) { - executeUpdateOperations(connection, updatePagesByCase, primaryKeys, columns); - } - - if (!insertPages.isEmpty()) { - executeInsertOperations(connection, insertPages, outputHandle); - } - } - - private void executeDeleteOperations( - Connection connection, - List deletePages, - List primaryKeys - ) throws SQLException { - String tableName = mergeHandle.getTableHandle().getRequiredNamedRelation().getRemoteTableName().getTableName(); - StringBuilder deleteSql = new StringBuilder("DELETE FROM "); - deleteSql.append(jdbcClient.quoted(tableName)); - deleteSql.append(" WHERE "); - - for (int i = 0; i < primaryKeys.size(); i++) { - if (i > 0) { - deleteSql.append(" AND "); - } - JdbcColumnHandle pk = primaryKeys.get(i); - deleteSql.append(jdbcClient.quoted(pk.getColumnName())); - deleteSql.append(" = ?"); - } - - String sql = remoteQueryModifier.apply(session, deleteSql.toString()); - - try (PreparedStatement stmt = connection.prepareStatement(sql)) { - for (Page page : deletePages) { - for (int pos = 0; pos < page.getPositionCount(); pos++) { - for (int i = 0; i < primaryKeys.size(); i++) { - Block block = page.getBlock(i); - Type type = primaryKeys.get(i).getColumnType(); - WriteFunction writer = jdbcClient.toWriteMapping(session, type).getWriteFunction(); - setParameter(stmt, i + 1, block, pos, type, writer); - } - stmt.addBatch(); - } - } - stmt.executeBatch(); - } - } - - private void executeUpdateOperations( - Connection connection, - Map> updatePagesByCase, - List primaryKeys, - List columns - ) throws SQLException { - String tableName = mergeHandle.getTableHandle().getRequiredNamedRelation().getRemoteTableName().getTableName(); - - for (Map.Entry> entry : updatePagesByCase.entrySet()) { - int caseNumber = entry.getKey(); - List pages = entry.getValue(); - - Collection updateColumns = mergeHandle.getUpdateCaseColumns().get(caseNumber); - if (updateColumns == null || updateColumns.isEmpty()) { - continue; - } - - StringBuilder updateSql = new StringBuilder("UPDATE "); - updateSql.append(jdbcClient.quoted(tableName)); - updateSql.append(" SET "); - - Set updateChannelsSet = updateColumns.stream() - .map(JdbcColumnHandle.class::cast) - .map(columns::indexOf) - .collect(toImmutableSet()); - - List updateColumnList = new ArrayList<>(); - for (int channel = 0; channel < columns.size(); channel++) { - if (updateChannelsSet.contains(channel)) { - updateColumnList.add(columns.get(channel)); - } - } - - for (int i = 0; i < updateColumnList.size(); i++) { - if (i > 0) { - updateSql.append(", "); - } - JdbcColumnHandle col = updateColumnList.get(i); - updateSql.append(jdbcClient.quoted(col.getColumnName())); - updateSql.append(" = ?"); - } - - updateSql.append(" WHERE "); - for (int i = 0; i < primaryKeys.size(); i++) { - if (i > 0) { - updateSql.append(" AND "); - } - JdbcColumnHandle pk = primaryKeys.get(i); - updateSql.append(jdbcClient.quoted(pk.getColumnName())); - updateSql.append(" = ?"); - } - - String sql = remoteQueryModifier.apply(session, updateSql.toString()); - - try (PreparedStatement stmt = connection.prepareStatement(sql)) { - for (Page page : pages) { - for (int pos = 0; pos < page.getPositionCount(); pos++) { - for (int i = 0; i < updateColumnList.size(); i++) { - Block block = page.getBlock(i); - Type type = updateColumnList.get(i).getColumnType(); - WriteFunction writer = jdbcClient.toWriteMapping(session, type).getWriteFunction(); - setParameter(stmt, i + 1, block, pos, type, writer); - } - - int offset = updateColumnList.size(); - for (int i = 0; i < primaryKeys.size(); i++) { - Block block = page.getBlock(offset + i); - Type type = primaryKeys.get(i).getColumnType(); - WriteFunction writer = jdbcClient.toWriteMapping(session, type).getWriteFunction(); - setParameter(stmt, offset + i + 1, block, pos, type, writer); - } - stmt.addBatch(); - } - } - stmt.executeBatch(); - } - } - } - - private void executeInsertOperations( - Connection connection, - List insertPages, - JdbcOutputTableHandle outputHandle - ) throws SQLException { - List columnTypes = outputHandle.getColumnTypes(); - List columnWriters = getColumnWriters(columnTypes); - String insertSql = jdbcClient.buildInsertSql(outputHandle, columnWriters); - insertSql = remoteQueryModifier.apply(session, insertSql); - - try (PreparedStatement insertStmt = connection.prepareStatement(insertSql)) { - int columnCount = outputHandle.getColumnNames().size(); - - for (Page page : insertPages) { - for (int pos = 0; pos < page.getPositionCount(); pos++) { - for (int channel = 0; channel < columnCount; channel++) { - setParameter(insertStmt, channel + 1, page.getBlock(channel), pos, columnTypes.get(channel), columnWriters.get(channel)); - } - insertStmt.addBatch(); - } - } - insertStmt.executeBatch(); - } - } - - private void separateMergeOperations( - Page page, - int columnCount, - List insertPages, - List deletePages, - Map> updatePagesByCase, - List columns - ) { - Block operationBlock = page.getBlock(columnCount); - Block updateCaseBlock = page.getBlock(columnCount + 1); - - List insertPositions = new ArrayList<>(); - List deletePositions = new ArrayList<>(); - Map> updatePositionsByCase = new HashMap<>(); - - for (int pos = 0; pos < page.getPositionCount(); pos++) { - int operation = TINYINT.getByte(operationBlock, pos); - switch (operation) { - case INSERT_OPERATION_NUMBER -> insertPositions.add(pos); - case DELETE_OPERATION_NUMBER -> deletePositions.add(pos); - case UPDATE_OPERATION_NUMBER -> { - int caseNumber = INTEGER.getInt(updateCaseBlock, pos); - updatePositionsByCase.computeIfAbsent(caseNumber, _ -> new ArrayList<>()).add(pos); - } - default -> throw new IllegalStateException(); - } - } - - if (!insertPositions.isEmpty()) { - int[] positions = insertPositions.stream().mapToInt(Integer::intValue).toArray(); - Page insertData = page.getColumns(IntStream.range(0, columnCount).toArray()) - .getPositions(positions, 0, positions.length); - insertPages.add(insertData); - } - - if (!deletePositions.isEmpty()) { - int[] positions = deletePositions.stream().mapToInt(Integer::intValue).toArray(); - Block rowIdBlock = page.getBlock(columnCount + 2); - List rowIdFields = RowBlock.getRowFieldsFromBlock(rowIdBlock); - Block[] deleteBlocks = new Block[rowIdFields.size()]; - for (int i = 0; i < rowIdFields.size(); i++) { - deleteBlocks[i] = rowIdFields.get(i).getPositions(positions, 0, positions.length); - } - deletePages.add(new Page(positions.length, deleteBlocks)); - } - - for (Map.Entry> entry : updatePositionsByCase.entrySet()) { - int caseNumber = entry.getKey(); - int[] positions = entry.getValue().stream().mapToInt(Integer::intValue).toArray(); - - Collection updateColumns = mergeHandle.getUpdateCaseColumns().get(caseNumber); - if (updateColumns == null) continue; - - Set updateChannelsSet = updateColumns.stream() - .map(JdbcColumnHandle.class::cast) - .map(columns::indexOf) - .collect(toImmutableSet()); - - List updateChannelsList = new ArrayList<>(); - for (int channel = 0; channel < columns.size(); channel++) { - if (updateChannelsSet.contains(channel)) { - updateChannelsList.add(channel); - } - } - - int[] updateChannels = updateChannelsList.stream().mapToInt(Integer::intValue).toArray(); - - Block rowIdBlock = page.getBlock(columnCount + 2); - List rowIdFields = RowBlock.getRowFieldsFromBlock(rowIdBlock); - Block[] updateBlocks = new Block[updateChannels.length + rowIdFields.size()]; - - int blockIdx = 0; - for (int channel : updateChannels) { - updateBlocks[blockIdx++] = page.getBlock(channel).getPositions(positions, 0, positions.length); - } - for (Block rowIdField : rowIdFields) { - updateBlocks[blockIdx++] = rowIdField.getPositions(positions, 0, positions.length); - } - - updatePagesByCase.computeIfAbsent(caseNumber, _ -> new ArrayList<>()) - .add(new Page(positions.length, updateBlocks)); - } - } - - private List getColumnWriters(List columnTypes) { - return columnTypes.stream() - .map(type -> jdbcClient.toWriteMapping(session, type).getWriteFunction()) - .collect(java.util.stream.Collectors.toList()); - } - - private void setParameter(PreparedStatement stmt, int index, Block block, int position, Type type, WriteFunction writer) throws SQLException { - if (block.isNull(position)) { - writer.setNull(stmt, index); - return; - } - - Class javaType = type.getJavaType(); - if (javaType == boolean.class) { - ((BooleanWriteFunction) writer).set(stmt, index, type.getBoolean(block, position)); - } - else if (javaType == long.class) { - ((LongWriteFunction) writer).set(stmt, index, type.getLong(block, position)); - } - else if (javaType == double.class) { - ((DoubleWriteFunction) writer).set(stmt, index, type.getDouble(block, position)); - } - else if (javaType == io.airlift.slice.Slice.class) { - ((SliceWriteFunction) writer).set(stmt, index, type.getSlice(block, position)); - } - else { - ((ObjectWriteFunction) writer).set(stmt, index, type.getObject(block, position)); - } - } - - private static boolean isRejectedDriverSessionAcquisition(Throwable error) { - for (Throwable cause = error; cause != null; cause = cause.getCause()) { - if (cause instanceof RejectedExecutionException) { - for (StackTraceElement frame : cause.getStackTrace()) { - if (frame.getClassName().equals("tech.ydb.table.impl.pool.SessionPool") && - frame.getMethodName().equals("acquire")) { - // Session acquisition was rejected locally before a statement could be sent to YDB. - return true; - } - } - } - } - return false; - } - - @Override - public void abort() { - finished = true; - bufferedPages.clear(); - } -} diff --git a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbPageSinkProvider.java b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbPageSinkProvider.java deleted file mode 100644 index 2880dcf0..00000000 --- a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbPageSinkProvider.java +++ /dev/null @@ -1,80 +0,0 @@ -package tech.ydb.trino; - -import com.google.inject.Inject; -import io.trino.plugin.jdbc.JdbcClient; -import io.trino.plugin.jdbc.JdbcOutputTableHandle; -import io.trino.plugin.jdbc.JdbcPageSink; -import io.trino.plugin.jdbc.QueryBuilder; -import io.trino.plugin.jdbc.logging.RemoteQueryModifier; -import io.trino.spi.connector.ConnectorInsertTableHandle; -import io.trino.spi.connector.ConnectorMergeSink; -import io.trino.spi.connector.ConnectorMergeTableHandle; -import io.trino.spi.connector.ConnectorOutputTableHandle; -import io.trino.spi.connector.ConnectorPageSink; -import io.trino.spi.connector.ConnectorPageSinkId; -import io.trino.spi.connector.ConnectorPageSinkProvider; -import io.trino.spi.connector.ConnectorSession; -import io.trino.spi.connector.ConnectorTableCredentials; -import io.trino.spi.connector.ConnectorTransactionHandle; - -import java.util.Optional; - -public record YdbPageSinkProvider( - JdbcClient jdbcClient, - RemoteQueryModifier remoteQueryModifier, - QueryBuilder queryBuilder -) implements ConnectorPageSinkProvider { - @Inject - public YdbPageSinkProvider { - - } - - @Override - public ConnectorPageSink createPageSink( - ConnectorTransactionHandle transactionHandle, - ConnectorSession session, - ConnectorOutputTableHandle tableHandle, - Optional tableCredentials, - ConnectorPageSinkId pageSinkId) { - return new JdbcPageSink( - session, - (JdbcOutputTableHandle) tableHandle, - jdbcClient, - pageSinkId, - remoteQueryModifier, - JdbcClient::buildInsertSql); - } - - @Override - public ConnectorPageSink createPageSink( - ConnectorTransactionHandle transactionHandle, - ConnectorSession session, - ConnectorInsertTableHandle tableHandle, - Optional tableCredentials, - ConnectorPageSinkId pageSinkId) { - return new JdbcPageSink( - session, - (JdbcOutputTableHandle) tableHandle, - jdbcClient, - pageSinkId, - remoteQueryModifier, - JdbcClient::buildInsertSql); - } - - @Override - public ConnectorMergeSink createMergeSink( - ConnectorTransactionHandle transactionHandle, - ConnectorSession session, - ConnectorMergeTableHandle mergeHandle, - Optional tableCredentials, - ConnectorPageSinkId pageSinkId) { - return new YdbMergeSink( - transactionHandle, - session, - mergeHandle, - jdbcClient, - pageSinkId, - remoteQueryModifier, - queryBuilder); - } -} diff --git a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbRetryUtils.java b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbRetryUtils.java deleted file mode 100644 index bd7af20b..00000000 --- a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbRetryUtils.java +++ /dev/null @@ -1,30 +0,0 @@ -package tech.ydb.trino; - -import tech.ydb.jdbc.exception.YdbStatusable; - -import java.util.concurrent.ThreadLocalRandom; - -public final class YdbRetryUtils { - private static final long BASE_DELAY_MS = 20; - private static final long MAX_DELAY_MS = 1000; - - private YdbRetryUtils() { - } - - public static boolean isRetryable(Throwable error) { - Throwable current = error; - while (current != null) { - if (current instanceof YdbStatusable statusable) { - return statusable.getStatus().getCode().isRetryable(false); - } - current = current.getCause(); - } - return false; - } - - public static long calculateBackoff(int attempt) { - long exponentialBackoff = BASE_DELAY_MS * (1L << Math.min(attempt, 10)); - long cappedBackoff = Math.min(exponentialBackoff, MAX_DELAY_MS); - return ThreadLocalRandom.current().nextLong(0, cappedBackoff + 1); - } -} diff --git a/ydb-trino-adapter/src/test/java/tech/ydb/trino/TestingYdbJdbcModule.java b/ydb-trino-adapter/src/test/java/tech/ydb/trino/TestingYdbJdbcModule.java index aba3ff8f..dc71a297 100644 --- a/ydb-trino-adapter/src/test/java/tech/ydb/trino/TestingYdbJdbcModule.java +++ b/ydb-trino-adapter/src/test/java/tech/ydb/trino/TestingYdbJdbcModule.java @@ -25,7 +25,6 @@ public void configure(Binder binder) { .to(YdbMetadataFactory.class) .in(Scopes.SINGLETON); - binder.bind(YdbPageSinkProvider.class).in(Scopes.SINGLETON); binder.bind(YdbConnector.class).in(Scopes.SINGLETON); }