From 3e628d2c7bbffe6780a758f110dd6547db601fef Mon Sep 17 00:00:00 2001 From: t8y2 <1156263951@qq.com> Date: Fri, 10 Jul 2026 22:16:46 +0800 Subject: [PATCH] fix(oceanbase): sync Oracle session query timeout --- .../java/com/dbx/agent/AbstractJdbcAgent.java | 15 +++- .../oceanbaseoracle/OceanBaseOracleAgent.java | 18 +++++ .../OceanBaseOracleAgentTest.java | 70 +++++++++++++++++++ 3 files changed, 100 insertions(+), 3 deletions(-) diff --git a/agents/common/src/main/java/com/dbx/agent/AbstractJdbcAgent.java b/agents/common/src/main/java/com/dbx/agent/AbstractJdbcAgent.java index 34d58434c..190db0d95 100644 --- a/agents/common/src/main/java/com/dbx/agent/AbstractJdbcAgent.java +++ b/agents/common/src/main/java/com/dbx/agent/AbstractJdbcAgent.java @@ -104,8 +104,10 @@ public abstract class AbstractJdbcAgent extends BaseDatabaseAgent { @Override public QueryResult executeQuery(String sql, String schema, ExecuteQueryOptions options) { + Connection conn = requireConnected(); + uncheckedVoid(() -> beforeQueryExecution(conn, options.getTimeoutSecs())); return JdbcExecutor.INSTANCE.execute( - requireConnected(), + conn, sql, schema, this::setSchemaSQL, @@ -118,8 +120,10 @@ public abstract class AbstractJdbcAgent extends BaseDatabaseAgent { @Override public QueryPageResult executeQueryPage(String sql, String schema, QueryPageOptions options) { + Connection conn = requireConnected(); + uncheckedVoid(() -> beforeQueryExecution(conn, options.getTimeoutSecs())); return JdbcExecutor.INSTANCE.executePage( - requireConnected(), + conn, sql, schema, this::setSchemaSQL, @@ -140,8 +144,10 @@ public abstract class AbstractJdbcAgent extends BaseDatabaseAgent { @Override public QueryPageResult startTableRead(String sql, String schema, QueryPageOptions options) { + Connection conn = requireConnected(); + uncheckedVoid(() -> beforeQueryExecution(conn, options.getTimeoutSecs())); return JdbcExecutor.INSTANCE.startTableRead( - requireConnected(), + conn, sql, schema, this::setSchemaSQL, @@ -186,6 +192,9 @@ public abstract class AbstractJdbcAgent extends BaseDatabaseAgent { protected void afterConnect(ConnectParams params, Connection connection) throws Exception { } + protected void beforeQueryExecution(Connection connection, int timeoutSecs) throws Exception { + } + protected String getConfiguredDatabase() { return configuredDatabase; } diff --git a/agents/drivers/oceanbase-oracle/src/main/java/com/dbx/agent/oceanbaseoracle/OceanBaseOracleAgent.java b/agents/drivers/oceanbase-oracle/src/main/java/com/dbx/agent/oceanbaseoracle/OceanBaseOracleAgent.java index 2e1356fe7..18d496472 100644 --- a/agents/drivers/oceanbase-oracle/src/main/java/com/dbx/agent/oceanbaseoracle/OceanBaseOracleAgent.java +++ b/agents/drivers/oceanbase-oracle/src/main/java/com/dbx/agent/oceanbaseoracle/OceanBaseOracleAgent.java @@ -15,6 +15,7 @@ import com.dbx.agent.ObjectInfo; import com.dbx.agent.TableInfo; import com.dbx.agent.TriggerInfo; +import java.sql.Connection; import java.sql.ResultSet; import java.sql.SQLException; import java.util.ArrayList; @@ -27,6 +28,7 @@ import java.util.Set; import java.util.stream.Collectors; public final class OceanBaseOracleAgent extends ConfiguredJdbcAgent { + private static final long MICROS_PER_SECOND = 1_000_000L; private static final String COMPATIBLE_OJDBC_VERSION = "compatibleOjdbcVersion"; private static final String DEFAULT_COMPATIBLE_OJDBC_VERSION = "compatibleOjdbcVersion=8"; private static final Set SYSTEM_SCHEMAS = Set.of( @@ -66,6 +68,22 @@ public final class OceanBaseOracleAgent extends ConfiguredJdbcAgent { return appendDefaultCompatibilityOption(OCEANBASE_ORACLE_PROFILE.buildUrl(params)); } + @Override + protected void beforeQueryExecution(Connection connection, int timeoutSecs) throws SQLException { + // Connector/J's Statement timeout does not update OceanBase's stricter + // session variable, so synchronize both limits before every execution. + try (var stmt = connection.createStatement()) { + stmt.execute(queryTimeoutSql(timeoutSecs)); + } + } + + static String queryTimeoutSql(int timeoutSecs) { + if (timeoutSecs < 0) { + throw new IllegalArgumentException("Query timeout cannot be negative: " + timeoutSecs); + } + return "ALTER SESSION SET ob_query_timeout = " + timeoutSecs * MICROS_PER_SECOND; + } + @Override public List listDatabases() { return unchecked(() -> { diff --git a/agents/drivers/oceanbase-oracle/src/test/java/com/dbx/agent/oceanbaseoracle/OceanBaseOracleAgentTest.java b/agents/drivers/oceanbase-oracle/src/test/java/com/dbx/agent/oceanbaseoracle/OceanBaseOracleAgentTest.java index 8bfd4c13a..0530b5ca7 100644 --- a/agents/drivers/oceanbase-oracle/src/test/java/com/dbx/agent/oceanbaseoracle/OceanBaseOracleAgentTest.java +++ b/agents/drivers/oceanbase-oracle/src/test/java/com/dbx/agent/oceanbaseoracle/OceanBaseOracleAgentTest.java @@ -2,8 +2,10 @@ package com.dbx.agent.oceanbaseoracle; import com.dbx.agent.ColumnInfo; import com.dbx.agent.ConnectParams; +import com.dbx.agent.ExecuteQueryOptions; import com.dbx.agent.MetadataListConstraints; import com.dbx.agent.ObjectInfo; +import com.dbx.agent.QueryPageOptions; import com.dbx.agent.TableInfo; import com.dbx.agent.test.TestSupport; import org.junit.jupiter.api.Assertions; @@ -15,6 +17,7 @@ import java.lang.reflect.Proxy; import java.sql.Connection; import java.sql.PreparedStatement; import java.sql.ResultSet; +import java.sql.Statement; import java.util.ArrayList; import java.util.List; @@ -71,6 +74,47 @@ class OceanBaseOracleAgentTest { ); } + @Test + void convertsQueryTimeoutToOceanBaseSessionMicroseconds() { + Assertions.assertEquals( + "ALTER SESSION SET ob_query_timeout = 300000000", + OceanBaseOracleAgent.queryTimeoutSql(300) + ); + Assertions.assertEquals( + "ALTER SESSION SET ob_query_timeout = 0", + OceanBaseOracleAgent.queryTimeoutSql(0) + ); + Assertions.assertEquals( + "ALTER SESSION SET ob_query_timeout = 2147483647000000", + OceanBaseOracleAgent.queryTimeoutSql(Integer.MAX_VALUE) + ); + } + + @Test + void rejectsNegativeQueryTimeout() { + Assertions.assertThrows(IllegalArgumentException.class, () -> OceanBaseOracleAgent.queryTimeoutSql(-1)); + } + + @Test + void synchronizesSessionTimeoutForEveryQueryEntryPoint() { + List sql = new ArrayList<>(); + OceanBaseOracleAgent agent = new OceanBaseOracleAgent(); + TestSupport.setPrivateConnection(agent, executionConnection(sql)); + + agent.executeQuery("SELECT 1 FROM DUAL", null, new ExecuteQueryOptions(10, null, 12)); + agent.executeQueryPage("SELECT 2 FROM DUAL", null, new QueryPageOptions(10, null, 10, 13)); + agent.startTableRead("SELECT 3 FROM DUAL", null, new QueryPageOptions(10, null, 10, 14)); + + Assertions.assertEquals(List.of( + "ALTER SESSION SET ob_query_timeout = 12000000", + "SELECT 1 FROM DUAL", + "ALTER SESSION SET ob_query_timeout = 13000000", + "SELECT 2 FROM DUAL", + "ALTER SESSION SET ob_query_timeout = 14000000", + "SELECT 3 FROM DUAL" + ), sql); + } + @Test void constrainedListTablesUsesOceanBaseOracleMetadataSql() { List sql = new ArrayList<>(); @@ -213,6 +257,32 @@ class OceanBaseOracleAgentTest { }); } + private static Connection executionConnection(List sql) { + Statement statement = proxy(Statement.class, (method, args) -> { + if ("execute".equals(method.getName())) { + sql.add(String.valueOf(args[0])); + return false; + } + if ("getUpdateCount".equals(method.getName())) { + return 0; + } + if ("close".equals(method.getName()) || "setMaxRows".equals(method.getName()) + || "setFetchSize".equals(method.getName()) || "setQueryTimeout".equals(method.getName())) { + return null; + } + return defaultValue(method.getReturnType()); + }); + return proxy(Connection.class, (method, args) -> { + if ("createStatement".equals(method.getName())) { + return statement; + } + if ("isClosed".equals(method.getName())) { + return false; + } + return defaultValue(method.getReturnType()); + }); + } + private static ResultSet resultSet(String[] columns, Object[][] rows) { int[] index = {-1}; return proxy(ResultSet.class, (method, args) -> {