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]

Reply via email to