This is an automated email from the ASF dual-hosted git repository.

liaoxin01 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new cd3c0780422 [fix](streaming-job) Optimize snapshot offset persistence 
(#66238)
cd3c0780422 is described below

commit cd3c078042267998a9d12b8b263df0736d447c7d
Author: wudi <[email protected]>
AuthorDate: Thu Aug 6 15:35:54 2026 +0800

    [fix](streaming-job) Optimize snapshot offset persistence (#66238)
    
    ### What problem does this PR solve?
    
    Problem Summary:
    
    Snapshot tasks persisted the complete offset provider after every split
    commit, producing frequent edit logs while the snapshot state kept
    growing. The previous default snapshot split size also produced too many
    small snapshot tasks. Natural terminal transitions needed explicit
    persistence so finished jobs replay with their terminal metadata.
    
    This change:
    - Limits snapshot offset edit logs with a mutable 300-second interval
    while keeping binlog offset commits immediate.
    - Increases the default snapshot split size from 8,096 to 40,960 rows,
    while preserving explicitly configured values.
    - Removes obsolete snapshot recovery state after the first committed
    binlog offset for FROM-TO and CDC stream TVF jobs.
    - Persists natural FINISHED transitions and restores terminal replay
    fields and callbacks consistently.
    - Adds focused unit coverage and a MySQL snapshot-only FE restart
    regression case.
---
 .../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 b626cba8cf1..114112d9e58 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
@@ -1174,6 +1174,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