This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 58a885dcc18 branch-4.1: [fix](streaming-job) Optimize snapshot offset
persistence #66238 (#66539)
58a885dcc18 is described below
commit 58a885dcc18db3d43afcb6fb358bdefc81f454ab
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Tue Aug 11 14:34:25 2026 +0800
branch-4.1: [fix](streaming-job) Optimize snapshot offset persistence
#66238 (#66539)
Cherry-picked from #66238
Co-authored-by: wudi <[email protected]>
---
.../main/java/org/apache/doris/common/Config.java | 4 +
.../apache/doris/job/cdc/DataSourceConfigKeys.java | 1 +
.../insert/streaming/StreamingInsertJob.java | 30 ++-
.../streaming/StreamingJobSchedulerTask.java | 1 +
.../job/offset/jdbc/JdbcSourceOffsetProvider.java | 30 +++
.../offset/jdbc/JdbcTvfSourceOffsetProvider.java | 9 +-
.../StreamingInsertJobOffsetPersistenceTest.java | 263 +++++++++++++++++++++
.../jdbc/JdbcSourceOffsetProviderOffsetTest.java | 211 +++++++++++++++++
.../source/reader/mysql/MySqlSourceReader.java | 9 +-
.../reader/postgres/PostgresSourceReader.java | 10 +-
...g_mysql_job_snapshot_finished_restart_fe.groovy | 160 +++++++++++++
11 files changed, 708 insertions(+), 20 deletions(-)
diff --git a/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
b/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
index 5e757ac58ed..7f5038c1da5 100644
--- a/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
+++ b/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
@@ -1382,6 +1382,10 @@ public class Config extends ConfigBase {
@ConfField(mutable = true, masterOnly = true)
public static int streaming_task_min_timeout_sec = 300;
+ @ConfField(mutable = true, masterOnly = true, description = {
+ "Minimum interval in seconds between snapshot offset persistence
operations"})
+ public static int streaming_job_snapshot_offset_persist_interval_sec = 300;
+
@ConfField(mutable = true, masterOnly = true)
public static int streaming_cdc_light_rpc_timeout_sec = 90;
diff --git
a/fe/fe-common/src/main/java/org/apache/doris/job/cdc/DataSourceConfigKeys.java
b/fe/fe-common/src/main/java/org/apache/doris/job/cdc/DataSourceConfigKeys.java
index fb0f1825324..95956cbeb49 100644
---
a/fe/fe-common/src/main/java/org/apache/doris/job/cdc/DataSourceConfigKeys.java
+++
b/fe/fe-common/src/main/java/org/apache/doris/job/cdc/DataSourceConfigKeys.java
@@ -35,6 +35,7 @@ public class DataSourceConfigKeys {
public static final String OFFSET_LATEST = "latest";
public static final String OFFSET_SNAPSHOT = "snapshot";
public static final String SNAPSHOT_SPLIT_SIZE = "snapshot_split_size";
+ public static final String SNAPSHOT_SPLIT_SIZE_DEFAULT = "40960";
public static final String SNAPSHOT_SPLIT_KEY = "snapshot_split_key";
public static final String SNAPSHOT_PARALLELISM = "snapshot_parallelism";
public static final String SNAPSHOT_PARALLELISM_DEFAULT = "1";
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
index cc7e7e34a18..affb6d6c599 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
@@ -152,6 +152,7 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
@SerializedName("opp")
// The value to be persisted in offsetProvider
private String offsetProviderPersist;
+ private transient long lastOffsetPersistTimeMs;
@Setter
@Getter
private long lastScheduleTaskTimestamp = -1L;
@@ -899,6 +900,7 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
// offset provider has reached a natural end, mark job as
finished
log.info("Streaming insert job {} source data fully consumed,
marking job as FINISHED", getJobId());
updateJobStatus(JobStatus.FINISHED);
+ logUpdateOperation();
return;
}
AbstractStreamingTask nextTask = createStreamingTask();
@@ -1026,6 +1028,10 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
// insert TVF does not persist the running state.
// streaming multi task persists the running state when
commitOffset() is called.
setJobStatus(replayJob.getJobStatus());
+ if (isFinalStatus()) {
+ setFinishTimeMs(replayJob.getFinishTimeMs());
+
Env.getCurrentGlobalTransactionMgr().getCallbackFactory().removeCallback(getJobId());
+ }
}
try {
modifyPropertiesInternal(replayJob.getProperties());
@@ -1059,6 +1065,7 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
setFailedTaskCount(replayJob.getFailedTaskCount());
setCanceledTaskCount(replayJob.getCanceledTaskCount());
setLastTaskSuccessTime(replayJob.getLastTaskSuccessTime());
+ setStartTimeMs(replayJob.getStartTimeMs());
this.boundBackendId = replayJob.boundBackendId;
}
@@ -1548,6 +1555,7 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
throw new JobException("Unsupported commit offset for offset
provider type: "
+ offsetProvider.getClass().getSimpleName());
}
+ JdbcSourceOffsetProvider jdbcOffsetProvider =
(JdbcSourceOffsetProvider) offsetProvider;
writeLock();
try {
@@ -1575,10 +1583,9 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
updateNoTxnJobStatisticAndOffset(offsetRequest);
offsetProvider.onTaskCommitted(offsetRequest.getScannedRows(),
offsetRequest.getLoadBytes());
if (offsetRequest.getTableSchemas() != null) {
- JdbcSourceOffsetProvider op = (JdbcSourceOffsetProvider)
offsetProvider;
- op.setTableSchemas(offsetRequest.getTableSchemas());
+
jdbcOffsetProvider.setTableSchemas(offsetRequest.getTableSchemas());
}
- persistOffsetProviderIfNeed();
+ persistOffsetProviderIfNeed(jdbcOffsetProvider,
System.currentTimeMillis());
log.info("Streaming multi table job {} task {} commit offset
successfully, offset: {}",
getJobId(), offsetRequest.getTaskId(),
offsetRequest.getOffset());
((StreamingMultiTblTask)
this.runningStreamTask).successCallback(offsetRequest);
@@ -1638,12 +1645,19 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
}
}
- private void persistOffsetProviderIfNeed() {
- // only for jdbc
- this.offsetProviderPersist = offsetProvider.getPersistInfo();
- if (this.offsetProviderPersist != null) {
- logUpdateOperation();
+ private void persistOffsetProviderIfNeed(
+ JdbcSourceOffsetProvider jdbcOffsetProvider, long currentTimeMs) {
+ this.offsetProviderPersist = jdbcOffsetProvider.getPersistInfo();
+ if (this.offsetProviderPersist == null) {
+ return;
}
+
+ if (!jdbcOffsetProvider.shouldPersistOffset(lastOffsetPersistTimeMs,
currentTimeMs)) {
+ return;
+ }
+
+ logUpdateOperation();
+ lastOffsetPersistTimeMs = currentTimeMs;
}
public void replayOffsetProviderIfNeed() throws JobException {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobSchedulerTask.java
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobSchedulerTask.java
index 0f6bcba892b..030c3bb2b8b 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobSchedulerTask.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobSchedulerTask.java
@@ -76,6 +76,7 @@ public class StreamingJobSchedulerTask extends AbstractTask {
// Source already fully consumed (e.g. snapshot-only mode
recovered after FE restart).
// Transition directly to FINISHED without creating a new task.
streamingInsertJob.updateJobStatus(JobStatus.FINISHED);
+ streamingInsertJob.logUpdateOperation();
return;
}
streamingInsertJob.createStreamingTask();
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProvider.java
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProvider.java
index 13874a35101..4eb5fad0890 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProvider.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProvider.java
@@ -256,8 +256,13 @@ public class JdbcSourceOffsetProvider implements
SourceOffsetProvider {
} else {
synchronized (splitsLock) {
BinlogSplit binlogSplit = (BinlogSplit)
newOffset.getSplits().get(0);
+ if (MapUtils.isEmpty(binlogSplit.getStartingOffset())) {
+ log.warn("Skip empty committed binlog offset for job {}",
getJobId());
+ return;
+ }
binlogOffsetPersist = new
HashMap<>(binlogSplit.getStartingOffset());
binlogOffsetPersist.put(SPLIT_ID, BinlogSplit.BINLOG_SPLIT_ID);
+ clearSnapshotState();
currentOffset = newOffset;
hasMoreData = true;
}
@@ -266,6 +271,31 @@ public class JdbcSourceOffsetProvider implements
SourceOffsetProvider {
this.currentOffset = newOffset;
}
+ protected void clearSnapshotState() {
+ if (MapUtils.isNotEmpty(chunkHighWatermarkMap)) {
+ chunkHighWatermarkMap = new HashMap<>();
+ }
+ remainingSplits.clear();
+ finishedSplits.clear();
+ if (committedSplitProgress != null) {
+ clearProgress(committedSplitProgress);
+ }
+ if (cdcSplitProgress != null) {
+ clearProgress(cdcSplitProgress);
+ }
+ }
+
+ public boolean shouldPersistOffset(long lastPersistTimeMs, long
currentTimeMs) {
+ synchronized (splitsLock) {
+ if (currentOffset == null || !currentOffset.snapshotSplit()) {
+ return true;
+ }
+ }
+ long intervalMs = Math.max(1L,
+ (long)
Config.streaming_job_snapshot_offset_persist_interval_sec) * 1000L;
+ return lastPersistTimeMs == 0L || currentTimeMs - lastPersistTimeMs >=
intervalMs;
+ }
+
@Override
public void setBoundBackendId(long boundBackendId) {
this.boundBackendId = boundBackendId;
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcTvfSourceOffsetProvider.java
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcTvfSourceOffsetProvider.java
index 0e5bb8fb753..e7324015d93 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcTvfSourceOffsetProvider.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcTvfSourceOffsetProvider.java
@@ -294,10 +294,13 @@ public class JdbcTvfSourceOffsetProvider extends
JdbcSourceOffsetProvider {
synchronized (splitsLock) {
// Mirror binlog offset into bop so it survives FE checkpoint
BinlogSplit bs = (BinlogSplit) newOffset.getSplits().get(0);
- if (MapUtils.isNotEmpty(bs.getStartingOffset())) {
- binlogOffsetPersist = new
HashMap<>(bs.getStartingOffset());
- binlogOffsetPersist.put(SPLIT_ID,
BinlogSplit.BINLOG_SPLIT_ID);
+ if (MapUtils.isEmpty(bs.getStartingOffset())) {
+ log.warn("Skip empty committed binlog offset for job {}",
getJobId());
+ return;
}
+ binlogOffsetPersist = new HashMap<>(bs.getStartingOffset());
+ binlogOffsetPersist.put(SPLIT_ID, BinlogSplit.BINLOG_SPLIT_ID);
+ clearSnapshotState();
currentOffset = newOffset;
hasMoreData = true;
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobOffsetPersistenceTest.java
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobOffsetPersistenceTest.java
new file mode 100644
index 00000000000..0b13d0ec745
--- /dev/null
+++
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobOffsetPersistenceTest.java
@@ -0,0 +1,263 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.job.extensions.insert.streaming;
+
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.job.cdc.request.CommitOffsetRequest;
+import org.apache.doris.job.cdc.split.SnapshotSplit;
+import org.apache.doris.job.common.JobStatus;
+import org.apache.doris.job.common.TaskStatus;
+import org.apache.doris.job.exception.JobException;
+import org.apache.doris.job.manager.JobManager;
+import org.apache.doris.job.manager.StreamingTaskManager;
+import org.apache.doris.job.offset.jdbc.JdbcSourceOffsetProvider;
+import org.apache.doris.transaction.GlobalTransactionMgrIface;
+import org.apache.doris.transaction.TxnStateCallbackFactory;
+
+import org.junit.Assert;
+import org.junit.Test;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.concurrent.locks.ReentrantReadWriteLock;
+
+public class StreamingInsertJobOffsetPersistenceTest {
+
+ @Test
+ public void testFirstSnapshotCommitPersistsImmediately() throws Exception {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ provider.getRemainingSplits().add(snapshotSplit("source_table:0"));
+ TestStreamingInsertJob job = newJob(provider, 1001L);
+
+ job.commitOffset(snapshotRequest(1001L, "source_table:0", null));
+
+ Assert.assertEquals(1, job.journalCount);
+ Assert.assertNotNull(job.getOffsetProviderPersist());
+ }
+
+ @Test
+ public void testSnapshotCommitWithinIntervalDoesNotPersistAgain() throws
Exception {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ provider.getRemainingSplits().add(snapshotSplit("source_table:0"));
+ TestStreamingInsertJob job = newJob(provider, 1003L);
+ job.commitOffset(snapshotRequest(1003L, "source_table:0", null));
+
+ provider.getRemainingSplits().add(snapshotSplit("source_table:1"));
+ job.commitOffset(snapshotRequest(1003L, "source_table:1", null));
+
+ Assert.assertEquals(1, job.journalCount);
+ Assert.assertNotNull(job.getOffsetProviderPersist());
+ }
+
+ @Test
+ public void testBinlogCommitPersistsImmediately() throws Exception {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ TestStreamingInsertJob job = newJob(provider, 1002L);
+
+ job.commitOffset(binlogRequest(1002L, "100"));
+ job.commitOffset(binlogRequest(1002L, "200"));
+
+ Assert.assertEquals(2, job.journalCount);
+ Assert.assertNotNull(job.getOffsetProviderPersist());
+ }
+
+ @Test
+ public void testSnapshotToBinlogTransitionPersistsCompactedState() throws
Exception {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ provider.getRemainingSplits().add(snapshotSplit("source_table:0"));
+ TestStreamingInsertJob job = newJob(provider, 1008L);
+ job.commitOffset(snapshotRequest(1008L, "source_table:0", null));
+ Assert.assertEquals(1, job.journalCount);
+
+ job.commitOffset(binlogRequest(1008L, "200"));
+
+ Assert.assertEquals(2, job.journalCount);
+
Assert.assertFalse(job.getOffsetProviderPersist().contains("source_table:0"));
+ Assert.assertTrue(provider.getFinishedSplits().isEmpty());
+ Assert.assertTrue(provider.getChunkHighWatermarkMap().isEmpty());
+ }
+
+ @Test
+ public void testSnapshotOffsetPersistsOnNextCommitAfterInterval() throws
Exception {
+ int oldInterval =
Config.streaming_job_snapshot_offset_persist_interval_sec;
+ Config.streaming_job_snapshot_offset_persist_interval_sec = 300;
+ try {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ provider.getRemainingSplits().add(snapshotSplit("source_table:0"));
+ TestStreamingInsertJob job = newJob(provider, 1011L);
+
+ job.commitOffset(snapshotRequest(1011L, "source_table:0", null));
+ Assert.assertEquals(1, job.journalCount);
+ Deencapsulation.setField(job, "lastOffsetPersistTimeMs",
+ System.currentTimeMillis() - 300_000L);
+ provider.getRemainingSplits().add(snapshotSplit("source_table:1"));
+ job.commitOffset(snapshotRequest(1011L, "source_table:1", null));
+
+ Assert.assertEquals(2, job.journalCount);
+ Assert.assertTrue((long) Deencapsulation.getField(job,
"lastOffsetPersistTimeMs") > 0L);
+ } finally {
+ Config.streaming_job_snapshot_offset_persist_interval_sec =
oldInterval;
+ }
+ }
+
+ @Test
+ public void testAlterOffsetReplacesSnapshotState() throws Exception {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ provider.getRemainingSplits().add(snapshotSplit("source_table:0"));
+ TestStreamingInsertJob job = newJob(provider, 1009L);
+ job.commitOffset(snapshotRequest(1009L, "source_table:0", null));
+
+ HashMap<String, String> properties = new HashMap<>();
+ properties.put(StreamingJobProperties.OFFSET_PROPERTY,
"{\"lsn\":\"300\"}");
+ Deencapsulation.invoke(job, "modifyPropertiesInternal", properties);
+
+ Assert.assertTrue(job.getOffsetProviderPersist().contains("300"));
+ Assert.assertTrue(provider.getFinishedSplits().isEmpty());
+ Assert.assertTrue(provider.getChunkHighWatermarkMap().isEmpty());
+ }
+
+ @Test
+ public void testNaturalFinishPersistsFinalState() throws Exception {
+ TestStreamingInsertJob job = newJob(new EndJdbcSourceOffsetProvider(),
1012L);
+ NoopStreamingMultiTblTask task =
+ (NoopStreamingMultiTblTask) Deencapsulation.getField(job,
"runningStreamTask");
+
+ try (MockedStatic<Env> envMockedStatic =
Mockito.mockStatic(Env.class)) {
+ Env env = Mockito.mock(Env.class);
+ JobManager<?, ?> jobManager = Mockito.mock(JobManager.class);
+ StreamingTaskManager streamingTaskManager =
Mockito.mock(StreamingTaskManager.class);
+ GlobalTransactionMgrIface transactionMgr =
Mockito.mock(GlobalTransactionMgrIface.class);
+ TxnStateCallbackFactory callbackFactory =
Mockito.mock(TxnStateCallbackFactory.class);
+ envMockedStatic.when(Env::getCurrentEnv).thenReturn(env);
+
envMockedStatic.when(Env::getCurrentGlobalTransactionMgr).thenReturn(transactionMgr);
+ Mockito.when(env.getJobManager()).thenReturn(jobManager);
+
Mockito.when(jobManager.getStreamingTaskManager()).thenReturn(streamingTaskManager);
+
Mockito.when(transactionMgr.getCallbackFactory()).thenReturn(callbackFactory);
+
+ long beforeFinish = System.currentTimeMillis();
+ job.onStreamTaskSuccess(task);
+
+ Assert.assertEquals(JobStatus.FINISHED, job.getJobStatus());
+ Assert.assertTrue(job.getFinishTimeMs() >= beforeFinish);
+ Assert.assertEquals(1, job.journalCount);
+ Mockito.verify(callbackFactory).removeCallback(9001L);
+ }
+ }
+
+ @Test
+ public void testReplayUpdatedRestoresFinalStateAndRemovesCallback() {
+ TestStreamingInsertJob job = newJob(new JdbcSourceOffsetProvider(),
1013L);
+ TestStreamingInsertJob replayJob = newJob(new
JdbcSourceOffsetProvider(), 1014L);
+ replayJob.setJobStatus(JobStatus.FINISHED);
+ replayJob.setFinishTimeMs(1234L);
+
+ try (MockedStatic<Env> envMockedStatic =
Mockito.mockStatic(Env.class)) {
+ GlobalTransactionMgrIface transactionMgr =
Mockito.mock(GlobalTransactionMgrIface.class);
+ TxnStateCallbackFactory callbackFactory =
Mockito.mock(TxnStateCallbackFactory.class);
+
envMockedStatic.when(Env::getCurrentGlobalTransactionMgr).thenReturn(transactionMgr);
+
Mockito.when(transactionMgr.getCallbackFactory()).thenReturn(callbackFactory);
+
+ job.replayOnUpdated(replayJob);
+
+ Assert.assertEquals(JobStatus.FINISHED, job.getJobStatus());
+ Assert.assertEquals(1234L, job.getFinishTimeMs());
+ Mockito.verify(callbackFactory).removeCallback(9001L);
+ }
+ }
+
+ @Test
+ public void testReplayUpdatedRestoresStartTime() {
+ TestStreamingInsertJob job = newJob(new JdbcSourceOffsetProvider(),
1015L);
+ TestStreamingInsertJob replayJob = newJob(new
JdbcSourceOffsetProvider(), 1016L);
+ replayJob.setStartTimeMs(1234L);
+
+ job.replayOnUpdated(replayJob);
+
+ Assert.assertEquals(1234L, job.getStartTimeMs());
+ }
+
+ private static TestStreamingInsertJob newJob(JdbcSourceOffsetProvider
provider, long taskId) {
+ TestStreamingInsertJob job = new TestStreamingInsertJob();
+ Deencapsulation.setField(job, "lock", new
ReentrantReadWriteLock(true));
+ Deencapsulation.setField(job, "jobId", 9001L);
+ Deencapsulation.setField(job, "jobName", "test_job");
+ Deencapsulation.setField(job, "jobStatus", JobStatus.RUNNING);
+ Deencapsulation.setField(job, "offsetProvider", provider);
+ Deencapsulation.setField(job, "properties", new HashMap<String,
String>());
+ Deencapsulation.setField(job, "targetProperties", new HashMap<String,
String>());
+ Deencapsulation.setField(job, "runningStreamTask", new
NoopStreamingMultiTblTask(taskId));
+ return job;
+ }
+
+ private static SnapshotSplit snapshotSplit(String splitId) {
+ return new SnapshotSplit(
+ splitId,
+ "source_db.source_table",
+ Collections.singletonList("id"),
+ new Object[]{1L},
+ new Object[]{2L},
+ null);
+ }
+
+ private static CommitOffsetRequest snapshotRequest(long taskId, String
splitId, String tableSchemas) {
+ CommitOffsetRequest request = new CommitOffsetRequest();
+ request.setTaskId(taskId);
+ request.setOffset("[{\"splitId\":\"" + splitId +
"\",\"lsn\":\"100\"}]");
+ request.setTableSchemas(tableSchemas);
+ return request;
+ }
+
+ private static CommitOffsetRequest binlogRequest(long taskId, String lsn) {
+ CommitOffsetRequest request = new CommitOffsetRequest();
+ request.setTaskId(taskId);
+ request.setOffset("[{\"splitId\":\"binlog-split\",\"lsn\":\"" + lsn +
"\"}]");
+ return request;
+ }
+
+ private static class TestStreamingInsertJob extends StreamingInsertJob {
+ private int journalCount;
+
+ @Override
+ public void logUpdateOperation() {
+ journalCount++;
+ }
+ }
+
+ private static class EndJdbcSourceOffsetProvider extends
JdbcSourceOffsetProvider {
+ @Override
+ public boolean hasReachedEnd() {
+ return true;
+ }
+ }
+
+ private static class NoopStreamingMultiTblTask extends
StreamingMultiTblTask {
+ NoopStreamingMultiTblTask(long taskId) {
+ super(9001L, taskId, null, null, null, null, null,
+ new StreamingJobProperties(new HashMap<>()), null, null);
+ Deencapsulation.setField(this, "status", TaskStatus.RUNNING);
+ }
+
+ @Override
+ public void successCallback(CommitOffsetRequest offsetRequest) throws
JobException {
+ }
+ }
+}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderOffsetTest.java
b/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderOffsetTest.java
index 6efb9959748..d71416fd3c5 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderOffsetTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderOffsetTest.java
@@ -17,16 +17,47 @@
package org.apache.doris.job.offset.jdbc;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.jmockit.Deencapsulation;
import org.apache.doris.job.cdc.split.BinlogSplit;
+import org.apache.doris.job.cdc.split.SnapshotSplit;
+import org.apache.doris.job.extensions.insert.streaming.StreamingInsertJob;
import org.junit.Assert;
import org.junit.Test;
import java.util.Collections;
+import java.util.HashMap;
import java.util.Map;
public class JdbcSourceOffsetProviderOffsetTest {
+ @Test
+ public void testSnapshotOffsetUsesConfiguredPersistInterval() {
+ int oldInterval =
Config.streaming_job_snapshot_offset_persist_interval_sec;
+ try {
+ Config.streaming_job_snapshot_offset_persist_interval_sec = 123;
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ provider.currentOffset = new JdbcOffset(
+
Collections.singletonList(snapshotSplit("source_table:0")));
+
+ Assert.assertTrue(provider.shouldPersistOffset(0L, 1_000L));
+ Assert.assertFalse(provider.shouldPersistOffset(1_000L, 123_999L));
+ Assert.assertTrue(provider.shouldPersistOffset(1_000L, 124_000L));
+ } finally {
+ Config.streaming_job_snapshot_offset_persist_interval_sec =
oldInterval;
+ }
+ }
+
+ @Test
+ public void testBinlogOffsetPersistsImmediately() {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ provider.currentOffset = new JdbcOffset(Collections.singletonList(
+ new BinlogSplit(Collections.singletonMap("lsn", "100"))));
+
+ Assert.assertTrue(provider.shouldPersistOffset(1_000L, 1_001L));
+ }
+
@Test
public void testEndOffsetAdvancesWhenCurrentOffsetIsAhead() {
assertEndOffsetAdvancesWhenCurrentOffsetIsAhead(new
TestJdbcSourceOffsetProvider(-1));
@@ -76,6 +107,71 @@ public class JdbcSourceOffsetProviderOffsetTest {
((BinlogSplit)
provider.currentOffset.getSplits().get(0)).getStartingOffset());
}
+ @Test
+ public void testValidBinlogOffsetClearsSnapshotState() {
+ assertValidBinlogOffsetClearsSnapshotState(new
TestJdbcSourceOffsetProvider(-1));
+ }
+
+ @Test
+ public void testTvfValidBinlogOffsetClearsSnapshotState() {
+ assertValidBinlogOffsetClearsSnapshotState(new
TestJdbcTvfSourceOffsetProvider(-1));
+ }
+
+ @Test
+ public void testEmptyBinlogOffsetKeepsPreviousState() {
+ assertEmptyBinlogOffsetKeepsPreviousState(new
TestJdbcSourceOffsetProvider(-1));
+ }
+
+ @Test
+ public void testTvfEmptyBinlogOffsetKeepsPreviousState() {
+ assertEmptyBinlogOffsetKeepsPreviousState(new
TestJdbcTvfSourceOffsetProvider(-1));
+ }
+
+ @Test
+ public void testRepeatedValidBinlogOffsetCleanupIsIdempotent() {
+ assertRepeatedValidBinlogOffsetCleanupIsIdempotent(new
TestJdbcSourceOffsetProvider(-1));
+ }
+
+ @Test
+ public void testTvfRepeatedValidBinlogOffsetCleanupIsIdempotent() {
+ assertRepeatedValidBinlogOffsetCleanupIsIdempotent(new
TestJdbcTvfSourceOffsetProvider(-1));
+ }
+
+ @Test
+ public void testBinlogOffsetRestoredFromPersistInfo() throws Exception {
+ JdbcSourceOffsetProvider source = new TestJdbcSourceOffsetProvider(-1);
+ source.updateOffset(new JdbcOffset(Collections.singletonList(
+ new BinlogSplit(Collections.singletonMap("lsn", "200")))));
+ StreamingInsertJob job =
mockJobWithPersistInfo(source.getPersistInfo());
+ JdbcSourceOffsetProvider restored = new JdbcSourceOffsetProvider();
+
+ restored.replayIfNeed(job);
+
+ Assert.assertNotNull(restored.currentOffset);
+ Assert.assertFalse(restored.currentOffset.snapshotSplit());
+ Assert.assertEquals("200", ((BinlogSplit)
restored.currentOffset.getSplits().get(0))
+ .getStartingOffset().get("lsn"));
+ Assert.assertTrue(restored.chunkHighWatermarkMap.isEmpty());
+ }
+
+ @Test
+ public void testTvfBinlogOffsetRestoredFromPersistInfo() throws Exception {
+ JdbcSourceOffsetProvider source = new
TestJdbcTvfSourceOffsetProvider(-1);
+ source.updateOffset(new JdbcOffset(Collections.singletonList(
+ new BinlogSplit(Collections.singletonMap("lsn", "200")))));
+ StreamingInsertJob job =
mockJobWithPersistInfo(source.getPersistInfo());
+ JdbcTvfSourceOffsetProvider restored = new
JdbcTvfSourceOffsetProvider();
+
+ restored.restoreFromPersistInfo(source.getPersistInfo());
+ restored.replayIfNeed(job);
+
+ Assert.assertNotNull(restored.currentOffset);
+ Assert.assertFalse(restored.currentOffset.snapshotSplit());
+ Assert.assertEquals("200", ((BinlogSplit)
restored.currentOffset.getSplits().get(0))
+ .getStartingOffset().get("lsn"));
+ Assert.assertTrue(restored.chunkHighWatermarkMap.isEmpty());
+ }
+
private static void
assertEndOffsetAdvancesWhenCurrentOffsetIsAhead(JdbcSourceOffsetProvider
provider) {
Map<String, String> staleEndOffset = Collections.singletonMap("lsn",
"100");
Map<String, String> committedOffset = Collections.singletonMap("lsn",
"200");
@@ -91,6 +187,115 @@ public class JdbcSourceOffsetProviderOffsetTest {
Assert.assertEquals("{\"lsn\":\"200\"}", provider.getShowMaxOffset());
}
+ private static void
assertValidBinlogOffsetClearsSnapshotState(JdbcSourceOffsetProvider provider) {
+ seedSnapshotState(provider);
+ Map<String, String> binlogOffset = Collections.singletonMap("lsn",
"200");
+
+ provider.updateOffset(new JdbcOffset(
+ Collections.singletonList(new BinlogSplit(binlogOffset))));
+
+ Assert.assertTrue(provider.chunkHighWatermarkMap.isEmpty());
+ Assert.assertTrue(provider.remainingSplits.isEmpty());
+ Assert.assertTrue(provider.finishedSplits.isEmpty());
+ assertProgressCleared(provider.committedSplitProgress);
+ assertProgressCleared(provider.cdcSplitProgress);
+ Assert.assertEquals("table-schemas", provider.tableSchemas);
+ Map<String, String> expectedPersist = new HashMap<>(binlogOffset);
+ expectedPersist.put(JdbcSourceOffsetProvider.SPLIT_ID,
BinlogSplit.BINLOG_SPLIT_ID);
+ Assert.assertEquals(expectedPersist, provider.binlogOffsetPersist);
+ String persistInfo = provider.getPersistInfo();
+ Assert.assertFalse(persistInfo.contains("source_table:0"));
+ Assert.assertFalse(persistInfo.contains("source_table:1"));
+ Assert.assertTrue(persistInfo.contains("table-schemas"));
+ }
+
+ private static void
assertEmptyBinlogOffsetKeepsPreviousState(JdbcSourceOffsetProvider provider) {
+ seedSnapshotState(provider);
+ JdbcOffset previousOffset = new JdbcOffset(
+ Collections.singletonList(snapshotSplit("source_table:0")));
+ provider.currentOffset = previousOffset;
+ provider.hasMoreData = false;
+
+ provider.updateOffset(new JdbcOffset(
+ Collections.singletonList(new
BinlogSplit(Collections.emptyMap()))));
+
+ Assert.assertSame(previousOffset, provider.currentOffset);
+ Assert.assertFalse(provider.hasMoreData);
+ Assert.assertFalse(provider.chunkHighWatermarkMap.isEmpty());
+ Assert.assertFalse(provider.remainingSplits.isEmpty());
+ Assert.assertFalse(provider.finishedSplits.isEmpty());
+ Assert.assertEquals("source_table",
provider.committedSplitProgress.getCurrentSplittingTable());
+ Assert.assertEquals("source_table",
provider.cdcSplitProgress.getCurrentSplittingTable());
+ Assert.assertNull(provider.binlogOffsetPersist);
+ }
+
+ private static void assertRepeatedValidBinlogOffsetCleanupIsIdempotent(
+ JdbcSourceOffsetProvider provider) {
+ seedSnapshotState(provider);
+ JdbcOffset binlogOffset = new JdbcOffset(Collections.singletonList(
+ new BinlogSplit(Collections.singletonMap("lsn", "200"))));
+
+ provider.updateOffset(binlogOffset);
+ String firstPersistInfo = provider.getPersistInfo();
+ Map<String, Map<String, Map<String, String>>> clearedHighWatermarkMap =
+ provider.chunkHighWatermarkMap;
+ provider.updateOffset(binlogOffset);
+
+ Assert.assertEquals(firstPersistInfo, provider.getPersistInfo());
+ Assert.assertSame(clearedHighWatermarkMap,
provider.chunkHighWatermarkMap);
+ Assert.assertTrue(provider.chunkHighWatermarkMap.isEmpty());
+ Assert.assertTrue(provider.remainingSplits.isEmpty());
+ Assert.assertTrue(provider.finishedSplits.isEmpty());
+ assertProgressCleared(provider.committedSplitProgress);
+ assertProgressCleared(provider.cdcSplitProgress);
+ }
+
+ private static StreamingInsertJob mockJobWithPersistInfo(String
persistInfo) {
+ StreamingInsertJob job = new ReplayStreamingInsertJob();
+ Deencapsulation.setField(job, "jobId", 9001L);
+ Deencapsulation.setField(job, "syncTables", Collections.emptyList());
+ job.setOffsetProviderPersist(persistInfo);
+ return job;
+ }
+
+ private static void seedSnapshotState(JdbcSourceOffsetProvider provider) {
+ SnapshotSplit remaining = snapshotSplit("source_table:1");
+ SnapshotSplit finished = snapshotSplit("source_table:0");
+ provider.remainingSplits.add(remaining);
+ provider.finishedSplits.add(finished);
+ provider.chunkHighWatermarkMap
+ .computeIfAbsent("source_db.source_table", key -> new
HashMap<>())
+ .put(finished.getSplitId(), finished.getHighWatermark());
+ provider.committedSplitProgress = splitProgress();
+ provider.cdcSplitProgress = splitProgress();
+ provider.tableSchemas = "table-schemas";
+ }
+
+ private static SnapshotSplit snapshotSplit(String splitId) {
+ return new SnapshotSplit(
+ splitId,
+ "source_db.source_table",
+ Collections.singletonList("id"),
+ new Object[]{1L},
+ new Object[]{2L},
+ Collections.singletonMap("lsn", "100"));
+ }
+
+ private static JdbcSourceOffsetProvider.SplitProgress splitProgress() {
+ JdbcSourceOffsetProvider.SplitProgress progress = new
JdbcSourceOffsetProvider.SplitProgress();
+ progress.setCurrentSplittingTable("source_table");
+ progress.setNextSplitStart(new Object[]{2L});
+ progress.setNextSplitId(2);
+ return progress;
+ }
+
+ private static void
assertProgressCleared(JdbcSourceOffsetProvider.SplitProgress progress) {
+ Assert.assertNotNull(progress);
+ Assert.assertNull(progress.getCurrentSplittingTable());
+ Assert.assertNull(progress.getNextSplitStart());
+ Assert.assertNull(progress.getNextSplitId());
+ }
+
private static class TestJdbcSourceOffsetProvider extends
JdbcSourceOffsetProvider {
private final int compareResult;
@@ -104,6 +309,12 @@ public class JdbcSourceOffsetProviderOffsetTest {
}
}
+ private static class ReplayStreamingInsertJob extends StreamingInsertJob {
+ ReplayStreamingInsertJob() {
+ super();
+ }
+ }
+
private static class TestJdbcTvfSourceOffsetProvider extends
JdbcTvfSourceOffsetProvider {
private final int compareResult;
diff --git
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java
index 11075ea2d81..0ad2629ce94 100644
---
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java
+++
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java
@@ -1039,10 +1039,11 @@ public class MySqlSourceReader extends
AbstractCdcSourceReader {
configFactory.debeziumProperties(dbzProps);
configFactory.heartbeatInterval(Duration.ofMillis(DEBEZIUM_HEARTBEAT_INTERVAL_MS));
- if (cdcConfig.containsKey(DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE)) {
- configFactory.splitSize(
-
Integer.parseInt(cdcConfig.get(DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE)));
- }
+ configFactory.splitSize(
+ Integer.parseInt(
+ cdcConfig.getOrDefault(
+ DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE,
+
DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE_DEFAULT)));
// todo: Currently, only one split key is supported; future will
require multiple split
// keys.
diff --git
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/postgres/PostgresSourceReader.java
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/postgres/PostgresSourceReader.java
index e9cc9e8d985..330f461510b 100644
---
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/postgres/PostgresSourceReader.java
+++
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/postgres/PostgresSourceReader.java
@@ -294,11 +294,11 @@ public class PostgresSourceReader extends
JdbcIncrementalSourceReader {
throw new RuntimeException("Unknown offset " + startupMode);
}
- // Set split size if provided
- if (cdcConfig.containsKey(DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE)) {
- configFactory.splitSize(
-
Integer.parseInt(cdcConfig.get(DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE)));
- }
+ configFactory.splitSize(
+ Integer.parseInt(
+ cdcConfig.getOrDefault(
+ DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE,
+
DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE_DEFAULT)));
if (cdcConfig.containsKey(DataSourceConfigKeys.SNAPSHOT_SPLIT_KEY)) {
configFactory.chunkKeyColumn(cdcConfig.get(DataSourceConfigKeys.SNAPSHOT_SPLIT_KEY));
diff --git
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_snapshot_finished_restart_fe.groovy
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_snapshot_finished_restart_fe.groovy
new file mode 100644
index 00000000000..6bf0a57eb82
--- /dev/null
+++
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_snapshot_finished_restart_fe.groovy
@@ -0,0 +1,160 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+import org.apache.doris.regression.suite.ClusterOptions
+import org.awaitility.Awaitility
+
+import static java.util.concurrent.TimeUnit.SECONDS
+
+suite("test_streaming_mysql_job_snapshot_finished_restart_fe",
+ "docker,mysql,external_docker,external_docker_mysql,nondatalake") {
+ def jobName = "test_streaming_mysql_job_snapshot_finished_restart_fe"
+ def tableName = "snapshot_finished_restart_fe"
+ def mysqlDb = "test_cdc_db"
+ def totalRows = 5
+ def options = new ClusterOptions()
+ options.setFeNum(1)
+ options.cloudMode = null
+
+ docker(options) {
+ def currentDb = (sql "select database()")[0][0]
+
+ sql """DROP JOB IF EXISTS where jobname = '${jobName}'"""
+ sql """DROP TABLE IF EXISTS ${currentDb}.${tableName} FORCE"""
+
+ String enabled = context.config.otherConfigs.get("enableJdbcTest")
+ if (enabled != null && enabled.equalsIgnoreCase("true")) {
+ String mysqlPort = context.config.otherConfigs.get("mysql_57_port")
+ String externalEnvIp =
context.config.otherConfigs.get("externalEnvIp")
+ String s3Endpoint = getS3Endpoint()
+ String bucket = getS3BucketName()
+ String driverUrl =
+
"https://${bucket}.${s3Endpoint}/regression/jdbc_driver/mysql-connector-j-8.4.0.jar"
+
+ connect("root", "123456",
"jdbc:mysql://${externalEnvIp}:${mysqlPort}") {
+ sql """CREATE DATABASE IF NOT EXISTS ${mysqlDb}"""
+ sql """DROP TABLE IF EXISTS ${mysqlDb}.${tableName}"""
+ sql """CREATE TABLE ${mysqlDb}.${tableName} (
+ `id` int NOT NULL,
+ `name` varchar(200),
+ PRIMARY KEY (`id`)
+ ) ENGINE=InnoDB"""
+ sql """INSERT INTO ${mysqlDb}.${tableName} (id, name) VALUES
+ (1, 'name_1'),
+ (2, 'name_2'),
+ (3, 'name_3'),
+ (4, 'name_4'),
+ (5, 'name_5')"""
+ }
+
+ sql """CREATE JOB ${jobName}
+ ON STREAMING
+ FROM MYSQL (
+ "jdbc_url" =
"jdbc:mysql://${externalEnvIp}:${mysqlPort}",
+ "driver_url" = "${driverUrl}",
+ "driver_class" = "com.mysql.cj.jdbc.Driver",
+ "user" = "root",
+ "password" = "123456",
+ "database" = "${mysqlDb}",
+ "include_tables" = "${tableName}",
+ "offset" = "snapshot",
+ "snapshot_split_size" = "1",
+ "snapshot_parallelism" = "1"
+ )
+ TO DATABASE ${currentDb} (
+ "table.create.properties.replication_num" = "1"
+ )
+ """
+
+ try {
+ Awaitility.await().atMost(300, SECONDS)
+ .pollInterval(2, SECONDS).until(
+ {
+ def jobStatus = sql """
+ SELECT Status
+ FROM jobs("type"="insert")
+ WHERE Name='${jobName}' AND
ExecuteType='STREAMING'
+ """
+ log.info("jobStatus before FE restart: " +
jobStatus)
+ jobStatus.size() == 1 && jobStatus.get(0).get(0)
== "FINISHED"
+ }
+ )
+
+ def jobIdRows = sql """
+ SELECT Id
+ FROM jobs("type"="insert")
+ WHERE Name='${jobName}' AND ExecuteType='STREAMING'
+ """
+ assert jobIdRows.size() == 1
+ def jobId = jobIdRows.get(0).get(0).toString()
+
+ def rowsBeforeRestart = sql """
+ SELECT COUNT(*), COUNT(DISTINCT id)
+ FROM ${currentDb}.${tableName}
+ """
+ assert rowsBeforeRestart.size() == 1
+ assert rowsBeforeRestart.get(0).get(0) == totalRows
+ assert rowsBeforeRestart.get(0).get(1) == totalRows
+
+ cluster.restartFrontends()
+ sleep(60000)
+ context.reconnectFe()
+
+ // A terminal job must survive replay with its finish time.
Otherwise
+ // the scheduler may treat it as an expired job and remove it.
+ Awaitility.await().atMost(120, SECONDS)
+ .pollInterval(2, SECONDS).until(
+ {
+ def jobAfterRestart = sql """
+ SELECT Id, Status
+ FROM jobs("type"="insert")
+ WHERE Name='${jobName}' AND
ExecuteType='STREAMING'
+ """
+ log.info("job after FE restart: " +
jobAfterRestart)
+ jobAfterRestart.size() == 1
+ &&
jobAfterRestart.get(0).get(0).toString() == jobId
+ && jobAfterRestart.get(0).get(1) ==
"FINISHED"
+ }
+ )
+
+ def rowsAfterRestart = sql """
+ SELECT COUNT(*), COUNT(DISTINCT id)
+ FROM ${currentDb}.${tableName}
+ """
+ assert rowsAfterRestart.size() == 1
+ assert rowsAfterRestart.get(0).get(0) == totalRows
+ assert rowsAfterRestart.get(0).get(1) == totalRows
+ } catch (Exception ex) {
+ def showJob = sql """
+ SELECT *
+ FROM jobs("type"="insert")
+ WHERE Name='${jobName}'
+ """
+ def showTask = sql """
+ SELECT *
+ FROM tasks("type"="insert")
+ WHERE JobName='${jobName}'
+ """
+ log.info("show job: " + showJob)
+ log.info("show task: " + showTask)
+ throw ex
+ } finally {
+ sql """DROP JOB IF EXISTS where jobname = '${jobName}'"""
+ }
+ }
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]