From c1fe64cefa1ddb60d2525203dadeeafd35817456 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=99=88=E5=AD=90=E9=BB=98?= <925456043@qq.com> Date: Tue, 1 Sep 2026 17:02:08 +0800 Subject: [PATCH] =?UTF-8?q?perf:=20=E4=BC=98=E5=8C=96=E5=A4=9A=E6=BA=90?= =?UTF-8?q?=E8=81=94=E9=82=A6=E5=88=86=E7=89=87=E8=B0=83=E5=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 增加 Engine 级有界并行调度与资源收口 - 补齐调度饱和、失败快速传播及六源基准测试 --- .../Jdbc100mCursorManualBenchmarkTest.java | 129 +++++ .../jdbc/JdbcFederatedQueryEngineTest.java | 146 +++++- .../jdbc/JdbcSixSourceQueryEngineTest.java | 185 ++++++++ .../runtime/DefaultFederationSqlEngine.java | 10 +- .../sql/runtime/FederatedResultCursor.java | 1 + .../runtime/FederationExecutionSession.java | 442 ++++++++++++++++-- .../runtime/FederationFragmentScheduler.java | 118 +++++ .../runtime/QueryCancellationRegistry.java | 36 +- .../FederationFragmentSchedulerTest.java | 118 +++++ 9 files changed, 1137 insertions(+), 48 deletions(-) create mode 100644 easy-agents-federation-sql/easy-agents-federation-sql-adapter-jdbc/src/test/java/com/easyagents/federation/sql/adapter/jdbc/Jdbc100mCursorManualBenchmarkTest.java create mode 100644 easy-agents-federation-sql/easy-agents-federation-sql-adapter-jdbc/src/test/java/com/easyagents/federation/sql/adapter/jdbc/JdbcSixSourceQueryEngineTest.java create mode 100644 easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/FederationFragmentScheduler.java create mode 100644 easy-agents-federation-sql/easy-agents-federation-sql-core/src/test/java/com/easyagents/federation/sql/runtime/FederationFragmentSchedulerTest.java diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-adapter-jdbc/src/test/java/com/easyagents/federation/sql/adapter/jdbc/Jdbc100mCursorManualBenchmarkTest.java b/easy-agents-federation-sql/easy-agents-federation-sql-adapter-jdbc/src/test/java/com/easyagents/federation/sql/adapter/jdbc/Jdbc100mCursorManualBenchmarkTest.java new file mode 100644 index 0000000..f400833 --- /dev/null +++ b/easy-agents-federation-sql/easy-agents-federation-sql-adapter-jdbc/src/test/java/com/easyagents/federation/sql/adapter/jdbc/Jdbc100mCursorManualBenchmarkTest.java @@ -0,0 +1,129 @@ +package com.easyagents.federation.sql.adapter.jdbc; + +import com.easyagents.federation.sql.api.FederationSqlEngine; +import com.easyagents.federation.sql.api.FederationSqlEngines; +import com.easyagents.federation.sql.api.SqlQueryCommand; +import com.easyagents.federation.sql.execute.FederationResultCursor; +import com.easyagents.federation.sql.federation.FederationExecutionPolicy; +import com.easyagents.federation.sql.federation.FederationQueryScopeDefinition; +import com.easyagents.federation.sql.federation.FederationSourceBindingDefinition; +import com.easyagents.federation.sql.source.FederationDataSourceHandles; +import com.easyagents.federation.sql.source.FederationSourceDefinition; +import com.easyagents.federation.sql.source.RuntimeFingerprint; +import com.easyagents.federation.sql.source.SourceApplyOptions; +import com.easyagents.federation.sql.source.SourceId; +import com.mysql.cj.jdbc.MysqlDataSource; +import org.junit.Assert; +import org.junit.Assume; +import org.junit.Test; + +import java.math.BigDecimal; +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.util.List; +import java.util.Map; + +/** + * 亿级数据集 100,000 行 Cursor 手工基准;默认跳过,必须显式启用。 + */ +public class Jdbc100mCursorManualBenchmarkTest { + + /** + * 从支付库顺序消费恰好 100,000 行,验证 Cursor 流式消费和资源上限。 + */ + @Test + public void shouldConsumeOneHundredThousandRows() { + Assume.assumeTrue(Boolean.getBoolean("federation.100m.enabled")); + String host = System.getProperty("federation.100m.mysql.host", "127.0.0.1"); + int port = Integer.getInteger("federation.100m.mysql.port", 33307); + String database = System.getProperty( + "federation.100m.mysql.database", "ef_bank_payment_perf"); + String username = System.getProperty("federation.100m.mysql.username", "ef_bench_ro"); + String password = System.getProperty("federation.100m.mysql.password"); + Assert.assertNotNull("必须通过系统属性提供只读密码", password); + + MysqlDataSource dataSource = new MysqlDataSource(); + dataSource.setUrl("jdbc:mysql://" + host + ':' + port + '/' + database + + "?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=Asia/Shanghai" + + "&useCursorFetch=true&defaultFetchSize=1000"); + dataSource.setUser(username); + dataSource.setPassword(password); + + SourceId sourceId = new SourceId("payment-100m"); + FederationSourceDefinition definition = new FederationSourceDefinition( + sourceId, 1, JdbcFederationSqlAdapterProvider.ADAPTER_ID, + List.of(new JdbcSchemaDefinition("MAIN", database, null)), Map.of()); + FederationExecutionPolicy policy = new FederationExecutionPolicy( + 6, 12, 2, 100_000, 64L * 1024L * 1024L, 60_000); + FederationQueryScopeDefinition scope = FederationQueryScopeDefinition.virtual( + "payment-cursor-100m", 1, + Map.of("PAY", FederationSourceBindingDefinition.of(sourceId, 1)), + "PAY", policy); + + CursorResult nativeResult = consumeNative(dataSource); + CursorResult federationResult; + long federationStarted = System.nanoTime(); + try (FederationSqlEngine engine = FederationSqlEngines.builder() + .dataSourceResolver(ignored -> FederationDataSourceHandles.shared( + dataSource, new RuntimeFingerprint("MySQL", "8.0", "JDBC", "8.4", "1"))) + .adapter(new JdbcFederationSqlAdapterProvider()) + .federationExecutionPolicy(policy) + .build()) { + engine.sources().apply(definition, SourceApplyOptions.prewarmNow()); + String sql = "SELECT payment_id, amount FROM PAY.MAIN.payment_order " + + "WHERE payment_id BETWEEN 1 AND 100000 ORDER BY payment_id"; + int rows = 0; + long amountCents = 0; + try (FederationResultCursor cursor = engine.query( + SqlQueryCommand.of(sql, scope, List.of()))) { + while (cursor.next()) { + rows++; + amountCents += ((BigDecimal) cursor.getObject(2)) + .movePointRight(2).longValue(); + } + } + federationResult = new CursorResult(rows, amountCents, + elapsedMillis(federationStarted)); + } + Assert.assertEquals(100_000, nativeResult.rows()); + Assert.assertEquals(100_000, federationResult.rows()); + Assert.assertEquals(nativeResult.amountCents(), federationResult.amountCents()); + Assert.assertTrue(federationResult.amountCents() > 0L); + System.out.println("BENCH_CURSOR_100K rows=" + federationResult.rows() + + " nativeMillis=" + nativeResult.elapsedMillis() + + " federationMillis=" + federationResult.elapsedMillis() + + " amountCents=" + federationResult.amountCents()); + } + + private CursorResult consumeNative(MysqlDataSource dataSource) { + String sql = "SELECT payment_id, amount FROM payment_order " + + "WHERE payment_id BETWEEN 1 AND 100000 ORDER BY payment_id"; + long started = System.nanoTime(); + int rows = 0; + long amountCents = 0; + try (Connection connection = dataSource.getConnection(); + PreparedStatement statement = connection.prepareStatement(sql, + ResultSet.TYPE_FORWARD_ONLY, ResultSet.CONCUR_READ_ONLY)) { + statement.setFetchSize(1000); + try (ResultSet resultSet = statement.executeQuery()) { + while (resultSet.next()) { + rows++; + amountCents += resultSet.getBigDecimal(2) + .movePointRight(2).longValue(); + } + } + } catch (SQLException exception) { + throw new IllegalStateException("原生 JDBC Cursor 基准执行失败", exception); + } + return new CursorResult(rows, amountCents, elapsedMillis(started)); + } + + private long elapsedMillis(long started) { + return (System.nanoTime() - started) / 1_000_000L; + } + + private record CursorResult(int rows, long amountCents, long elapsedMillis) { + } +} diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-adapter-jdbc/src/test/java/com/easyagents/federation/sql/adapter/jdbc/JdbcFederatedQueryEngineTest.java b/easy-agents-federation-sql/easy-agents-federation-sql-adapter-jdbc/src/test/java/com/easyagents/federation/sql/adapter/jdbc/JdbcFederatedQueryEngineTest.java index eac235d..de498bb 100644 --- a/easy-agents-federation-sql/easy-agents-federation-sql-adapter-jdbc/src/test/java/com/easyagents/federation/sql/adapter/jdbc/JdbcFederatedQueryEngineTest.java +++ b/easy-agents-federation-sql/easy-agents-federation-sql-adapter-jdbc/src/test/java/com/easyagents/federation/sql/adapter/jdbc/JdbcFederatedQueryEngineTest.java @@ -27,6 +27,7 @@ import com.easyagents.federation.sql.source.RuntimeFingerprint; import com.easyagents.federation.sql.source.SourceApplyOptions; import com.easyagents.federation.sql.source.SourceId; import java.math.BigDecimal; +import java.sql.SQLException; import java.time.Instant; import java.sql.Connection; import java.sql.Statement; @@ -34,6 +35,9 @@ import java.sql.Types; import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; import org.h2.jdbcx.JdbcDataSource; @@ -50,6 +54,13 @@ public class JdbcFederatedQueryEngineTest { private static final SourceId SALES_SOURCE = new SourceId("sales-source"); private static final SourceId BILLING_SOURCE = new SourceId("billing-source"); private static final SourceId REGION_SOURCE = new SourceId("region-source"); + private static final AtomicInteger ACTIVE_FRAGMENTS = new AtomicInteger(); + private static final AtomicInteger MAXIMUM_ACTIVE_FRAGMENTS = new AtomicInteger(); + private static final AtomicInteger FAILING_TRACKED_VALUE = new AtomicInteger( + Integer.MIN_VALUE + ); + private static final AtomicReference FRAGMENT_START_BARRIER = + new AtomicReference<>(new CountDownLatch(0)); private JdbcDataSource sales; private JdbcDataSource billing; @@ -80,7 +91,10 @@ public class JdbcFederatedQueryEngineTest { + "TIMESTAMP WITH TIME ZONE '2026-08-21 12:00:00.123456+08:00')", "CREATE TABLE PRECISE_EVENT (ID INT PRIMARY KEY, EVENT_AT TIMESTAMP(6))", "INSERT INTO PRECISE_EVENT VALUES " - + "(1, TIMESTAMP '2026-08-21 12:00:00.123456')"); + + "(1, TIMESTAMP '2026-08-21 12:00:00.123456')", + fragmentTrackingAlias(), + "CREATE VIEW TRACKED_VALUE AS SELECT TRACK_FRAGMENT(ID) AS TRACKED_KEY " + + "FROM CUSTOMER WHERE ID <= 2"); execute(billing, "CREATE TABLE ORDER_ITEM (ID INT PRIMARY KEY, CUSTOMER_ID INT NOT NULL, AMOUNT DECIMAL(12,2), CODE VARCHAR(64))", "INSERT INTO ORDER_ITEM VALUES (10, 1, 30.00, 'Alice'), (11, 1, 20.00, 'Alice'), (12, 2, 80.00, 'Bob')", @@ -90,10 +104,16 @@ public class JdbcFederatedQueryEngineTest { + "TIMESTAMP WITH TIME ZONE '2026-08-21 06:00:00.123456+02:00')", "CREATE TABLE PRECISE_EVENT (ID INT PRIMARY KEY, EVENT_AT TIMESTAMP(6))", "INSERT INTO PRECISE_EVENT VALUES " - + "(2, TIMESTAMP '2026-08-21 08:30:00.654321')"); + + "(2, TIMESTAMP '2026-08-21 08:30:00.654321')", + fragmentTrackingAlias(), + "CREATE VIEW TRACKED_VALUE AS SELECT TRACK_FRAGMENT(ID) AS TRACKED_KEY " + + "FROM ORDER_ITEM WHERE ID <= 11"); execute(region, "CREATE TABLE CUSTOMER_REGION (CUSTOMER_ID INT PRIMARY KEY, REGION VARCHAR(64))", - "INSERT INTO CUSTOMER_REGION VALUES (1, 'North'), (2, 'South'), (3, 'West')"); + "INSERT INTO CUSTOMER_REGION VALUES (1, 'North'), (2, 'South'), (3, 'West')", + fragmentTrackingAlias(), + "CREATE VIEW TRACKED_VALUE AS SELECT TRACK_FRAGMENT(CUSTOMER_ID) AS TRACKED_KEY " + + "FROM CUSTOMER_REGION WHERE CUSTOMER_ID <= 2"); RuntimeFingerprint fingerprint = new RuntimeFingerprint( "H2", "2", "H2 JDBC Driver", "2", "1" @@ -356,6 +376,91 @@ public class JdbcFederatedQueryEngineTest { ); } + /** + * 验证真实 JDBC Fragment 执行会重叠,且始终受查询级并发上限约束。 + */ + @Test + public void shouldBoundConcurrentJdbcFragments() { + ACTIVE_FRAGMENTS.set(0); + MAXIMUM_ACTIVE_FRAGMENTS.set(0); + FRAGMENT_START_BARRIER.set(new CountDownLatch(2)); + FederationQueryScopeDefinition threeSourceScope = + new FederationQueryScopeDefinition( + "tracked-fragments", + 1, + Map.of( + "SALES", FederationSourceBindingDefinition.of(SALES_SOURCE, 1), + "BILLING", FederationSourceBindingDefinition.of(BILLING_SOURCE, 1), + "REGION", FederationSourceBindingDefinition.of(REGION_SOURCE, 1) + ), + "SALES", + threeSourcePolicy() + ); + + try (FederationResultCursor cursor = engine.query(SqlQueryCommand.of( + "SELECT COUNT(*) FROM SALES.APP.TRACKED_VALUE s " + + "JOIN BILLING.APP.TRACKED_VALUE b ON s.TRACKED_KEY = b.TRACKED_KEY " + + "JOIN REGION.APP.TRACKED_VALUE r ON s.TRACKED_KEY = r.TRACKED_KEY", + threeSourceScope, + List.of() + ))) { + Assert.assertTrue(cursor.next()); + Assert.assertEquals(8L, cursor.getObject(1)); + Assert.assertFalse(cursor.next()); + } + + Assert.assertEquals(2, MAXIMUM_ACTIVE_FRAGMENTS.get()); + Assert.assertEquals(0, ACTIVE_FRAGMENTS.get()); + } + + /** + * 验证任一并行 Fragment 失败后立即传播原始执行错误,不等待未启动任务超时。 + */ + @Test + public void shouldFailFastWhenConcurrentFragmentFails() { + ACTIVE_FRAGMENTS.set(0); + MAXIMUM_ACTIVE_FRAGMENTS.set(0); + FRAGMENT_START_BARRIER.set(new CountDownLatch(0)); + FAILING_TRACKED_VALUE.set(10); + FederationQueryScopeDefinition threeSourceScope = + new FederationQueryScopeDefinition( + "tracked-fragment-failure", + 1, + Map.of( + "SALES", FederationSourceBindingDefinition.of(SALES_SOURCE, 1), + "BILLING", FederationSourceBindingDefinition.of(BILLING_SOURCE, 1), + "REGION", FederationSourceBindingDefinition.of(REGION_SOURCE, 1) + ), + "SALES", + new FederationExecutionPolicy( + 3, + 8, + 2, + 100_000, + 64L * 1024L * 1024L, + 1_000 + ) + ); + + try (FederationResultCursor cursor = engine.query(SqlQueryCommand.of( + "SELECT COUNT(*) FROM SALES.APP.TRACKED_VALUE s " + + "JOIN BILLING.APP.TRACKED_VALUE b ON s.TRACKED_KEY = b.TRACKED_KEY " + + "JOIN REGION.APP.TRACKED_VALUE r ON s.TRACKED_KEY = r.TRACKED_KEY", + threeSourceScope, + List.of() + ))) { + cursor.next(); + Assert.fail("failing fragment should terminate the federated query"); + } catch (FederationSqlException exception) { + Assert.assertEquals( + FederationSqlErrorCode.EXECUTION_FAILED, + exception.errorCode() + ); + } finally { + FAILING_TRACKED_VALUE.set(Integer.MIN_VALUE); + } + } + /** * 验证两个分片分别绑定原始查询中的动态参数。 */ @@ -717,6 +822,41 @@ public class JdbcFederatedQueryEngineTest { return dataSource; } + /** + * H2 视图调用的并发跟踪函数。 + * + * @param value 原始列值 + * @return 固定键值 + * @throws SQLException 分片未能在时限内重叠启动 + */ + public static int trackFragment(int value) throws SQLException { + int active = ACTIVE_FRAGMENTS.incrementAndGet(); + MAXIMUM_ACTIVE_FRAGMENTS.accumulateAndGet(active, Math::max); + CountDownLatch barrier = FRAGMENT_START_BARRIER.get(); + barrier.countDown(); + try { + if (value == FAILING_TRACKED_VALUE.get()) { + throw new SQLException("simulated federation fragment failure"); + } + if (!barrier.await(2, TimeUnit.SECONDS)) { + throw new SQLException("federation fragments did not overlap"); + } + Thread.sleep(100L); + return 1; + } catch (InterruptedException exception) { + Thread.currentThread().interrupt(); + throw new SQLException("fragment tracking was interrupted", exception); + } finally { + ACTIVE_FRAGMENTS.decrementAndGet(); + } + } + + private static String fragmentTrackingAlias() { + return "CREATE ALIAS TRACK_FRAGMENT FOR \"" + + JdbcFederatedQueryEngineTest.class.getName() + + ".trackFragment\""; + } + /** * 根据物理源选择测试数据库。 * diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-adapter-jdbc/src/test/java/com/easyagents/federation/sql/adapter/jdbc/JdbcSixSourceQueryEngineTest.java b/easy-agents-federation-sql/easy-agents-federation-sql-adapter-jdbc/src/test/java/com/easyagents/federation/sql/adapter/jdbc/JdbcSixSourceQueryEngineTest.java new file mode 100644 index 0000000..5c573f6 --- /dev/null +++ b/easy-agents-federation-sql/easy-agents-federation-sql-adapter-jdbc/src/test/java/com/easyagents/federation/sql/adapter/jdbc/JdbcSixSourceQueryEngineTest.java @@ -0,0 +1,185 @@ +package com.easyagents.federation.sql.adapter.jdbc; + +import com.easyagents.federation.sql.api.FederationSqlEngine; +import com.easyagents.federation.sql.api.FederationSqlEngines; +import com.easyagents.federation.sql.api.FederationSqlErrorCode; +import com.easyagents.federation.sql.api.FederationSqlException; +import com.easyagents.federation.sql.api.SqlQueryCommand; +import com.easyagents.federation.sql.compile.FederationSqlPlan; +import com.easyagents.federation.sql.compile.SqlCompileRequest; +import com.easyagents.federation.sql.execute.FederationResultCursor; +import com.easyagents.federation.sql.federation.FederationExecutionPolicy; +import com.easyagents.federation.sql.federation.FederationQueryMode; +import com.easyagents.federation.sql.federation.FederationQueryScopeDefinition; +import com.easyagents.federation.sql.federation.FederationSourceBindingDefinition; +import com.easyagents.federation.sql.source.FederationDataSourceHandles; +import com.easyagents.federation.sql.source.FederationSourceDefinition; +import com.easyagents.federation.sql.source.RuntimeFingerprint; +import com.easyagents.federation.sql.source.SourceApplyOptions; +import com.easyagents.federation.sql.source.SourceId; +import org.h2.jdbcx.JdbcDataSource; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; + +import java.sql.Connection; +import java.sql.Statement; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; + +/** + * 四源和六源 JDBC 多阶段 Join 集成测试。 + */ +public class JdbcSixSourceQueryEngineTest { + + private static final FederationExecutionPolicy SIX_SOURCE_POLICY = + new FederationExecutionPolicy(6, 12, 2, 100_000, + 64L * 1024L * 1024L, 60_000); + + private final Map dataSources = new LinkedHashMap<>(); + private final Map bindings = new LinkedHashMap<>(); + private FederationSqlEngine engine; + + /** + * 创建六个独立 H2 JDBC 数据源。 + */ + @Before + public void setUp() throws Exception { + for (int index = 1; index <= 6; index++) { + SourceId sourceId = new SourceId("source-" + index); + JdbcDataSource dataSource = dataSource("six-source-" + index); + execute(dataSource, + "DROP TABLE IF EXISTS SUBJECT_METRIC", + "CREATE TABLE SUBJECT_METRIC (SUBJECT_ID INT PRIMARY KEY, METRIC_VALUE INT NOT NULL)", + "INSERT INTO SUBJECT_METRIC VALUES (1, " + (index * 10) + + "), (2, " + (index * 100) + ")"); + dataSources.put(sourceId, dataSource); + bindings.put("S" + index, FederationSourceBindingDefinition.of(sourceId, 1)); + } + RuntimeFingerprint fingerprint = new RuntimeFingerprint( + "H2", "2", "H2 JDBC Driver", "2", "1"); + engine = FederationSqlEngines.builder() + .dataSourceResolver(definition -> FederationDataSourceHandles.shared( + dataSources.get(definition.sourceId()), fingerprint)) + .adapter(new JdbcFederationSqlAdapterProvider()) + .federationExecutionPolicy(SIX_SOURCE_POLICY) + .maximumPlanCacheEntries(32) + .build(); + for (SourceId sourceId : dataSources.keySet()) { + engine.sources().apply(definition(sourceId), SourceApplyOptions.prewarmNow()); + } + } + + /** + * 关闭 Engine。 + */ + @After + public void tearDown() { + if (engine != null) { + engine.close(); + } + } + + /** + * 验证同一六源 Scope 中只引用四源时生成四个 Fragment。 + */ + @Test + public void shouldCompileFourReferencedSources() { + FederationSqlPlan plan = engine.compile(SqlCompileRequest.of( + selectSql(4), scope(SIX_SOURCE_POLICY))); + + Assert.assertEquals(FederationQueryMode.FEDERATED, plan.queryMode()); + Assert.assertEquals(4, plan.referencedSources().size()); + Assert.assertEquals(4, plan.fragments().size()); + Assert.assertEquals(3, plan.joinOptimizations().size()); + } + + /** + * 验证六个独立 JDBC 数据源可完成五阶段 Join 并返回稳定结果。 + */ + @Test + public void shouldExecuteJoinAcrossSixSources() { + String sql = selectSql(6); + FederationSqlPlan plan = engine.compile(SqlCompileRequest.of( + sql, scope(SIX_SOURCE_POLICY))); + + Assert.assertEquals(6, plan.referencedSources().size()); + Assert.assertEquals(6, plan.fragments().size()); + Assert.assertEquals(5, plan.joinOptimizations().size()); + + try (FederationResultCursor cursor = engine.query( + SqlQueryCommand.of(sql, scope(SIX_SOURCE_POLICY), List.of()))) { + Assert.assertTrue(cursor.next()); + Assert.assertEquals(1, cursor.getObject(1)); + Assert.assertEquals(10, cursor.getObject(2)); + Assert.assertEquals(60, cursor.getObject(7)); + Assert.assertTrue(cursor.next()); + Assert.assertEquals(2, cursor.getObject(1)); + Assert.assertEquals(600, cursor.getObject(7)); + Assert.assertFalse(cursor.next()); + Assert.assertEquals(6, cursor.metrics().fragments().size()); + } + } + + /** + * 验证 Scope 收紧为四源时,六源 SQL 在编译期稳定拒绝。 + */ + @Test + public void shouldRejectSixSourcesWhenScopeLimitIsFour() { + FederationExecutionPolicy fourSourcePolicy = new FederationExecutionPolicy( + 4, 12, 2, 100_000, 64L * 1024L * 1024L, 60_000); + + try { + engine.compile(SqlCompileRequest.of(selectSql(6), scope(fourSourcePolicy))); + Assert.fail("expected source limit rejection"); + } catch (FederationSqlException exception) { + Assert.assertEquals(FederationSqlErrorCode.FEDERATION_RESOURCE_LIMIT_EXCEEDED, + exception.errorCode()); + } + } + + private FederationQueryScopeDefinition scope(FederationExecutionPolicy policy) { + return new FederationQueryScopeDefinition( + "six-source-scope", 1, bindings, "S1", policy); + } + + private static String selectSql(int sourceCount) { + StringBuilder sql = new StringBuilder("SELECT s1.SUBJECT_ID"); + for (int index = 1; index <= sourceCount; index++) { + sql.append(", s").append(index).append(".METRIC_VALUE"); + } + sql.append(" FROM S1.APP.SUBJECT_METRIC s1"); + for (int index = 2; index <= sourceCount; index++) { + sql.append(" JOIN S").append(index).append(".APP.SUBJECT_METRIC s") + .append(index).append(" ON s").append(index) + .append(".SUBJECT_ID = s1.SUBJECT_ID"); + } + return sql.append(" ORDER BY s1.SUBJECT_ID").toString(); + } + + private static JdbcDataSource dataSource(String name) { + JdbcDataSource dataSource = new JdbcDataSource(); + dataSource.setURL("jdbc:h2:mem:" + name + ";DB_CLOSE_DELAY=-1"); + dataSource.setUser("sa"); + dataSource.setPassword(""); + return dataSource; + } + + private static void execute(JdbcDataSource dataSource, String... statements) + throws Exception { + try (Connection connection = dataSource.getConnection(); + Statement statement = connection.createStatement()) { + for (String sql : statements) { + statement.execute(sql); + } + } + } + + private static FederationSourceDefinition definition(SourceId sourceId) { + return new FederationSourceDefinition( + sourceId, 1, JdbcFederationSqlAdapterProvider.ADAPTER_ID, + List.of(new JdbcSchemaDefinition("APP", null, "PUBLIC")), Map.of()); + } +} diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/DefaultFederationSqlEngine.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/DefaultFederationSqlEngine.java index 7d2abf8..bc6a332 100644 --- a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/DefaultFederationSqlEngine.java +++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/DefaultFederationSqlEngine.java @@ -70,6 +70,7 @@ public final class DefaultFederationSqlEngine implements FederationSqlEngine { private final CalciteFederationSqlCompiler compiler; private final AdapterFederationStatisticsProvider automaticStatisticsProvider; private final NodeMemoryAdmissionController nodeMemoryAdmission; + private final FederationFragmentScheduler fragmentScheduler; private final CalciteSqlCompleter completer = new CalciteSqlCompleter(); private final QueryCancellationRegistry cancellations = new QueryCancellationRegistry(); private final ScheduledThreadPoolExecutor deadlineScheduler = deadlineScheduler(); @@ -297,6 +298,7 @@ public final class DefaultFederationSqlEngine implements FederationSqlEngine { this.nodeMemoryAdmission = new NodeMemoryAdmissionController( maximumNodeIntermediateBytes ); + this.fragmentScheduler = new FederationFragmentScheduler(executionPolicy); this.planCache = new BoundedPlanCache( maximumPlanCacheEntries, maximumConcurrentCompilations, @@ -650,7 +652,8 @@ public final class DefaultFederationSqlEngine implements FederationSqlEngine { registration, compiler.effectivePolicy(plan.queryScope()), metrics, - queryDeadline + queryDeadline, + fragmentScheduler ); queryCursor = new FederatedResultCursor( plan, @@ -1080,6 +1083,11 @@ public final class DefaultFederationSqlEngine implements FederationSqlEngine { } catch (RuntimeException exception) { failure = exception; } + try { + fragmentScheduler.close(); + } catch (RuntimeException exception) { + failure = append(failure, exception); + } try { sourceManager.close(); } catch (RuntimeException exception) { diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/FederatedResultCursor.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/FederatedResultCursor.java index 82f80ee..5725a34 100644 --- a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/FederatedResultCursor.java +++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/FederatedResultCursor.java @@ -62,6 +62,7 @@ final class FederatedResultCursor this.session = session; this.metrics = metrics; this.maximumRows = maximumRows; + session.prepareFragments(); this.enumerator = plan.localBindable() .bind(new FederationDataContext( plan.localRootSchema(), diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/FederationExecutionSession.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/FederationExecutionSession.java index 2f9a58a..0d58136 100644 --- a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/FederationExecutionSession.java +++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/FederationExecutionSession.java @@ -20,14 +20,21 @@ import java.time.LocalDateTime; import java.time.LocalTime; import java.time.OffsetDateTime; import java.time.ZoneOffset; +import java.util.ArrayDeque; import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; +import java.util.concurrent.FutureTask; import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; import org.apache.calcite.avatica.util.ByteString; import org.apache.calcite.linq4j.AbstractEnumerable; import org.apache.calcite.linq4j.Enumerable; @@ -50,9 +57,19 @@ final class FederationExecutionSession implements AutoCloseable { private final FederationExecutionPolicy policy; private final FederationQueryMetricsTracker metrics; private final QueryDeadline deadline; + private final FederationFragmentScheduler fragmentScheduler; private final Semaphore fragmentSlots; private final Set openFragmentCursors = ConcurrentHashMap.newKeySet(); + private final Set> fragmentTasks = ConcurrentHashMap.newKeySet(); + private final Map> fragmentPreparations = + new ConcurrentHashMap<>(); + private final Object preparationLock = new Object(); + private final ArrayDeque pendingPreparations = + new ArrayDeque<>(); private final AtomicBoolean closed = new AtomicBoolean(); + private final AtomicReference preparationFailure = new AtomicReference<>(); + private int runningPreparations; + private boolean preparationFailed; /** * 创建查询执行会话。 @@ -64,6 +81,7 @@ final class FederationExecutionSession implements AutoCloseable { * @param policy 有效资源策略 * @param metrics 指标跟踪器 * @param deadline 请求级统一截止时间 + * @param fragmentScheduler Engine 级有界分片调度器 */ FederationExecutionSession( DefaultFederationSqlPlan plan, @@ -72,7 +90,8 @@ final class FederationExecutionSession implements AutoCloseable { QueryCancellationRegistry.QueryRegistration registration, FederationExecutionPolicy policy, FederationQueryMetricsTracker metrics, - QueryDeadline deadline + QueryDeadline deadline, + FederationFragmentScheduler fragmentScheduler ) { this.plan = plan; this.context = context; @@ -81,9 +100,41 @@ final class FederationExecutionSession implements AutoCloseable { this.policy = policy; this.metrics = metrics; this.deadline = deadline; + this.fragmentScheduler = fragmentScheduler; this.fragmentSlots = new Semaphore(policy.maximumConcurrentFragments(), true); } + /** + * 异步启动多源物理分片,在查询级中间结果预算内完成有界物化。 + */ + void prepareFragments() { + if (plan.fragments().size() < 2 || policy.maximumConcurrentFragments() < 2) { + return; + } + try { + for (FederationFragmentPlan fragment : plan.fragments()) { + CompiledFederationFragment compiled = plan.compiledFragments().get( + fragment.fragmentId() + ); + if (compiled == null) { + throw new FederationSqlException( + FederationSqlErrorCode.EXECUTION_FAILED, + "compiled fragment is missing: " + fragment.fragmentId() + ); + } + CompletableFuture preparation = new CompletableFuture<>(); + fragmentPreparations.put(fragment.fragmentId(), preparation); + synchronized (preparationLock) { + pendingPreparations.addLast(compiled); + } + } + schedulePreparations(); + } catch (RuntimeException | Error exception) { + cancelPreparation(); + throw exception; + } + } + /** * 为 Calcite ScannableTable 创建一个按需执行物理分片的 Enumerable。 * @@ -101,6 +152,13 @@ final class FederationExecutionSession implements AutoCloseable { return new AbstractEnumerable<>() { @Override public Enumerator enumerator() { + CompletableFuture preparation = + fragmentPreparations.get(fragmentId); + if (preparation != null) { + PreparedFragment prepared = awaitPrepared(preparation); + fragmentPreparations.remove(fragmentId, preparation); + return prepared.enumerator(); + } return openFragment(compiled); } }; @@ -138,6 +196,7 @@ final class FederationExecutionSession implements AutoCloseable { } registration.ensureNotCancelled(); } + throwPreparationFailure(); deadline.ensureAllowed(); } @@ -166,38 +225,260 @@ final class FederationExecutionSession implements AutoCloseable { private Enumerator openFragment(CompiledFederationFragment compiled) { ensureExecutionAllowed(); acquireFragmentSlot(); - FederationFragmentPlan fragment = compiled.plan(); - SourceRuntime runtime = runtimeSnapshot.runtime(fragment.bindingName()); try { - FederationResultCursor cursor = runtime.adapter().fragmentExecutor().execute( - new FederationFragmentExecutionContext( - context.queryId(), - fragment.executableSql(), - remapParameters(fragment, context.parameters()), - fragmentOptions(context.options()), - runtime.handle().dataSource(), - runtime.compatibility(), - runtime.definition().adapterOptions(), - registration.statementLifecycle(), - metrics.fragmentObserver(fragment.fragmentId()), - deadline - ) - ); - if (cursor == null || !context.queryId().equals(cursor.queryId())) { - closeQuietly(cursor); - throw new FederationSqlException( - FederationSqlErrorCode.EXECUTION_FAILED, - "adapter returned an invalid fragment cursor" - ); - } - openFragmentCursors.add(cursor); - return new FragmentEnumerator(compiled, cursor); + return new FragmentEnumerator(compiled, createFragmentCursor(compiled), true); } catch (RuntimeException exception) { fragmentSlots.release(); throw exception; } } + private void schedulePreparations() { + List tasks = new ArrayList<>(); + synchronized (preparationLock) { + while (!closed.get() + && !preparationFailed + && runningPreparations < policy.maximumConcurrentFragments() + && !pendingPreparations.isEmpty()) { + CompiledFederationFragment compiled = pendingPreparations.removeFirst(); + CompletableFuture preparation = + fragmentPreparations.get( + compiled.plan().fragmentId() + ); + if (preparation == null || preparation.isCancelled()) { + continue; + } + runningPreparations++; + tasks.add(new PreparationTask(compiled, preparation)); + } + } + for (PreparationTask task : tasks) { + submitPreparation(task); + } + } + + private void submitPreparation(PreparationTask preparationTask) { + FutureTask task = new FutureTask<>(() -> { + runPreparation(preparationTask); + return null; + }) { + @Override + protected void done() { + fragmentTasks.remove(this); + } + }; + fragmentTasks.add(task); + try { + fragmentScheduler.execute(task); + } catch (RuntimeException | Error exception) { + fragmentTasks.remove(task); + task.cancel(true); + preparationTask.preparation().completeExceptionally(exception); + failPreparations(exception); + preparationFinished(); + } + } + + private void runPreparation(PreparationTask task) { + try { + task.preparation().complete(prepareFragment(task.compiled())); + } catch (RuntimeException | Error exception) { + task.preparation().completeExceptionally(exception); + failPreparations(exception); + } finally { + preparationFinished(); + } + } + + private void preparationFinished() { + boolean scheduleNext; + synchronized (preparationLock) { + runningPreparations--; + scheduleNext = !closed.get() && !preparationFailed; + } + if (scheduleNext) { + schedulePreparations(); + } + } + + private void failPreparations(Throwable failure) { + if (!preparationFailure.compareAndSet(null, failure)) { + return; + } + synchronized (preparationLock) { + preparationFailed = true; + pendingPreparations.clear(); + } + for (CompletableFuture preparation : fragmentPreparations.values()) { + preparation.completeExceptionally(failure); + } + registration.terminateStatementsForCleanup(); + } + + private void throwPreparationFailure() { + Throwable failure = preparationFailure.get(); + if (failure == null) { + return; + } + if (failure instanceof RuntimeException runtimeException) { + throw runtimeException; + } + if (failure instanceof Error error) { + throw error; + } + throw new FederationSqlException( + FederationSqlErrorCode.EXECUTION_FAILED, + "federation fragment preparation failed", + failure + ); + } + + private PreparedFragment awaitPrepared(Future preparation) { + try { + while (true) { + ensureExecutionAllowed(); + long remaining = deadline.remainingNanos(); + if (remaining <= 0L) { + ensureExecutionAllowed(); + } + try { + return preparation.get( + Math.min(remaining, TimeUnit.MILLISECONDS.toNanos(50)), + TimeUnit.NANOSECONDS + ); + } catch (TimeoutException ignored) { + // 周期性复核共享 Deadline 和取消状态。 + } + } + } catch (InterruptedException exception) { + Thread.currentThread().interrupt(); + cancelPreparation(); + throw new FederationSqlException( + FederationSqlErrorCode.EXECUTION_FAILED, + "waiting for federation fragment preparation was interrupted", + exception + ); + } catch (ExecutionException exception) { + cancelPreparation(); + Throwable cause = exception.getCause(); + if (cause instanceof RuntimeException runtimeException) { + throw runtimeException; + } + if (cause instanceof Error error) { + throw error; + } + throw new FederationSqlException( + FederationSqlErrorCode.EXECUTION_FAILED, + "federation fragment preparation failed", + cause + ); + } catch (RuntimeException | Error exception) { + cancelPreparation(); + throw exception; + } + } + + private PreparedFragment prepareFragment(CompiledFederationFragment compiled) { + ensureExecutionAllowed(); + FederationResultCursor cursor = null; + try { + cursor = createFragmentCursor(compiled); + List rows = new ArrayList<>(); + while (true) { + ensureExecutionAllowed(); + if (!cursor.next()) { + releaseFragment(compiled, cursor, true, false); + cursor = null; + return new PreparedFragment(rows); + } + Object[] row = normalizeRow(cursor, compiled.rowType()); + recordFragmentIntermediate(compiled, row); + rows.add(row); + } + } catch (RuntimeException | Error exception) { + if (cursor != null) { + releaseFragment(compiled, cursor, false, false); + } + throw exception; + } + } + + private FederationResultCursor createFragmentCursor( + CompiledFederationFragment compiled + ) { + FederationFragmentPlan fragment = compiled.plan(); + SourceRuntime runtime = runtimeSnapshot.runtime(fragment.bindingName()); + FederationResultCursor cursor = runtime.adapter().fragmentExecutor().execute( + new FederationFragmentExecutionContext( + context.queryId(), + fragment.executableSql(), + remapParameters(fragment, context.parameters()), + fragmentOptions(context.options()), + runtime.handle().dataSource(), + runtime.compatibility(), + runtime.definition().adapterOptions(), + registration.statementLifecycle(), + metrics.fragmentObserver(fragment.fragmentId()), + deadline + ) + ); + if (cursor == null || !context.queryId().equals(cursor.queryId())) { + closeQuietly(cursor); + throw new FederationSqlException( + FederationSqlErrorCode.EXECUTION_FAILED, + "adapter returned an invalid fragment cursor" + ); + } + openFragmentCursors.add(cursor); + try { + ensureExecutionAllowed(); + return cursor; + } catch (RuntimeException | Error exception) { + openFragmentCursors.remove(cursor); + closeQuietly(cursor); + throw exception; + } + } + + private void recordFragmentIntermediate( + CompiledFederationFragment compiled, + Object[] row + ) { + FederationQueryMetricsTracker.IntermediateUsage usage = + metrics.recordIntermediate(compiled.plan().fragmentId(), row); + if (usage.rows() > policy.maximumIntermediateRows() + || usage.bytes() > policy.maximumIntermediateBytes()) { + throw new FederationSqlException( + FederationSqlErrorCode.FEDERATION_RESOURCE_LIMIT_EXCEEDED, + "federation intermediate result limit was exceeded" + ); + } + } + + private void cancelPreparation() { + synchronized (preparationLock) { + preparationFailed = true; + pendingPreparations.clear(); + } + for (Future task : List.copyOf(fragmentTasks)) { + task.cancel(true); + } + fragmentTasks.clear(); + for (FederationResultCursor cursor : List.copyOf(openFragmentCursors)) { + openFragmentCursors.remove(cursor); + try { + cursor.close(); + } catch (RuntimeException ignored) { + // 保留触发取消的原始分片异常;Statement 取消通道继续收口剩余资源。 + } + } + for (CompletableFuture preparation : fragmentPreparations.values()) { + preparation.cancel(true); + } + fragmentPreparations.clear(); + registration.terminateStatementsForCleanup(); + } + private void acquireFragmentSlot() { while (true) { ensureExecutionAllowed(); @@ -259,6 +540,21 @@ final class FederationExecutionSession implements AutoCloseable { if (!closed.compareAndSet(false, true)) { return; } + boolean terminationRequired = !fragmentTasks.isEmpty() + || !openFragmentCursors.isEmpty(); + synchronized (preparationLock) { + terminationRequired = terminationRequired || runningPreparations > 0; + preparationFailed = true; + pendingPreparations.clear(); + } + for (Future task : List.copyOf(fragmentTasks)) { + task.cancel(true); + } + fragmentTasks.clear(); + for (CompletableFuture preparation : fragmentPreparations.values()) { + preparation.cancel(true); + } + fragmentPreparations.clear(); RuntimeException failure = null; for (FederationResultCursor cursor : List.copyOf(openFragmentCursors)) { try { @@ -273,6 +569,9 @@ final class FederationExecutionSession implements AutoCloseable { openFragmentCursors.remove(cursor); } } + if (terminationRequired) { + registration.terminateStatementsForCleanup(); + } if (failure != null) { throw failure; } @@ -282,15 +581,18 @@ final class FederationExecutionSession implements AutoCloseable { private final CompiledFederationFragment compiled; private final FederationResultCursor cursor; + private final boolean fragmentSlotHeld; private final AtomicBoolean released = new AtomicBoolean(); private Object[] current; private FragmentEnumerator( CompiledFederationFragment compiled, - FederationResultCursor cursor + FederationResultCursor cursor, + boolean fragmentSlotHeld ) { this.compiled = compiled; this.cursor = cursor; + this.fragmentSlotHeld = fragmentSlotHeld; } /** {@inheritDoc} */ @@ -309,15 +611,7 @@ final class FederationExecutionSession implements AutoCloseable { return false; } Object[] row = normalizeRow(cursor, compiled.rowType()); - FederationQueryMetricsTracker.IntermediateUsage usage = - metrics.recordIntermediate(compiled.plan().fragmentId(), row); - if (usage.rows() > policy.maximumIntermediateRows() - || usage.bytes() > policy.maximumIntermediateBytes()) { - throw new FederationSqlException( - FederationSqlErrorCode.FEDERATION_RESOURCE_LIMIT_EXCEEDED, - "federation intermediate result limit was exceeded" - ); - } + recordFragmentIntermediate(compiled, row); current = row; return true; } catch (RuntimeException exception) { @@ -346,7 +640,9 @@ final class FederationExecutionSession implements AutoCloseable { try { cursor.close(); } finally { - fragmentSlots.release(); + if (fragmentSlotHeld) { + fragmentSlots.release(); + } if (exhausted) { metrics.finishFragment(compiled.plan().fragmentId()); } @@ -354,6 +650,78 @@ final class FederationExecutionSession implements AutoCloseable { } } + private static final class PreparedFragment { + + private final List rows; + + private PreparedFragment(List rows) { + this.rows = rows; + } + + private Enumerator enumerator() { + return new Enumerator<>() { + + private int index = -1; + private Object[] current; + + @Override + public Object[] current() { + return current; + } + + @Override + public boolean moveNext() { + int next = index + 1; + if (next >= rows.size()) { + current = null; + return false; + } + index = next; + current = rows.set(index, null); + return true; + } + + @Override + public void reset() { + throw new UnsupportedOperationException( + "fragment cursor cannot be reset" + ); + } + + @Override + public void close() { + rows.clear(); + current = null; + } + }; + } + } + + private record PreparationTask( + CompiledFederationFragment compiled, + CompletableFuture preparation + ) { + } + + private void releaseFragment( + CompiledFederationFragment compiled, + FederationResultCursor cursor, + boolean exhausted, + boolean fragmentSlotHeld + ) { + openFragmentCursors.remove(cursor); + try { + cursor.close(); + } finally { + if (fragmentSlotHeld) { + fragmentSlots.release(); + } + if (exhausted) { + metrics.finishFragment(compiled.plan().fragmentId()); + } + } + } + private static Object[] normalizeRow( FederationResultCursor cursor, RelDataType rowType diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/FederationFragmentScheduler.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/FederationFragmentScheduler.java new file mode 100644 index 0000000..9a9b30e --- /dev/null +++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/FederationFragmentScheduler.java @@ -0,0 +1,118 @@ +package com.easyagents.federation.sql.runtime; + +import com.easyagents.federation.sql.api.FederationSqlErrorCode; +import com.easyagents.federation.sql.api.FederationSqlException; +import com.easyagents.federation.sql.federation.FederationExecutionPolicy; +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.Executor; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * Engine 级有界分片工作线程池;查询级并发度仍由有效执行策略单独控制。 + */ +final class FederationFragmentScheduler implements Executor, AutoCloseable { + + private static final int MAXIMUM_WORKERS = 32; + private static final int MINIMUM_QUEUE_CAPACITY = 64; + private static final int QUEUE_CAPACITY_PER_WORKER = 8; + + private final ThreadPoolExecutor executor; + private final AtomicBoolean closed = new AtomicBoolean(); + + /** + * 创建 Engine 独占的有界分片调度器。 + * + * @param policy Engine 联邦执行硬上限 + */ + FederationFragmentScheduler(FederationExecutionPolicy policy) { + int availableProcessors = Math.max(1, Runtime.getRuntime().availableProcessors()); + int workers = Math.min( + MAXIMUM_WORKERS, + Math.max(policy.maximumConcurrentFragments(), availableProcessors) + ); + int queueCapacity = Math.max( + MINIMUM_QUEUE_CAPACITY, + workers * QUEUE_CAPACITY_PER_WORKER + ); + this.executor = executor(workers, queueCapacity); + } + + FederationFragmentScheduler(int workers, int queueCapacity) { + if (workers <= 0) { + throw new IllegalArgumentException("workers must be positive"); + } + if (queueCapacity <= 0) { + throw new IllegalArgumentException("queueCapacity must be positive"); + } + this.executor = executor(workers, queueCapacity); + } + + private static ThreadPoolExecutor executor(int workers, int queueCapacity) { + AtomicInteger threadSequence = new AtomicInteger(); + return new ThreadPoolExecutor( + workers, + workers, + 0L, + TimeUnit.MILLISECONDS, + new ArrayBlockingQueue<>(queueCapacity), + runnable -> { + Thread thread = new Thread( + runnable, + "easy-agents-federation-fragment-" + threadSequence.incrementAndGet() + ); + thread.setDaemon(true); + return thread; + }, + new ThreadPoolExecutor.AbortPolicy() + ); + } + + /** + * 提交分片任务;队列饱和时明确拒绝,禁止在调用线程绕过 Worker 上限执行 JDBC。 + * + * @param command 分片任务 + */ + @Override + public void execute(Runnable command) { + if (closed.get()) { + throw engineClosed(); + } + try { + executor.execute(command); + } catch (RejectedExecutionException exception) { + if (closed.get() || executor.isShutdown()) { + throw engineClosed(); + } + throw schedulerOverloaded(exception); + } + } + + /** + * 中断排队任务并关闭工作线程。 + */ + @Override + public void close() { + if (closed.compareAndSet(false, true)) { + executor.shutdownNow(); + } + } + + private static FederationSqlException engineClosed() { + return new FederationSqlException( + FederationSqlErrorCode.ENGINE_CLOSED, + "federation fragment scheduler is closed" + ); + } + + private static FederationSqlException schedulerOverloaded(Throwable cause) { + return new FederationSqlException( + FederationSqlErrorCode.QUERY_ADMISSION_TIMEOUT, + "federation fragment scheduler is saturated", + cause + ); + } +} diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/QueryCancellationRegistry.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/QueryCancellationRegistry.java index 876b3b3..1b46374 100644 --- a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/QueryCancellationRegistry.java +++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/runtime/QueryCancellationRegistry.java @@ -301,6 +301,13 @@ final class QueryCancellationRegistry implements AutoCloseable { return state.cancellationRequested(); } + /** + * 不改变查询终态,仅终止当前已登记的 Statement,用于提前关闭异步分片。 + */ + void terminateStatementsForCleanup() { + state.terminateStatements(); + } + /** * 已取消时立即抛出稳定错误。 */ @@ -485,13 +492,18 @@ final class QueryCancellationRegistry implements AutoCloseable { } private void terminate(QueryId queryId) { + if (isClosed()) { + return; + } SQLException failure = null; if (cancelIssued.compareAndSet(false, true)) { try { statement.cancel(); } catch (SQLException exception) { - cancelIssued.set(false); - failure = exception; + if (!isClosed()) { + cancelIssued.set(false); + failure = exception; + } } } // close() 同时覆盖 register 与 executeQuery 之间的 JDBC 取消空窗。 @@ -499,11 +511,13 @@ final class QueryCancellationRegistry implements AutoCloseable { try { statement.close(); } catch (SQLException exception) { - closeIssued.set(false); - if (failure == null) { - failure = exception; - } else { - failure.addSuppressed(exception); + if (!isClosed()) { + closeIssued.set(false); + if (failure == null) { + failure = exception; + } else { + failure.addSuppressed(exception); + } } } } @@ -515,6 +529,14 @@ final class QueryCancellationRegistry implements AutoCloseable { ); } } + + private boolean isClosed() { + try { + return statement.isClosed(); + } catch (SQLException ignored) { + return false; + } + } } private static FederationSqlException cancelled(QueryId queryId, Throwable cause) { diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/test/java/com/easyagents/federation/sql/runtime/FederationFragmentSchedulerTest.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/test/java/com/easyagents/federation/sql/runtime/FederationFragmentSchedulerTest.java new file mode 100644 index 0000000..fcda3c7 --- /dev/null +++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/test/java/com/easyagents/federation/sql/runtime/FederationFragmentSchedulerTest.java @@ -0,0 +1,118 @@ +package com.easyagents.federation.sql.runtime; + +import com.easyagents.federation.sql.api.FederationSqlErrorCode; +import com.easyagents.federation.sql.api.FederationSqlException; +import com.easyagents.federation.sql.federation.FederationExecutionPolicy; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorCompletionService; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import org.junit.Assert; +import org.junit.Test; + +/** + * Engine 级有界分片调度器测试。 + */ +public class FederationFragmentSchedulerTest { + + /** + * 验证两个分片任务能够真实重叠执行,且调度器关闭后拒绝新任务。 + * + * @throws Exception 等待任务失败 + */ + @Test + public void shouldExecuteFragmentTasksConcurrentlyAndRejectAfterClose() throws Exception { + FederationFragmentScheduler scheduler = new FederationFragmentScheduler( + FederationExecutionPolicy.basic() + ); + CountDownLatch started = new CountDownLatch(2); + CountDownLatch release = new CountDownLatch(1); + AtomicInteger active = new AtomicInteger(); + AtomicInteger maximumActive = new AtomicInteger(); + ExecutorCompletionService completions = new ExecutorCompletionService<>(scheduler); + try { + Future first = completions.submit( + () -> runBlockingTask(started, release, active, maximumActive) + ); + Future second = completions.submit( + () -> runBlockingTask(started, release, active, maximumActive) + ); + + Assert.assertTrue(started.await(2, TimeUnit.SECONDS)); + Assert.assertEquals(2, maximumActive.get()); + release.countDown(); + first.get(2, TimeUnit.SECONDS); + second.get(2, TimeUnit.SECONDS); + } finally { + release.countDown(); + scheduler.close(); + } + + try { + scheduler.execute(() -> { }); + Assert.fail("closed scheduler should reject new tasks"); + } catch (FederationSqlException exception) { + Assert.assertEquals(FederationSqlErrorCode.ENGINE_CLOSED, exception.errorCode()); + } + } + + /** + * 验证队列饱和时明确拒绝,且不会在提交线程绕过 Worker 上限执行任务。 + * + * @throws Exception 等待任务失败 + */ + @Test + public void shouldRejectSaturationWithoutRunningTaskOnCaller() throws Exception { + FederationFragmentScheduler scheduler = new FederationFragmentScheduler(1, 1); + CountDownLatch firstStarted = new CountDownLatch(1); + CountDownLatch releaseFirst = new CountDownLatch(1); + CountDownLatch queuedCompleted = new CountDownLatch(1); + AtomicBoolean rejectedTaskRan = new AtomicBoolean(); + try { + scheduler.execute(() -> { + firstStarted.countDown(); + try { + releaseFirst.await(2, TimeUnit.SECONDS); + } catch (InterruptedException exception) { + Thread.currentThread().interrupt(); + } + }); + Assert.assertTrue(firstStarted.await(2, TimeUnit.SECONDS)); + scheduler.execute(queuedCompleted::countDown); + + try { + scheduler.execute(() -> rejectedTaskRan.set(true)); + Assert.fail("saturated scheduler should reject the task"); + } catch (FederationSqlException exception) { + Assert.assertEquals( + FederationSqlErrorCode.QUERY_ADMISSION_TIMEOUT, + exception.errorCode() + ); + } + Assert.assertFalse(rejectedTaskRan.get()); + } finally { + releaseFirst.countDown(); + Assert.assertTrue(queuedCompleted.await(2, TimeUnit.SECONDS)); + scheduler.close(); + } + } + + private static Void runBlockingTask( + CountDownLatch started, + CountDownLatch release, + AtomicInteger active, + AtomicInteger maximumActive + ) throws InterruptedException { + int current = active.incrementAndGet(); + maximumActive.accumulateAndGet(current, Math::max); + started.countDown(); + try { + Assert.assertTrue(release.await(2, TimeUnit.SECONDS)); + return null; + } finally { + active.decrementAndGet(); + } + } +}