diff --git a/CHANGELOG.md b/CHANGELOG.md index a114e10..ee724de 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -33,6 +33,14 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/). ### Fixed +- Data task streaming memory and completeness: the single statement now executes + in an explicit transaction (auto-commit restored when handing the connection + back, DML committed first) so PostgreSQL-family drivers — including Hologres — + honor `fetch-size` and pull rows in bounded batches instead of buffering the + whole result set in heap, which caused `Java heap space` / repeated + `GC overhead limit exceeded` failures on large exports. Also fixed a dropped + trailing partial batch (fewer than 1000 rows) whenever a result set ended + normally; both are locked by `DataTaskJobEngineStreamQueryTest`. - Deployment timezone alignment: the compose MySQL now starts with `--default-time-zone=+08:00` (`build-docker/install`, `.devcontainer`) and the devcontainer PostgreSQL sets `timezone=Asia/Shanghai`, matching the JDBC `serverTimezone=Asia/Shanghai` diff --git a/build-docker/datapoly/datapoly-release/bin/datapolyctl.sh b/build-docker/datapoly/datapoly-release/bin/datapolyctl.sh index 4ed37a9..b70273c 100644 --- a/build-docker/datapoly/datapoly-release/bin/datapolyctl.sh +++ b/build-docker/datapoly/datapoly-release/bin/datapolyctl.sh @@ -21,10 +21,13 @@ echo "Base Directory:${APP_HOME}" export APP_DRIVERS_PATH=$APP_HOME/drivers # JVM参数可以在这里设置 +# 堆 4G、年轻代/老年代 1:3:长驻对象(Hazelcast token/API 响应缓存、Eureka、Spring 框架 +# 对象)占堆内大头,老年代空间优先;年轻代 1G 足以容纳数据任务流式批次与 ≤200 行的 +# 调试/预览结果集(线上观测:16 分钟仅 7 次 minor GC,老年代却 30 秒内 5→25 次 major GC)。 # -XX:+PerfDisableSharedMem: the JDK perfdata file lands in java.io.tmpdir; when the host # bind-mounts /tmp (macOS Docker Desktop virtiofs), zeroing that 32KB mmap SIGBUSes the JVM # at startup. Disabling shared-mem perfdata removes the mmap entirely. -JVMFLAGS="-server -Xms1024m -Xmx1024m -Xmn1024m -XX:+DisableExplicitGC -XX:+PerfDisableSharedMem -Djava.awt.headless=true -Dfile.encoding=UTF-8 " +JVMFLAGS="-server -Xms4096m -Xmx4096m -Xmn1024m -XX:+DisableExplicitGC -XX:+PerfDisableSharedMem -Djava.awt.headless=true -Dfile.encoding=UTF-8 " if [ "$JAVA_HOME" != "" ]; then JAVA="$JAVA_HOME/bin/java" diff --git a/datapoly-core/src/main/java/com/cs/core/datatask/DataTaskJobEngine.java b/datapoly-core/src/main/java/com/cs/core/datatask/DataTaskJobEngine.java index 3f88dec..cf5867b 100644 --- a/datapoly-core/src/main/java/com/cs/core/datatask/DataTaskJobEngine.java +++ b/datapoly-core/src/main/java/com/cs/core/datatask/DataTaskJobEngine.java @@ -267,10 +267,20 @@ public Map debugPreview(DataTaskDefEntity def, Map executeBeforeQuery = spec.getProduct().getContext().getExecuteBeforeQuery(); LambdaUtils.ifDo(null != executeBeforeQuery, () -> executeBeforeQuery.accept(connection)); + // PostgreSQL-family drivers (including Hologres) ignore the ResultSet fetch + // size while autocommit is on and buffer the whole result set in heap. Move + // the single statement into an explicit transaction so setFetchSize below + // bounds every network fetch to that many rows. + if (connection.getAutoCommit()) { + connection.setAutoCommit(false); + autoCommitChanged = true; + } + PreparedStatement statement = connection.prepareStatement(spec.getSqlMeta().getSql()); try { statement.setQueryTimeout(spec.getTimeoutSeconds()); @@ -281,6 +291,9 @@ protected StreamResult streamQuery(StreamSpec spec, ResultChannel channel) throw } boolean hasResult = statement.execute(); if (!hasResult) { + // DML ran outside autocommit: commit it, otherwise handing the + // connection back rolls the transaction (and the rows) away. + connection.commit(); return StreamResult.update(statement.getUpdateCount()); } return drainResultSet(statement.getResultSet(), spec, channel); @@ -288,6 +301,13 @@ protected StreamResult streamQuery(StreamSpec spec, ResultChannel channel) throw statement.close(); } } finally { + if (autoCommitChanged) { + try { + connection.setAutoCommit(true); + } catch (SQLException e) { + // the connection is about to be closed; nothing left to salvage + } + } connection.close(); } } @@ -338,8 +358,10 @@ private StreamResult drainResultSet(ResultSet rs, StreamSpec spec, ResultChannel } } } - if (!reachedEnd && !stoppedByChannel && !truncated && !buffer.isEmpty()) { - // trailing partial batch smaller than WRITE_BATCH_ROWS + if (!stoppedByChannel && !truncated && !buffer.isEmpty()) { + // trailing partial batch smaller than WRITE_BATCH_ROWS: this also fires + // when the result set ends normally, otherwise the final rows < batch + // size are silently dropped from the delivery if (channel.batch(buffer)) { tickProgress(spec, total, lastFlush); } diff --git a/datapoly-core/src/test/java/com/cs/core/datatask/DataTaskJobEngineStreamQueryTest.java b/datapoly-core/src/test/java/com/cs/core/datatask/DataTaskJobEngineStreamQueryTest.java new file mode 100644 index 0000000..ca76148 --- /dev/null +++ b/datapoly-core/src/test/java/com/cs/core/datatask/DataTaskJobEngineStreamQueryTest.java @@ -0,0 +1,350 @@ +// Use of this source code is governed by a BSD-style license +package com.cs.core.datatask; + +import com.cs.common.datatask.ColumnMetadata; +import com.cs.common.enums.NamingStrategyEnum; +import com.cs.common.enums.ProductTypeEnum; +import com.cs.template.SqlMeta; +import com.zaxxer.hikari.HikariDataSource; +import org.junit.Assert; +import org.junit.Test; + +import java.lang.reflect.InvocationHandler; +import java.lang.reflect.Method; +import java.lang.reflect.Proxy; +import java.sql.Connection; +import java.sql.DatabaseMetaData; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.ResultSetMetaData; +import java.sql.Types; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +/** + * Locks the JDBC-level contract of {@link DataTaskJobEngine#streamQuery} using + * proxy doubles: the single statement runs inside an explicit transaction so + * PostgreSQL-family drivers honor the fetch size instead of buffering the whole + * result set, DML is committed before the connection goes back to the pool, and + * auto-commit is restored only when the engine flipped it. + */ +public class DataTaskJobEngineStreamQueryTest { + + /** Shared state between the proxy doubles plus a sequential call log. */ + private static class Trace { + final List calls = new ArrayList<>(); + final List fetchSizes = new ArrayList<>(); + boolean autoCommit = true; + boolean autoCommitAtExecute; + boolean committed; + boolean closed; + boolean rowset; + boolean executed; + int updateCount; + int remainingRows; + Object cellValue; + String productName = "PostgreSQL"; + } + + private static class CaptureChannel implements DataTaskJobEngine.ResultChannel { + List columns; + final List rows = new ArrayList<>(); + + @Override + public void start(List columns, List metadata) { + this.columns = new ArrayList<>(columns); + } + + @Override + public boolean batch(List batch) { + rows.addAll(batch); + return true; + } + } + + private static final class FakeJdbc { + final Trace trace; + final Connection connection; + final HikariDataSource dataSource; + + FakeJdbc(Trace trace) { + this.trace = trace; + this.connection = connectionProxy(trace); + this.dataSource = new HikariDataSource() { + @Override + public Connection getConnection() { + return connection; + } + }; + } + } + + private static Connection connectionProxy(final Trace trace) { + final Connection[] self = new Connection[1]; + InvocationHandler handler = new InvocationHandler() { + @Override + public Object invoke(Object proxy, Method method, Object[] args) throws Throwable { + String name = method.getName(); + trace.calls.add("conn." + name); + if ("getAutoCommit".equals(name)) { + return trace.autoCommit; + } + if ("setAutoCommit".equals(name)) { + trace.autoCommit = (Boolean) args[0]; + return null; + } + if ("commit".equals(name)) { + trace.committed = true; + return null; + } + if ("close".equals(name)) { + trace.closed = true; + return null; + } + if ("prepareStatement".equals(name)) { + return statementProxy((String) args[0], trace); + } + if ("getMetaData".equals(name)) { + return databaseMetaDataProxy(trace); + } + return scalarDefault(method.getReturnType()); + } + }; + self[0] = (Connection) Proxy.newProxyInstance( + DataTaskJobEngineStreamQueryTest.class.getClassLoader(), + new Class[]{Connection.class}, handler); + return self[0]; + } + + private static PreparedStatement statementProxy(final String sql, final Trace trace) { + final ResultSet[] resultSet = new ResultSet[1]; + InvocationHandler handler = new InvocationHandler() { + @Override + public Object invoke(Object proxy, Method method, Object[] args) throws Throwable { + String name = method.getName(); + trace.calls.add("stmt." + name); + if ("setQueryTimeout".equals(name)) { + return null; + } + if ("setFetchSize".equals(name)) { + trace.fetchSizes.add((Integer) args[0]); + return null; + } + if ("setObject".equals(name)) { + return null; + } + if ("execute".equals(name)) { + trace.executed = true; + trace.autoCommitAtExecute = trace.autoCommit; + return trace.rowset; + } + if ("getUpdateCount".equals(name)) { + return trace.updateCount; + } + if ("getResultSet".equals(name)) { + resultSet[0] = resultSetProxy(trace); + return resultSet[0]; + } + if ("close".equals(name)) { + return null; + } + if ("toString".equals(name)) { + return "statement(" + sql + ")"; + } + return scalarDefault(method.getReturnType()); + } + }; + return (PreparedStatement) Proxy.newProxyInstance( + DataTaskJobEngineStreamQueryTest.class.getClassLoader(), + new Class[]{PreparedStatement.class}, handler); + } + + private static ResultSet resultSetProxy(final Trace trace) { + InvocationHandler handler = new InvocationHandler() { + @Override + public Object invoke(Object proxy, Method method, Object[] args) throws Throwable { + String name = method.getName(); + trace.calls.add("rs." + name); + if ("next".equals(name)) { + return trace.remainingRows-- > 0; + } + if ("getObject".equals(name)) { + return trace.cellValue; + } + if ("getMetaData".equals(name)) { + return resultSetMetaDataProxy(trace); + } + if ("close".equals(name)) { + return null; + } + return scalarDefault(method.getReturnType()); + } + }; + return (ResultSet) Proxy.newProxyInstance( + DataTaskJobEngineStreamQueryTest.class.getClassLoader(), + new Class[]{ResultSet.class}, handler); + } + + private static ResultSetMetaData resultSetMetaDataProxy(final Trace trace) { + InvocationHandler handler = new InvocationHandler() { + @Override + public Object invoke(Object proxy, Method method, Object[] args) throws Throwable { + String name = method.getName(); + trace.calls.add("meta." + name); + if ("getColumnCount".equals(name)) { + return 1; + } + if ("getColumnLabel".equals(name)) { + return "c1"; + } + if ("getColumnType".equals(name)) { + return Types.VARCHAR; + } + if ("getColumnClassName".equals(name)) { + return "java.lang.String"; + } + return scalarDefault(method.getReturnType()); + } + }; + return (ResultSetMetaData) Proxy.newProxyInstance( + DataTaskJobEngineStreamQueryTest.class.getClassLoader(), + new Class[]{ResultSetMetaData.class}, handler); + } + + private static DatabaseMetaData databaseMetaDataProxy(final Trace trace) { + InvocationHandler handler = new InvocationHandler() { + @Override + public Object invoke(Object proxy, Method method, Object[] args) throws Throwable { + String name = method.getName(); + trace.calls.add("dbmd." + name); + if ("getDatabaseProductName".equals(name)) { + return trace.productName; + } + return scalarDefault(method.getReturnType()); + } + }; + return (DatabaseMetaData) Proxy.newProxyInstance( + DataTaskJobEngineStreamQueryTest.class.getClassLoader(), + new Class[]{DatabaseMetaData.class}, handler); + } + + private static Object scalarDefault(Class type) { + if (type == boolean.class) { + return Boolean.FALSE; + } + if (type == int.class) { + return Integer.valueOf(0); + } + if (type == long.class) { + return Long.valueOf(0); + } + return null; + } + + private DataTaskJobEngine.StreamSpec spec(HikariDataSource dataSource, int fetchSize) { + return DataTaskJobEngine.StreamSpec.builder() + .dataSource(dataSource) + .product(ProductTypeEnum.POSTGRESQL) + .sqlMeta(new SqlMeta("SELECT c1 FROM t", Collections.emptyList())) + .naming(NamingStrategyEnum.NONE) + .cancelSupplier(new java.util.function.BooleanSupplier() { + @Override + public boolean getAsBoolean() { + return false; + } + }) + .rowLimit(10L) + .flushIntervalMs(5000L) + .fetchSize(fetchSize) + .timeoutSeconds(1800) + .progress(null) + .build(); + } + + @Test + public void rowsetRunsInLocalTransactionWithBoundedFetchAndRestoresAutoCommit() throws Exception { + Trace trace = new Trace(); + trace.rowset = true; + trace.remainingRows = 1; + trace.cellValue = "v1"; + FakeJdbc jdbc = new FakeJdbc(trace); + CaptureChannel channel = new CaptureChannel(); + DataTaskJobEngine engine = new DataTaskJobEngine(); + + DataTaskJobEngine.StreamResult result = engine.streamQuery(spec(jdbc.dataSource, 1000), channel); + + Assert.assertTrue(result.isRowset()); + Assert.assertEquals(1L, result.getRows()); + Assert.assertEquals(Collections.singletonList("c1"), channel.columns); + Assert.assertEquals(1, channel.rows.size()); + Assert.assertArrayEquals(new Object[]{"v1"}, channel.rows.get(0)); + + // auto-commit flipped off before the statement ran, so the PG-family driver + // must treat fetch size as a statement cursor bound instead of buffering all + Assert.assertFalse(trace.autoCommitAtExecute); + Assert.assertTrue(trace.executed); + Assert.assertEquals(Collections.singletonList(Integer.valueOf(1000)), trace.fetchSizes); + // reads never commit, the engine restores auto-commit before closing, and + // restore + close happen after the statement finished + Assert.assertFalse(trace.committed); + Assert.assertTrue(indexOf(trace.calls, "conn.setAutoCommit") + < indexOf(trace.calls, "conn.prepareStatement")); + Assert.assertTrue(lastIndexOf(trace.calls, "conn.setAutoCommit") > indexOf(trace.calls, "stmt.close")); + Assert.assertTrue(lastIndexOf(trace.calls, "conn.setAutoCommit") < indexOf(trace.calls, "conn.close")); + Assert.assertTrue(trace.closed); + } + + @Test + public void updateStatementCommitsBeforeRestoringAutoCommit() throws Exception { + Trace trace = new Trace(); + trace.rowset = false; + trace.updateCount = 7; + FakeJdbc jdbc = new FakeJdbc(trace); + CaptureChannel channel = new CaptureChannel(); + DataTaskJobEngine engine = new DataTaskJobEngine(); + + DataTaskJobEngine.StreamResult result = engine.streamQuery(spec(jdbc.dataSource, 1000), channel); + + Assert.assertFalse(result.isRowset()); + Assert.assertEquals(7, result.getUpdateCount()); + Assert.assertTrue(channel.rows.isEmpty()); + // DML now runs outside autocommit: rows must be committed, and the commit + // lands after execute but before restore + close + Assert.assertTrue(trace.committed); + Assert.assertTrue(indexOf(trace.calls, "conn.commit") > indexOf(trace.calls, "stmt.execute")); + Assert.assertTrue(indexOf(trace.calls, "conn.commit") < lastIndexOf(trace.calls, "conn.setAutoCommit")); + Assert.assertTrue(lastIndexOf(trace.calls, "conn.setAutoCommit") < indexOf(trace.calls, "conn.close")); + } + + @Test + public void autocommitAlreadyDisabledIsLeftAlone() throws Exception { + Trace trace = new Trace(); + trace.autoCommit = false; + trace.rowset = false; + trace.updateCount = 3; + FakeJdbc jdbc = new FakeJdbc(trace); + DataTaskJobEngine engine = new DataTaskJobEngine(); + + DataTaskJobEngine.StreamResult result = engine.streamQuery(spec(jdbc.dataSource, 1000), new CaptureChannel()); + + Assert.assertEquals(3, result.getUpdateCount()); + // the pool owns the autocommit state here: the engine must neither flip nor + // restore it, only commit its own single-statement transaction + Assert.assertFalse(contains(trace.calls, "conn.setAutoCommit")); + Assert.assertTrue(trace.committed); + Assert.assertTrue(trace.closed); + } + + private static int indexOf(List calls, String name) { + return calls.indexOf(name); + } + + private static int lastIndexOf(List calls, String name) { + return calls.lastIndexOf(name); + } + + private static boolean contains(List calls, String name) { + return calls.contains(name); + } +} \ No newline at end of file diff --git a/datapoly-dist/src/main/assembly/bin/datapolyctl.sh b/datapoly-dist/src/main/assembly/bin/datapolyctl.sh index 826bb41..62c7d5f 100644 --- a/datapoly-dist/src/main/assembly/bin/datapolyctl.sh +++ b/datapoly-dist/src/main/assembly/bin/datapolyctl.sh @@ -56,7 +56,10 @@ export DATAPOLY_MANAGER_URL=$(get_config_value "DATAPOLY_MANAGER_URL" "${APP_CON export DATAPOLY_GATEWAY_URL=$(get_config_value "DATAPOLY_GATEWAY_URL" "${APP_CONF_PATH}/config.ini") # JVM参数可以在这里设置 -JVMFLAGS="-server -Xms1024m -Xmx1024m -Xmn1024m -XX:+DisableExplicitGC -Djava.awt.headless=true -Dfile.encoding=UTF-8 " +# 堆 4G、年轻代/老年代 1:3:长驻对象(Hazelcast token/API 响应缓存、Eureka、Spring 框架 +# 对象)占堆内大头,老年代空间优先;年轻代 1G 足以容纳数据任务流式批次与 ≤200 行的 +# 调试/预览结果集(线上观测:16 分钟仅 7 次 minor GC,老年代却 30 秒内 5→25 次 major GC)。 +JVMFLAGS="-server -Xms4096m -Xmx4096m -Xmn1024m -XX:+DisableExplicitGC -Djava.awt.headless=true -Dfile.encoding=UTF-8 " if [ "$JAVA_HOME" != "" ]; then JAVA="$JAVA_HOME/bin/java" diff --git a/datapoly-dist/src/main/assembly/bin/executor_startup.cmd b/datapoly-dist/src/main/assembly/bin/executor_startup.cmd index 3b04a67..882f21a 100644 --- a/datapoly-dist/src/main/assembly/bin/executor_startup.cmd +++ b/datapoly-dist/src/main/assembly/bin/executor_startup.cmd @@ -24,8 +24,8 @@ for /f "delims=" %%i in ('type "%APP_HOME%\conf\config.ini"^| find /i "="') do s ::设置DEBUG端口 set DEBUG_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=n,address=18092 -::java虚拟机启动参数 -set JAVA_OPTS=-server -Xms4096m -Xmx4096m -Xmn2048m -XX:+DisableExplicitGC %DEBUG_OPTS% -Djava.awt.headless=true -Dfile.encoding=UTF-8 -Doracle.jdbc.J2EE13Compliant=true +::java虚拟机启动参数(堆 4G、年轻代/老年代 1:3,理由见 datapolyctl.sh 注释) +set JAVA_OPTS=-server -Xms4096m -Xmx4096m -Xmn1024m -XX:+DisableExplicitGC %DEBUG_OPTS% -Djava.awt.headless=true -Dfile.encoding=UTF-8 -Doracle.jdbc.J2EE13Compliant=true ::打印环境信息 echo System Information: diff --git a/datapoly-dist/src/main/assembly/bin/gateway_startup.cmd b/datapoly-dist/src/main/assembly/bin/gateway_startup.cmd index c8671c1..e75502f 100644 --- a/datapoly-dist/src/main/assembly/bin/gateway_startup.cmd +++ b/datapoly-dist/src/main/assembly/bin/gateway_startup.cmd @@ -24,8 +24,8 @@ for /f "delims=" %%i in ('type "%APP_HOME%\conf\config.ini"^| find /i "="') do s ::设置DEBUG端口 set DEBUG_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=n,address=18091 -::java虚拟机启动参数 -set JAVA_OPTS=-server -Xms4096m -Xmx4096m -Xmn2048m -XX:+DisableExplicitGC %DEBUG_OPTS% -Djava.awt.headless=true -Dfile.encoding=UTF-8 -Doracle.jdbc.J2EE13Compliant=true +::java虚拟机启动参数(堆 4G、年轻代/老年代 1:3,理由见 datapolyctl.sh 注释) +set JAVA_OPTS=-server -Xms4096m -Xmx4096m -Xmn1024m -XX:+DisableExplicitGC %DEBUG_OPTS% -Djava.awt.headless=true -Dfile.encoding=UTF-8 -Doracle.jdbc.J2EE13Compliant=true ::打印环境信息 echo System Information: diff --git a/datapoly-dist/src/main/assembly/bin/manager_startup.cmd b/datapoly-dist/src/main/assembly/bin/manager_startup.cmd index cfbb30d..f058253 100644 --- a/datapoly-dist/src/main/assembly/bin/manager_startup.cmd +++ b/datapoly-dist/src/main/assembly/bin/manager_startup.cmd @@ -24,8 +24,8 @@ for /f "delims=" %%i in ('type "%APP_HOME%\conf\config.ini"^| find /i "="') do s ::设置DEBUG端口 set DEBUG_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=n,address=18090 -::java虚拟机启动参数 -set JAVA_OPTS=-server -Xms4096m -Xmx4096m -Xmn2048m -XX:+DisableExplicitGC %DEBUG_OPTS% -Djava.awt.headless=true -Dfile.encoding=UTF-8 -Doracle.jdbc.J2EE13Compliant=true +::java虚拟机启动参数(堆 4G、年轻代/老年代 1:3,理由见 datapolyctl.sh 注释) +set JAVA_OPTS=-server -Xms4096m -Xmx4096m -Xmn1024m -XX:+DisableExplicitGC %DEBUG_OPTS% -Djava.awt.headless=true -Dfile.encoding=UTF-8 -Doracle.jdbc.J2EE13Compliant=true ::打印环境信息 echo System Information: diff --git a/docs/en/data-task.md b/docs/en/data-task.md index 039a39c..ae4fab6 100644 --- a/docs/en/data-task.md +++ b/docs/en/data-task.md @@ -289,7 +289,7 @@ variable in `application.yaml`): | `reap-interval-ms` | `30000` | lost-worker reap cadence | | `lease-seconds` | `600` | run lease length, renewed via heartbeats; expired leases are reaped to FAILED | | `flush-interval-ms` | `5000` | minimum interval between progress/cancel-check/lease-renewal ticks | -| `fetch-size` | `1000` | JDBC fetchSize (MySQL dialects switch to streaming `Integer.MIN_VALUE`) | +| `fetch-size` | `1000` | per-fetch row bound: the engine executes the single statement in an explicit transaction (DML is committed afterwards) so PostgreSQL-family drivers (incl. Hologres) honor it and pull rows in bounded batches instead of buffering the whole result set under autocommit; MySQL dialects switch to streaming `Integer.MIN_VALUE` | | `query-timeout-seconds` | `1800` | statement-level timeout backstop | | `max-rows-default` | `1000000` | global row cap when a definition sets no `maxRows` (or ≤ 0) | @@ -323,4 +323,5 @@ nodes; jobs of lost workers are marked FAILED once the lease expires and callers | `artifactInfo.truncated=true` | output hit the row cap (definition `maxRows` or `max-rows-default`); raise the cap or split the SQL into batches for full extracts | | Cancellation is slow to take effect | cancelling a RUNNING job is cooperative and lands within ≤ `flush-interval-ms` (5 s by default); PENDING jobs cancel immediately | | Preview works but the submitted job fails | preview never touches the sink — the failure is almost certainly on the delivery side (sink missing, invalid `sinkConfig`, target-side auth); check `errorMessage` and executor logs | +| FAILED with `Java heap space` / executor repeated `GC overhead limit exceeded` | first check executor heap vs host memory overcommit (the release script `datapolyctl.sh` defaults to a 4G heap per service — leave headroom when co-locating); the engine now keeps PostgreSQL-family drivers fetching in `fetch-size` batches, while the row cap (`maxRows`) and the sink's own memory behavior stay under your control | | Duplicate deliveries after deploying multiple executors? | Not possible: claiming is an atomic `FOR UPDATE SKIP LOCKED` operation in the meta store — a job is RUNNING on at most one node | diff --git a/docs/zh/data-task.md b/docs/zh/data-task.md index c7b214d..0ead322 100644 --- a/docs/zh/data-task.md +++ b/docs/zh/data-task.md @@ -271,7 +271,7 @@ executor 侧前缀 `datapoly.data-task.*`(`application.yaml` 均可用 `DATAPO | `reap-interval-ms` | `30000` | 失联任务回收检查周期 | | `lease-seconds` | `600` | 运行租约时长,随心跳续期;到期未续即被 reaper 记 FAILED | | `flush-interval-ms` | `5000` | 进度刷新/取消检测/租约续期的最小间隔 | -| `fetch-size` | `1000` | JDBC fetchSize(MySQL 方言自动改用流式 `Integer.MIN_VALUE`)| +| `fetch-size` | `1000` | JDBC 单次网络抓取行数上限:引擎把单条语句放进显式事务执行(DML 事后自动提交),PostgreSQL 系驱动(含 Hologres)才会真正按该值分批拉取,而不是 autocommit 下把整个结果集缓冲进堆;MySQL 方言自动改用流式 `Integer.MIN_VALUE`| | `query-timeout-seconds` | `1800` | 语句级超时兜底 | | `max-rows-default` | `1000000` | 定义未显式设置 `maxRows`(或 ≤0)时生效的全局行数上限 | @@ -303,4 +303,5 @@ executor 侧前缀 `datapoly.data-task.*`(`application.yaml` 均可用 `DATAPO | `artifactInfo.truncated=true` | 命中行数上限被截断(定义 `maxRows` 或 `max-rows-default`);需要全量就调大上限或改写 SQL 分批 | | 取消迟迟不生效 | RUNNING 任务的取消是协作式的,生效延迟 ≤ `flush-interval-ms`(默认 5 秒);PENDING 任务取消立即生效 | | preview 正常但正式任务失败 | preview 不触碰 sink——失败几乎必然在投递侧(sink 未部署、`sinkConfig` 不合法、目标端鉴权失败),看 `errorMessage` 与 executor 日志 | +| FAILED,`Java heap space` / executor 反复 `GC overhead limit exceeded` | 先核对 executor 堆与宿主内存是否超配(发行脚本 `datapolyctl.sh` 默认每服务 4G 堆,同机多服务需留余量);引擎已保证 PostgreSQL 系驱动按 `fetch-size` 分批拉取,行数上限(`maxRows`)与 sink 自身内存行为仍需控制 | | 多 executor 部署后任务重复投递? | 不会。认领是元库事务内 `FOR UPDATE SKIP LOCKED` 原子操作,一条任务只会被一个节点置为 RUNNING |