From 2238d823d7d0ebb1b0252067e00513a3d227a9ec Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Fri, 11 Sep 2026 13:55:12 +0300 Subject: [PATCH 1/6] Remove YDB MERGE replay --- ydb-trino-adapter/ROADMAP.md | 23 ++--- .../java/tech/ydb/trino/YdbMergeSink.java | 94 ++++++------------- .../java/tech/ydb/trino/YdbRetryUtils.java | 30 ------ 3 files changed, 37 insertions(+), 110 deletions(-) delete mode 100644 ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbRetryUtils.java diff --git a/ydb-trino-adapter/ROADMAP.md b/ydb-trino-adapter/ROADMAP.md index 89fe9311..d3c96f75 100644 --- a/ydb-trino-adapter/ROADMAP.md +++ b/ydb-trino-adapter/ROADMAP.md @@ -64,8 +64,7 @@ have been verified locally against a real YDB test container: - real JDBC primary-key metadata, including `KEY_SEQ` ordering; - concurrent updates identified by the test-only hidden primary key; - smoke MERGE and row-level UPDATE without disabling their behavior flags; -- retry classification through the YDB SDK status model, with fresh merge - connections and rollback-before-close. +- one connector attempt per MERGE, with rollback-before-close on failure. GitHub Actions is green on PR #240: 314 tests run, 0 failed, 0 errors, 84 skipped. This includes 36 smoke tests (4 skipped) and 278 connector tests (80 @@ -89,22 +88,18 @@ failure. 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, and rollback/close failure paths. **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. -## P1 — retry and transaction hardening +## P1 — transaction and replay 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 - [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. +- Do not replay buffered MERGE pages after any driver or YDB failure, including + session-acquisition rejection or an ambiguous commit result. The connector + performs one transaction attempt and preserves the original failure; a future + retry design requires staging or an operation ID that proves replay safety. - 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,8 +108,8 @@ 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. +- Add unit tests for rollback failure suppression, connection cleanup, and a + failure after commit. - Define memory/backpressure limits for buffered merge pages; memory usage must not remain unreported. 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 index d1f29ddd..1f836be3 100644 --- a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbMergeSink.java +++ b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbMergeSink.java @@ -37,7 +37,6 @@ 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; @@ -50,7 +49,6 @@ 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 { @@ -91,65 +89,44 @@ public void storeMergedRows(Page page) { 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); - } - } + Connection connection = null; + boolean committed = false; + try { + connection = openConnection(); + executeMergeInTransaction(connection); + connection.commit(); + committed = true; + connection.close(); + connection = null; + } + catch (Exception e) { + if (connection != null) { + if (!committed) { try { - connection.close(); + connection.rollback(); } - catch (SQLException closeError) { - e.addSuppressed(closeError); + catch (SQLException rollbackError) { + e.addSuppressed(rollbackError); } } - - 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); + connection.close(); } - catch (InterruptedException ie) { - Thread.currentThread().interrupt(); - throw new TrinoException(JDBC_ERROR, ie); + catch (SQLException closeError) { + e.addSuppressed(closeError); } } + + if (committed) { + throw new TrinoException(JDBC_ERROR, "YDB MERGE committed, but closing its connection failed", e); + } + + throw new TrinoException(JDBC_ERROR, e); } - throw new TrinoException(JDBC_ERROR, lastException); + Slice value = Slices.allocate(Long.BYTES); + value.setLong(0, pageSinkId.getId()); + return completedFuture(ImmutableList.of(value)); } private Connection openConnection() throws SQLException { @@ -453,21 +430,6 @@ else if (javaType == io.airlift.slice.Slice.class) { } } - 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; 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); - } -} From eeb1a587140b6301263ed093afe025e52a486260 Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Fri, 11 Sep 2026 13:59:46 +0300 Subject: [PATCH 2/6] Clarify historical MERGE baseline --- ydb-trino-adapter/ROADMAP.md | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/ydb-trino-adapter/ROADMAP.md b/ydb-trino-adapter/ROADMAP.md index d3c96f75..1e35cb3b 100644 --- a/ydb-trino-adapter/ROADMAP.md +++ b/ydb-trino-adapter/ROADMAP.md @@ -63,8 +63,7 @@ have been verified locally against a real YDB test container: - `Float`/`Double` write mappings during MERGE; - real JDBC primary-key metadata, including `KEY_SEQ` ordering; - concurrent updates identified by the test-only hidden primary key; -- smoke MERGE and row-level UPDATE without disabling their behavior flags; -- one connector attempt per MERGE, with rollback-before-close on failure. +- smoke MERGE and row-level UPDATE without disabling their behavior flags. GitHub Actions is green on PR #240: 314 tests run, 0 failed, 0 errors, 84 skipped. This includes 36 smoke tests (4 skipped) and 278 connector tests (80 From c3f765ef8348e6b6c92ac4df7af5778e1e196da3 Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Fri, 11 Sep 2026 15:00:35 +0300 Subject: [PATCH 3/6] Revert "Clarify historical MERGE baseline" This reverts commit eeb1a587140b6301263ed093afe025e52a486260. --- ydb-trino-adapter/ROADMAP.md | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/ydb-trino-adapter/ROADMAP.md b/ydb-trino-adapter/ROADMAP.md index 1e35cb3b..d3c96f75 100644 --- a/ydb-trino-adapter/ROADMAP.md +++ b/ydb-trino-adapter/ROADMAP.md @@ -63,7 +63,8 @@ have been verified locally against a real YDB test container: - `Float`/`Double` write mappings during MERGE; - real JDBC primary-key metadata, including `KEY_SEQ` ordering; - concurrent updates identified by the test-only hidden primary key; -- smoke MERGE and row-level UPDATE without disabling their behavior flags. +- smoke MERGE and row-level UPDATE without disabling their behavior flags; +- one connector attempt per MERGE, with rollback-before-close on failure. GitHub Actions is green on PR #240: 314 tests run, 0 failed, 0 errors, 84 skipped. This includes 36 smoke tests (4 skipped) and 278 connector tests (80 From 92d415e493436996e0aa1ed4da01edf433fa084c Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Fri, 11 Sep 2026 15:00:35 +0300 Subject: [PATCH 4/6] Revert "Remove YDB MERGE replay" This reverts commit 2238d823d7d0ebb1b0252067e00513a3d227a9ec. --- ydb-trino-adapter/ROADMAP.md | 23 +++-- .../java/tech/ydb/trino/YdbMergeSink.java | 94 +++++++++++++------ .../java/tech/ydb/trino/YdbRetryUtils.java | 30 ++++++ 3 files changed, 110 insertions(+), 37 deletions(-) create mode 100644 ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbRetryUtils.java diff --git a/ydb-trino-adapter/ROADMAP.md b/ydb-trino-adapter/ROADMAP.md index d3c96f75..89fe9311 100644 --- a/ydb-trino-adapter/ROADMAP.md +++ b/ydb-trino-adapter/ROADMAP.md @@ -64,7 +64,8 @@ have been verified locally against a real YDB test container: - real JDBC primary-key metadata, including `KEY_SEQ` ordering; - concurrent updates identified by the test-only hidden primary key; - smoke MERGE and row-level UPDATE without disabling their behavior flags; -- one connector attempt per MERGE, with rollback-before-close on failure. +- retry classification through the YDB SDK status model, with fresh merge + connections and rollback-before-close. GitHub Actions is green on PR #240: 314 tests run, 0 failed, 0 errors, 84 skipped. This includes 36 smoke tests (4 skipped) and 278 connector tests (80 @@ -88,18 +89,22 @@ failure. 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, and rollback/close failure paths. + column, physical-key updates, rollback/close, and fresh-state retries. **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. -## P1 — transaction and replay hardening +## P1 — retry and transaction hardening -- Do not replay buffered MERGE pages after any driver or YDB failure, including - session-acquisition rejection or an ambiguous commit result. The connector - performs one transaction attempt and preserves the original failure; a future - retry design requires staging or an operation ID that proves replay safety. +- 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 + [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 @@ -108,8 +113,8 @@ 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 rollback failure suppression, connection cleanup, and a - failure after commit. +- 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. 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 index 1f836be3..d1f29ddd 100644 --- a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbMergeSink.java +++ b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbMergeSink.java @@ -37,6 +37,7 @@ 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; @@ -49,6 +50,7 @@ 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 { @@ -89,44 +91,65 @@ public void storeMergedRows(Page page) { public CompletableFuture> finish() { finished = true; - Connection connection = null; - boolean committed = false; - try { - connection = openConnection(); - executeMergeInTransaction(connection); - connection.commit(); - committed = true; - connection.close(); - connection = null; - } - catch (Exception e) { - if (connection != null) { - if (!committed) { + 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.rollback(); + connection.close(); } - catch (SQLException rollbackError) { - e.addSuppressed(rollbackError); + catch (SQLException closeError) { + e.addSuppressed(closeError); } } - try { - connection.close(); + + if (committed) { + throw new TrinoException(JDBC_ERROR, "YDB MERGE committed, but closing its connection failed", e); } - catch (SQLException closeError) { - e.addSuppressed(closeError); + + if (!YdbRetryUtils.isRetryable(e) && !isRejectedDriverSessionAcquisition(e)) { + throw new TrinoException(JDBC_ERROR, e); } - } - if (committed) { - throw new TrinoException(JDBC_ERROR, "YDB MERGE committed, but closing its connection failed", e); - } + lastException = e; - throw new TrinoException(JDBC_ERROR, e); + long delay = YdbRetryUtils.calculateBackoff(attempt); + + try { + Thread.sleep(delay); + } + catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + throw new TrinoException(JDBC_ERROR, ie); + } + } } - Slice value = Slices.allocate(Long.BYTES); - value.setLong(0, pageSinkId.getId()); - return completedFuture(ImmutableList.of(value)); + throw new TrinoException(JDBC_ERROR, lastException); } private Connection openConnection() throws SQLException { @@ -430,6 +453,21 @@ else if (javaType == io.airlift.slice.Slice.class) { } } + 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; 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 new file mode 100644 index 00000000..bd7af20b --- /dev/null +++ b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbRetryUtils.java @@ -0,0 +1,30 @@ +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); + } +} From f4a29356d99f9286d4946d66c3d61b21c4aabade Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Fri, 11 Sep 2026 14:59:46 +0300 Subject: [PATCH 5/6] Use Trino JDBC merge sink --- ydb-trino-adapter/ROADMAP.md | 39 +++++++------------ .../main/java/tech/ydb/trino/YdbClient.java | 2 +- .../tech/ydb/trino/YdbPageSinkProvider.java | 4 +- 3 files changed, 18 insertions(+), 27 deletions(-) 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/YdbPageSinkProvider.java b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbPageSinkProvider.java index 2880dcf0..e82d1def 100644 --- a/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbPageSinkProvider.java +++ b/ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbPageSinkProvider.java @@ -2,6 +2,7 @@ import com.google.inject.Inject; import io.trino.plugin.jdbc.JdbcClient; +import io.trino.plugin.jdbc.JdbcMergeSink; import io.trino.plugin.jdbc.JdbcOutputTableHandle; import io.trino.plugin.jdbc.JdbcPageSink; import io.trino.plugin.jdbc.QueryBuilder; @@ -68,8 +69,7 @@ public ConnectorMergeSink createMergeSink( ConnectorMergeTableHandle mergeHandle, Optional tableCredentials, ConnectorPageSinkId pageSinkId) { - return new YdbMergeSink( - transactionHandle, + return new JdbcMergeSink( session, mergeHandle, jdbcClient, From 3c7d813d7df953e51ac4220dacba7c2063a633b7 Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Fri, 11 Sep 2026 15:18:01 +0300 Subject: [PATCH 6/6] Remove custom YDB merge sink --- .../java/tech/ydb/trino/YdbClientModule.java | 1 - .../java/tech/ydb/trino/YdbConnector.java | 2 +- .../java/tech/ydb/trino/YdbMergeSink.java | 476 ------------------ .../tech/ydb/trino/YdbPageSinkProvider.java | 80 --- .../java/tech/ydb/trino/YdbRetryUtils.java | 30 -- .../tech/ydb/trino/TestingYdbJdbcModule.java | 1 - 6 files changed, 1 insertion(+), 589 deletions(-) delete mode 100644 ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbMergeSink.java delete mode 100644 ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbPageSinkProvider.java delete mode 100644 ydb-trino-adapter/src/main/java/tech/ydb/trino/YdbRetryUtils.java 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 e82d1def..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.JdbcMergeSink; -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 JdbcMergeSink( - 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); }