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

JNSimba 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 a3f9949f37b [improvement](streaming) Support one-time S3 streaming 
ingestion (#68007)
a3f9949f37b is described below

commit a3f9949f37bfcc7d3921ff3285db642cdce021cd
Author: wudi <[email protected]>
AuthorDate: Fri Oct 9 11:38:56 2026 +0800

    [improvement](streaming) Support one-time S3 streaming ingestion (#68007)
    
    ### What problem does this PR solve?
    
    Problem Summary: S3 streaming insert jobs currently keep polling for new
    files after consuming all files matched by the path. This change adds
    `s3.ingestion_mode=ONCE`, which imports matching files in lexical
    batches and marks the job as `FINISHED` after the final batch commits
    successfully. Recovered jobs probe from their committed offset and
    finish when no files remain. ONCE jobs reject user-specified offsets,
    including offset changes through ALTER JOB.
---
 .../insert/streaming/StreamingInsertJob.java       |  76 +++++++----
 .../insert/streaming/StreamingJobProperties.java   |  25 +++-
 .../streaming/StreamingJobSchedulerTask.java       |   4 +-
 .../job/offset/SourceOffsetProviderFactory.java    |   7 +-
 .../org/apache/doris/job/offset/s3/S3Offset.java   |   2 +
 .../job/offset/s3/S3SourceOffsetProvider.java      |  30 ++++-
 .../job/scheduler/StreamingTaskScheduler.java      |   6 +-
 .../trees/plans/commands/AlterJobCommand.java      |  19 ++-
 .../StreamingInsertJobOffsetPersistenceTest.java   | 104 ++++++++++++++-
 .../StreamingInsertJobStatusTransitionTest.java    |   9 ++
 .../streaming/StreamingJobPropertiesTest.java      |  43 ++++++
 .../job/offset/s3/S3SourceOffsetProviderTest.java  | 147 +++++++++++++++++++++
 .../job/scheduler/StreamingTaskSchedulerTest.java  |  70 ++++++++++
 .../test_streaming_insert_job_s3_once.groovy       | 107 +++++++++++++++
 14 files changed, 610 insertions(+), 39 deletions(-)

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 14c08866633..995fcfca6aa 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
@@ -404,7 +404,7 @@ public class StreamingInsertJob extends 
AbstractJob<StreamingJobSchedulerTask, M
     private SourceOffsetProvider createOffsetProvider(Map<String, String> 
jdbcSourceProps) {
         SourceOffsetProvider provider;
         if (tvfType != null) {
-            provider = 
SourceOffsetProviderFactory.createSourceOffsetProvider(tvfType);
+            provider = 
SourceOffsetProviderFactory.createSourceOffsetProvider(tvfType, jobProperties);
         } else {
             provider = new JdbcSourceOffsetProvider(getJobId(), 
dataSourceType, jdbcSourceProps);
         }
@@ -555,6 +555,9 @@ public class StreamingInsertJob extends 
AbstractJob<StreamingJobSchedulerTask, M
     public void updateJobStatus(JobStatus status) throws JobException {
         lock.writeLock().lock();
         try {
+            if (isFinalStatus() && !getJobStatus().equals(status)) {
+                throw new JobException("Can't update final job status " + 
getJobStatus() + " to " + status);
+            }
             super.updateJobStatus(status);
             if (JobStatus.PAUSED.equals(getJobStatus())) {
                 clearRunningStreamTask(status);
@@ -592,6 +595,41 @@ public class StreamingInsertJob extends 
AbstractJob<StreamingJobSchedulerTask, M
         }
     }
 
+    public boolean tryFinishJob() throws JobException {
+        lock.writeLock().lock();
+        try {
+            if (!isActive()) {
+                return false;
+            }
+            if (runningStreamTask != null && 
TaskStatus.PENDING.equals(runningStreamTask.getStatus())) {
+                // Cancel the waiting task when a metadata scan detects the 
end of the source.
+                cancelAllTasks(false);
+            }
+            resetFailureInfo(null);
+            updateJobStatus(JobStatus.FINISHED);
+            logUpdateOperation();
+            return true;
+        } finally {
+            lock.writeLock().unlock();
+        }
+    }
+
+    private boolean tryPauseJob(FailureReason reason) throws JobException {
+        lock.writeLock().lock();
+        try {
+            if (!isActive()
+                    || (getFailureReason() != null
+                    && 
InternalErrorCode.MANUAL_PAUSE_ERR.equals(getFailureReason().getCode()))) {
+                return false;
+            }
+            updateJobStatus(JobStatus.PAUSED);
+            setFailureReason(reason);
+            return true;
+        } finally {
+            lock.writeLock().unlock();
+        }
+    }
+
     public void resetFailureInfo(FailureReason reason) {
         this.setFailureReason(reason);
         // Currently, only delayMsg is present here, which needs to be cleared 
when the status changes.
@@ -766,17 +804,8 @@ public class StreamingInsertJob extends 
AbstractJob<StreamingJobSchedulerTask, M
             offsetProvider.fetchRemoteMeta(props);
         } catch (Exception ex) {
             log.warn("fetch remote meta failed, job id: {}", getJobId(), ex);
-            if (this.getFailureReason() == null
-                    || 
!InternalErrorCode.MANUAL_PAUSE_ERR.equals(this.getFailureReason().getCode())) {
-                // When a job is manually paused, it does not need to be set 
again,
-                // otherwise, it may be woken up by auto resume.
-                // Pause before setting the reason: updateJobStatus's 
writeLock orders this after any
-                // task-success callback that clears failureReason, so a 
success can't wipe the reason.
-                this.updateJobStatus(JobStatus.PAUSED);
-                this.setFailureReason(
-                        new 
FailureReason(InternalErrorCode.GET_REMOTE_DATA_ERROR,
-                                "Failed to fetch meta, " + ex.getMessage()));
-
+            if (tryPauseJob(new 
FailureReason(InternalErrorCode.GET_REMOTE_DATA_ERROR,
+                    "Failed to fetch meta, " + ex.getMessage()))) {
                 if (MetricRepo.isInit) {
                     
MetricRepo.COUNTER_STREAMING_JOB_GET_META_FAIL_COUNT.increase(1L);
                 }
@@ -814,26 +843,24 @@ public class StreamingInsertJob extends 
AbstractJob<StreamingJobSchedulerTask, M
             }
         } catch (Exception ex) {
             log.warn("advance splits failed, job id: {}", getJobId(), ex);
-            if (this.getFailureReason() == null
-                    || 
!InternalErrorCode.MANUAL_PAUSE_ERR.equals(this.getFailureReason().getCode())) {
-                this.setFailureReason(new FailureReason(
-                        InternalErrorCode.GET_REMOTE_DATA_ERROR,
-                        "Failed to advance splits, " + ex.getMessage()));
-                this.updateJobStatus(JobStatus.PAUSED);
-            }
+            tryPauseJob(new 
FailureReason(InternalErrorCode.GET_REMOTE_DATA_ERROR,
+                    "Failed to advance splits, " + ex.getMessage()));
         }
     }
 
     public boolean needScheduleTask() {
         readLock();
         try {
-            return (getJobStatus().equals(JobStatus.RUNNING)
-                    || getJobStatus().equals(JobStatus.PENDING));
+            return isActive();
         } finally {
             readUnlock();
         }
     }
 
+    private boolean isActive() {
+        return JobStatus.PENDING.equals(getJobStatus()) || 
JobStatus.RUNNING.equals(getJobStatus());
+    }
+
     public void clearRunningStreamTask(JobStatus newJobStatus) {
         if (runningStreamTask != null) {
             log.info("clear running streaming insert task for job {}, task {}, 
status {} ",
@@ -972,6 +999,7 @@ public class StreamingInsertJob extends 
AbstractJob<StreamingJobSchedulerTask, M
         this.jobStatistic.setFileSize(attachment.getFileBytes());
         this.jobStatistic.setFilteredRows(attachment.getFilteredRows());
         
offsetProvider.updateOffset(offsetProvider.deserializeOffset(attachment.getOffset()));
+        this.offsetProviderPersist = offsetProvider.getPersistInfo();
 
         //update metric
         if (MetricRepo.isInit && !isReplay) {
@@ -1488,6 +1516,9 @@ public class StreamingInsertJob extends 
AbstractJob<StreamingJobSchedulerTask, M
 
     @Override
     public void gsonPostProcess() throws IOException {
+        if (jobProperties == null && properties != null) {
+            jobProperties = new StreamingJobProperties(properties);
+        }
         if (offsetProvider == null) {
             offsetProvider = createOffsetProvider(sourceProperties);
             if (tvfType != null) {
@@ -1495,9 +1526,6 @@ public class StreamingInsertJob extends 
AbstractJob<StreamingJobSchedulerTask, M
             }
         }
 
-        if (jobProperties == null && properties != null) {
-            jobProperties = new StreamingJobProperties(properties);
-        }
         recomputeDerivedFields();
 
         if (null == getSucceedTaskCount()) {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobProperties.java
 
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobProperties.java
index 1bde1bf6205..d0ea8ec6055 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobProperties.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobProperties.java
@@ -33,6 +33,7 @@ import java.util.Arrays;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
+import java.util.Locale;
 import java.util.Map;
 import java.util.Set;
 
@@ -41,12 +42,16 @@ public class StreamingJobProperties implements 
JobProperties {
     public static final String MAX_INTERVAL_SECOND_PROPERTY = "max_interval";
     public static final String S3_MAX_BATCH_FILES_PROPERTY = 
"s3.max_batch_files";
     public static final String S3_MAX_BATCH_BYTES_PROPERTY = 
"s3.max_batch_bytes";
+    public static final String S3_INGESTION_MODE_PROPERTY = 
"s3.ingestion_mode";
+    public static final String S3_INGESTION_MODE_LEXICAL = "LEXICAL";
+    public static final String S3_INGESTION_MODE_ONCE = "ONCE";
     public static final String SESSION_VAR_PREFIX = "session.";
     public static final String INTERNAL_KEY_PREFIX = "__";
     public static final String OFFSET_PROPERTY = "offset";
     public static final String COMPUTE_GROUP_PROPERTY = "compute_group";
     public static final List<String> SUPPORT_STREAM_JOB_PROPS = 
Arrays.asList(MAX_INTERVAL_SECOND_PROPERTY,
-            S3_MAX_BATCH_FILES_PROPERTY, S3_MAX_BATCH_BYTES_PROPERTY, 
OFFSET_PROPERTY, COMPUTE_GROUP_PROPERTY);
+            S3_MAX_BATCH_FILES_PROPERTY, S3_MAX_BATCH_BYTES_PROPERTY, 
S3_INGESTION_MODE_PROPERTY,
+            OFFSET_PROPERTY, COMPUTE_GROUP_PROPERTY);
 
     public static final long DEFAULT_MAX_INTERVAL_SECOND = 10;
     public static final long DEFAULT_MAX_S3_BATCH_FILES = 256;
@@ -123,6 +128,15 @@ public class StreamingJobProperties implements 
JobProperties {
                         + " must be between 100 MB and 10 GB, but was "
                         + 
properties.get(StreamingJobProperties.S3_MAX_BATCH_BYTES_PROPERTY));
 
+        String ingestionMode = getS3IngestionMode();
+        if (!S3_INGESTION_MODE_LEXICAL.equals(ingestionMode)
+                && !S3_INGESTION_MODE_ONCE.equals(ingestionMode)) {
+            throw new AnalysisException("Unsupported s3.ingestion_mode: " + 
ingestionMode);
+        }
+        if (S3_INGESTION_MODE_ONCE.equals(ingestionMode) && 
properties.containsKey(OFFSET_PROPERTY)) {
+            throw new AnalysisException("offset is not supported when 
s3.ingestion_mode is ONCE");
+        }
+
         // validate session variables
         try {
             Map<String, String> sessionVarMap = parseSessionVarMap();
@@ -204,4 +218,13 @@ public class StreamingJobProperties implements 
JobProperties {
     public String getComputeGroup() {
         return properties.get(COMPUTE_GROUP_PROPERTY);
     }
+
+    public String getS3IngestionMode() {
+        return properties.getOrDefault(S3_INGESTION_MODE_PROPERTY, 
S3_INGESTION_MODE_LEXICAL)
+                .trim().toUpperCase(Locale.ROOT);
+    }
+
+    public boolean isS3OnceMode() {
+        return S3_INGESTION_MODE_ONCE.equals(getS3IngestionMode());
+    }
 }
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 0a4b7ca09f6..d183050d302 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
@@ -74,9 +74,7 @@ public class StreamingJobSchedulerTask extends AbstractTask {
         }
         if (streamingInsertJob.hasReachedEnd()) {
             // 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();
+            streamingInsertJob.tryFinishJob();
             return;
         }
         streamingInsertJob.createStreamingTask();
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/SourceOffsetProviderFactory.java
 
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/SourceOffsetProviderFactory.java
index 30f9d0edd57..631162d27fe 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/SourceOffsetProviderFactory.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/SourceOffsetProviderFactory.java
@@ -18,6 +18,7 @@
 package org.apache.doris.job.offset;
 
 import org.apache.doris.job.exception.JobException;
+import org.apache.doris.job.extensions.insert.streaming.StreamingJobProperties;
 import org.apache.doris.job.offset.jdbc.JdbcTvfSourceOffsetProvider;
 import org.apache.doris.job.offset.s3.S3SourceOffsetProvider;
 
@@ -35,8 +36,12 @@ public class SourceOffsetProviderFactory {
         map.put("cdc_stream", JdbcTvfSourceOffsetProvider.class);
     }
 
-    public static SourceOffsetProvider createSourceOffsetProvider(String 
sourceType) {
+    public static SourceOffsetProvider createSourceOffsetProvider(
+            String sourceType, StreamingJobProperties jobProperties) {
         try {
+            if ("s3".equalsIgnoreCase(sourceType) && 
jobProperties.isS3OnceMode()) {
+                return new S3SourceOffsetProvider(jobProperties);
+            }
             Class<? extends SourceOffsetProvider> cla = 
map.get(sourceType.toLowerCase());
             if (cla == null) {
                 throw new JobException("Unsupported source type: " + 
sourceType);
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3Offset.java 
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3Offset.java
index ebd33803490..44df1756e02 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3Offset.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3Offset.java
@@ -34,6 +34,8 @@ public class S3Offset implements Offset {
     // s3://bucket/path/{1.csv,2.csv}
     String fileLists;
     int fileNum;
+    @SerializedName("lastBatch")
+    boolean lastBatch;
 
     @Override
     public String toSerializedJson() {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3SourceOffsetProvider.java
 
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3SourceOffsetProvider.java
index 2777013471c..f58287ee9a5 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3SourceOffsetProvider.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3SourceOffsetProvider.java
@@ -49,8 +49,18 @@ import java.util.stream.Collectors;
 
 @Log4j2
 public class S3SourceOffsetProvider implements SourceOffsetProvider {
-    S3Offset currentOffset;
-    String maxEndFile;
+    private final boolean onceMode;
+    private volatile S3Offset noMoreFilesAfterOffset;
+    volatile S3Offset currentOffset;
+    volatile String maxEndFile;
+
+    public S3SourceOffsetProvider() {
+        this.onceMode = false;
+    }
+
+    public S3SourceOffsetProvider(StreamingJobProperties jobProperties) {
+        this.onceMode = jobProperties.isS3OnceMode();
+    }
 
     @Override
     public String getSourceType() {
@@ -95,6 +105,7 @@ public class S3SourceOffsetProvider implements 
SourceOffsetProvider {
                 offset.setEndFile(lastFile);
                 offset.setFileNum(rfiles.size());
                 maxEndFile = globListing.getMaxFile();
+                offset.setLastBatch(onceMode && 
lastFile.equals(globListing.getMaxFile()));
             } else {
                 throw new RuntimeException("No new files found in path: " + 
filePath);
             }
@@ -159,7 +170,8 @@ public class S3SourceOffsetProvider implements 
SourceOffsetProvider {
         Map<String, String> copiedProps = 
Maps.newTreeMap(String.CASE_INSENSITIVE_ORDER);
         copiedProps.putAll(properties);
         StorageAdapter storageAdapter = StorageAdapter.of(copiedProps);
-        String startFile = currentOffset == null ? null : 
currentOffset.endFile;
+        S3Offset offsetAtScan = currentOffset;
+        String startFile = offsetAtScan == null ? null : offsetAtScan.endFile;
         try (FileSystem fileSystem = 
FileSystemFactory.getFileSystem(storageAdapter)) {
             String uri = storageAdapter.validateAndGetUri(copiedProps);
             String filePath = storageAdapter.validateAndNormalizeUri(uri);
@@ -168,6 +180,9 @@ public class S3SourceOffsetProvider implements 
SourceOffsetProvider {
                 throw new java.io.IOException("debug point: simulated S3 auth 
error");
             }
             GlobListing globListing = 
fileSystem.globListWithLimit(Location.of(filePath), startFile, 1, 1);
+            if (onceMode) {
+                noMoreFilesAfterOffset = globListing.getFiles().isEmpty() ? 
offsetAtScan : null;
+            }
             if (!globListing.getFiles().isEmpty() && 
StringUtils.isNotEmpty(globListing.getMaxFile())) {
                 maxEndFile = globListing.getMaxFile();
             }
@@ -176,6 +191,9 @@ public class S3SourceOffsetProvider implements 
SourceOffsetProvider {
 
     @Override
     public boolean hasMoreDataToConsume() {
+        if (hasReachedEnd()) {
+            return false;
+        }
         if (currentOffset == null || currentOffset.endFile == null) {
             return true;
         }
@@ -186,6 +204,12 @@ public class S3SourceOffsetProvider implements 
SourceOffsetProvider {
         return false;
     }
 
+    @Override
+    public boolean hasReachedEnd() {
+        S3Offset offset = currentOffset;
+        return onceMode && offset != null && (offset.isLastBatch() || 
noMoreFilesAfterOffset == offset);
+    }
+
     @Override
     public String getPersistInfo() {
         if (currentOffset == null) {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/job/scheduler/StreamingTaskScheduler.java
 
b/fe/fe-core/src/main/java/org/apache/doris/job/scheduler/StreamingTaskScheduler.java
index 91d5fb1a658..62d324061b0 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/job/scheduler/StreamingTaskScheduler.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/job/scheduler/StreamingTaskScheduler.java
@@ -104,7 +104,7 @@ public class StreamingTaskScheduler extends MasterDaemon {
         }
     }
 
-    private void scheduleOneTask(AbstractStreamingTask task) {
+    private void scheduleOneTask(AbstractStreamingTask task) throws 
JobException {
         if (DebugPointUtil.isEnable("StreamingJob.scheduleTask.exception")) {
             throw new RuntimeException("debug point 
StreamingJob.scheduleTask.exception");
         }
@@ -122,6 +122,10 @@ public class StreamingTaskScheduler extends MasterDaemon {
         }
         // reject task if no more data to consume
         if (!job.hasMoreDataToConsume()) {
+            if (job.hasReachedEnd()) {
+                job.tryFinishJob();
+                return;
+            }
             String delayMsg = "No data available for consumption at the 
moment, will retry after "
                     + (System.currentTimeMillis() + DELAY_SCHEDULER_MS);
             job.setJobRuntimeMsg(delayMsg);
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/AlterJobCommand.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/AlterJobCommand.java
index 0a529d55424..e31af87bb36 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/AlterJobCommand.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/AlterJobCommand.java
@@ -41,6 +41,7 @@ import org.apache.doris.qe.StmtExecutor;
 import com.google.common.base.Preconditions;
 import org.apache.commons.lang3.StringUtils;
 
+import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
@@ -277,11 +278,19 @@ public class AlterJobCommand extends AlterCommand 
implements ForwardWithSync, Ne
     }
 
     private void validateProps(StreamingInsertJob streamingJob) throws 
AnalysisException {
-        StreamingJobProperties jobProperties = new 
StreamingJobProperties(properties);
-        jobProperties.validate();
-        if (jobProperties.getOffsetProperty() != null) {
-            
streamingJob.validateAlterOffset(jobProperties.getOffsetProperty());
-            streamingJob.validateOffset(jobProperties.getOffsetProperty());
+        StreamingJobProperties originJobProperties =
+                new StreamingJobProperties(streamingJob.getProperties());
+        Map<String, String> mergedProperties = new 
HashMap<>(streamingJob.getProperties());
+        mergedProperties.putAll(properties);
+        StreamingJobProperties updatedJobProperties = new 
StreamingJobProperties(mergedProperties);
+        updatedJobProperties.validate();
+        if 
(!originJobProperties.getS3IngestionMode().equals(updatedJobProperties.getS3IngestionMode()))
 {
+            throw new AnalysisException("s3.ingestion_mode cannot be altered");
+        }
+        String offset = properties.get(StreamingJobProperties.OFFSET_PROPERTY);
+        if (offset != null) {
+            streamingJob.validateAlterOffset(offset);
+            streamingJob.validateOffset(offset);
         }
     }
 
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
index e3246f1b530..1a724b2330a 100644
--- 
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
@@ -27,7 +27,11 @@ 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.SourceOffsetProvider;
 import org.apache.doris.job.offset.jdbc.JdbcSourceOffsetProvider;
+import org.apache.doris.job.offset.s3.S3Offset;
+import org.apache.doris.job.offset.s3.S3SourceOffsetProvider;
+import org.apache.doris.persist.gson.GsonUtils;
 import org.apache.doris.transaction.GlobalTransactionMgrIface;
 import org.apache.doris.transaction.TxnStateCallbackFactory;
 
@@ -38,6 +42,7 @@ import org.mockito.Mockito;
 
 import java.util.Collections;
 import java.util.HashMap;
+import java.util.Map;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 
 public class StreamingInsertJobOffsetPersistenceTest {
@@ -163,6 +168,103 @@ public class StreamingInsertJobOffsetPersistenceTest {
         }
     }
 
+    @Test
+    public void testS3OnceLastBatchFinishesOnlyAfterSuccess() throws Exception 
{
+        StreamingJobProperties properties = new 
StreamingJobProperties(Map.of("s3.ingestion_mode", "ONCE"));
+        TestStreamingInsertJob failedJob = newJob(new 
S3SourceOffsetProvider(properties), 1017L);
+        NoopStreamingMultiTblTask failedTask =
+                (NoopStreamingMultiTblTask) 
Deencapsulation.getField(failedJob, "runningStreamTask");
+        S3Offset failedOffset = new S3Offset();
+        failedOffset.setLastBatch(true);
+        Deencapsulation.setField(failedTask, "runningOffset", failedOffset);
+        failedTask.setErrMsg("failed");
+
+        TestStreamingInsertJob succeededJob = newJob(new 
S3SourceOffsetProvider(properties), 1018L);
+        NoopStreamingMultiTblTask succeededTask =
+                (NoopStreamingMultiTblTask) 
Deencapsulation.getField(succeededJob, "runningStreamTask");
+        S3Offset succeededOffset = new S3Offset();
+        succeededOffset.setEndFile("data/b.csv");
+        succeededOffset.setLastBatch(true);
+        Deencapsulation.setField(succeededTask, "runningOffset", 
succeededOffset);
+
+        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);
+
+            Assertions.assertEquals(JobStatus.RUNNING, 
failedJob.getJobStatus());
+            Assertions.assertFalse(failedJob.hasReachedEnd());
+            failedJob.onStreamTaskFail(failedTask);
+            Assertions.assertEquals(JobStatus.PAUSED, 
failedJob.getJobStatus());
+            Assertions.assertEquals(0, failedJob.journalCount);
+
+            Assertions.assertEquals(JobStatus.RUNNING, 
succeededJob.getJobStatus());
+            Assertions.assertFalse(succeededJob.hasReachedEnd());
+            Deencapsulation.invoke(succeededJob, "updateJobStatisticAndOffset",
+                    new StreamingTaskTxnCommitAttachment(9001L, 1018L, 0, 0, 
0, 0, 0,
+                            succeededOffset.toSerializedJson()), false);
+            succeededJob.onStreamTaskSuccess(succeededTask);
+            Assertions.assertEquals(JobStatus.FINISHED, 
succeededJob.getJobStatus());
+            Assertions.assertEquals(1, succeededJob.journalCount);
+        }
+    }
+
+    @Test
+    public void testRecoveredSourceFinishesBeforeCreatingTask() throws 
Exception {
+        TestStreamingInsertJob job = newJob(new EndJdbcSourceOffsetProvider(), 
1019L);
+        job.setJobStatus(JobStatus.PENDING);
+        Deencapsulation.setField(job, "runningStreamTask", null);
+        job.setJobRuntimeMsg("will retry");
+
+        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);
+
+            Assertions.assertTrue(job.tryFinishJob());
+        }
+
+        Assertions.assertEquals(JobStatus.FINISHED, job.getJobStatus());
+        Assertions.assertNull(job.getRunningStreamTask());
+        Assertions.assertEquals("", job.getJobRuntimeMsg());
+        Assertions.assertEquals(1, job.journalCount);
+    }
+
+    @Test
+    public void testCloudReplayRefreshesPersistedOffset() {
+        Map<String, String> properties = Map.of("s3.ingestion_mode", "ONCE");
+        S3SourceOffsetProvider provider = new S3SourceOffsetProvider(new 
StreamingJobProperties(properties));
+        StreamingInsertJob job = new StreamingInsertJob();
+        job.offsetProvider = provider;
+        job.setJobStatus(JobStatus.PENDING);
+        Deencapsulation.setField(job, "properties", properties);
+        Deencapsulation.setField(job, "tvfType", "s3");
+        job.setOffsetProviderPersist("{\"endFile\":\"data/a.csv\"}");
+        S3Offset offset = new S3Offset();
+        offset.setEndFile("data/b.csv");
+        offset.setLastBatch(true);
+        StreamingTaskTxnCommitAttachment attachment = new 
StreamingTaskTxnCommitAttachment(
+                9001L, 1020L, 0, 0, 0, 0, 0, offset.toSerializedJson());
+
+        Deencapsulation.invoke(job, "updateCloudJobStatisticAndOffset", 
attachment, true);
+
+        Assertions.assertEquals(provider.getPersistInfo(), 
job.getOffsetProviderPersist());
+        
Assertions.assertTrue(job.getOffsetProviderPersist().contains("data/b.csv"));
+        StreamingInsertJob recovered = GsonUtils.GSON.fromJson(
+                GsonUtils.GSON.toJson(job), StreamingInsertJob.class);
+        Assertions.assertEquals(JobStatus.PENDING, recovered.getJobStatus());
+        Assertions.assertTrue(recovered.hasReachedEnd());
+        Assertions.assertFalse(recovered.hasMoreDataToConsume());
+    }
+
     @Test
     public void testReplayUpdatedRestoresFinalStateAndRemovesCallback() {
         TestStreamingInsertJob job = newJob(new JdbcSourceOffsetProvider(), 
1013L);
@@ -195,7 +297,7 @@ public class StreamingInsertJobOffsetPersistenceTest {
         Assertions.assertEquals(1234L, job.getStartTimeMs());
     }
 
-    private static TestStreamingInsertJob newJob(JdbcSourceOffsetProvider 
provider, long taskId) {
+    private static TestStreamingInsertJob newJob(SourceOffsetProvider 
provider, long taskId) {
         TestStreamingInsertJob job = new TestStreamingInsertJob();
         Deencapsulation.setField(job, "lock", new 
ReentrantReadWriteLock(true));
         Deencapsulation.setField(job, "jobId", 9001L);
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobStatusTransitionTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobStatusTransitionTest.java
index 3e5c97c61d7..23b35983d84 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobStatusTransitionTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobStatusTransitionTest.java
@@ -19,6 +19,7 @@ package org.apache.doris.job.extensions.insert.streaming;
 
 import org.apache.doris.common.jmockit.Deencapsulation;
 import org.apache.doris.job.common.JobStatus;
+import org.apache.doris.job.exception.JobException;
 
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
@@ -63,4 +64,12 @@ public class StreamingInsertJobStatusTransitionTest {
         Assertions.assertFalse(job.updateJobStatusIfCurrent(JobStatus.PENDING, 
JobStatus.RUNNING));
         Assertions.assertEquals(JobStatus.STOPPED, job.getJobStatus());
     }
+
+    @Test
+    public void testFinishedStatusCannotBeOverwritten() {
+        StreamingInsertJob job = newJob(JobStatus.FINISHED);
+
+        Assertions.assertThrows(JobException.class, () -> 
job.updateJobStatus(JobStatus.PAUSED));
+        Assertions.assertEquals(JobStatus.FINISHED, job.getJobStatus());
+    }
 }
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobPropertiesTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobPropertiesTest.java
index 64f4515816e..56fc54fc91f 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobPropertiesTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingJobPropertiesTest.java
@@ -18,8 +18,11 @@
 package org.apache.doris.job.extensions.insert.streaming;
 
 import org.apache.doris.common.AnalysisException;
+import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.job.common.JobStatus;
 import org.apache.doris.job.exception.JobException;
 import org.apache.doris.job.extensions.insert.InsertTask;
+import org.apache.doris.nereids.trees.plans.commands.AlterJobCommand;
 import org.apache.doris.qe.ConnectContext;
 import org.apache.doris.qe.SessionVariable;
 
@@ -27,9 +30,49 @@ import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
 import java.util.HashMap;
+import java.util.Map;
 
 public class StreamingJobPropertiesTest {
 
+    @Test
+    public void testS3OnceModeValidationAndAlter() throws Exception {
+        StreamingJobProperties properties = new 
StreamingJobProperties(Map.of("s3.ingestion_mode", " once "));
+        properties.validate();
+        Assertions.assertTrue(properties.isS3OnceMode());
+        Assertions.assertThrows(AnalysisException.class,
+                () -> new StreamingJobProperties(Map.of("s3.ingestion_mode", 
"invalid")).validate());
+        Assertions.assertThrows(AnalysisException.class,
+                () -> new StreamingJobProperties(
+                        Map.of("s3.ingestion_mode", "ONCE", "offset", 
"{\"fileName\":\"a.csv\"}"))
+                        .validate());
+
+        StreamingInsertJob job = new StreamingInsertJob();
+        Deencapsulation.setField(job, "properties", 
properties.getProperties());
+        AlterJobCommand alterBatch = new AlterJobCommand("job", 
Map.of("s3.max_batch_files", "1"),
+                null, null, null, Map.of(), Map.of());
+        Deencapsulation.invoke(alterBatch, "validateProps", job);
+        AlterJobCommand alterMode = new AlterJobCommand("job", 
Map.of("s3.ingestion_mode", "LEXICAL"),
+                null, null, null, Map.of(), Map.of());
+        Assertions.assertThrows(AnalysisException.class,
+                () -> Deencapsulation.invoke(alterMode, "validateProps", job));
+        AlterJobCommand alterOffset = new AlterJobCommand("job",
+                Map.of("offset", "{\"fileName\":\"a.csv\"}"), null, null, 
null, Map.of(), Map.of());
+        Assertions.assertThrows(AnalysisException.class,
+                () -> Deencapsulation.invoke(alterOffset, "validateProps", 
job));
+    }
+
+    @Test
+    public void testS3OnceModeRestoredWithLegacyOffset() throws Exception {
+        StreamingInsertJob job = new StreamingInsertJob();
+        Deencapsulation.setField(job, "properties", 
Map.of("s3.ingestion_mode", "ONCE"));
+        Deencapsulation.setField(job, "tvfType", "s3");
+        job.setOffsetProviderPersist("{\"endFile\":\"data/a.csv\"}");
+        job.gsonPostProcess();
+        job.setJobStatus(JobStatus.RUNNING);
+        Assertions.assertFalse(job.hasMoreDataToConsume());
+        Assertions.assertFalse(job.hasReachedEnd());
+    }
+
     /**
      * Simulate FE restart: constructor is called without validate().
      * Before the fix, maxIntervalSecond would be 0 when properties is 
non-empty,
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/job/offset/s3/S3SourceOffsetProviderTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/job/offset/s3/S3SourceOffsetProviderTest.java
new file mode 100644
index 00000000000..4f72cfd909e
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/job/offset/s3/S3SourceOffsetProviderTest.java
@@ -0,0 +1,147 @@
+// 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.offset.s3;
+
+import org.apache.doris.datasource.storage.StorageAdapter;
+import org.apache.doris.filesystem.FileEntry;
+import org.apache.doris.filesystem.FileSystem;
+import org.apache.doris.filesystem.GlobListing;
+import org.apache.doris.filesystem.Location;
+import org.apache.doris.fs.FileSystemFactory;
+import org.apache.doris.job.extensions.insert.streaming.StreamingJobProperties;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+import java.io.IOException;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+
+public class S3SourceOffsetProviderTest {
+    private static final Map<String, String> TVF_PROPS = Map.of(
+            "uri", "s3://bucket/data/*.csv", "s3.endpoint", 
"s3.us-east-1.amazonaws.com",
+            "s3.region", "us-east-1", "s3.access_key", "ak", "s3.secret_key", 
"sk");
+    private static final StreamingJobProperties ONCE_PROPS = new 
StreamingJobProperties(
+            Map.of("s3.ingestion_mode", "ONCE", "s3.max_batch_files", "1"));
+
+    @Test
+    public void testOnceMarksLastBatchAndRecovers() throws Exception {
+        FileSystem fs = Mockito.mock(FileSystem.class);
+        try (MockedStatic<FileSystemFactory> factory = 
Mockito.mockStatic(FileSystemFactory.class)) {
+            factory.when(() -> 
FileSystemFactory.getFileSystem(Mockito.any(StorageAdapter.class))).thenReturn(fs);
+            Mockito.when(fs.globListWithLimit(Mockito.any(Location.class), 
Mockito.isNull(),
+                    Mockito.eq(1L), 
Mockito.eq(1L))).thenReturn(page("data/a.csv", "data/b.csv"));
+            Mockito.when(fs.globListWithLimit(Mockito.any(Location.class), 
Mockito.isNull(),
+                    Mockito.eq(ONCE_PROPS.getS3BatchBytes()), Mockito.eq(1L)))
+                    .thenReturn(page("data/a.csv", "data/b.csv"));
+            Mockito.when(fs.globListWithLimit(Mockito.any(Location.class), 
Mockito.eq("data/a.csv"),
+                    Mockito.eq(1L), Mockito.eq(1L)))
+                    .thenReturn(page("data/b.csv", "data/b.csv"));
+            Mockito.when(fs.globListWithLimit(Mockito.any(Location.class), 
Mockito.eq("data/a.csv"),
+                    Mockito.eq(ONCE_PROPS.getS3BatchBytes()), Mockito.eq(1L)))
+                    .thenReturn(page("data/b.csv", "data/b.csv"));
+            Mockito.when(fs.globListWithLimit(Mockito.any(Location.class), 
Mockito.eq("data/b.csv"),
+                    Mockito.eq(1L), Mockito.eq(1L)))
+                    .thenReturn(new GlobListing(Collections.emptyList(), 
"bucket", "data/", ""));
+
+            S3SourceOffsetProvider provider = new 
S3SourceOffsetProvider(ONCE_PROPS);
+            Assertions.assertTrue(provider.hasMoreDataToConsume());
+            provider.fetchRemoteMeta(TVF_PROPS);
+            Assertions.assertTrue(provider.hasMoreDataToConsume());
+            S3Offset first = provider.getNextOffset(ONCE_PROPS, TVF_PROPS);
+            Assertions.assertEquals("data/a.csv", first.getEndFile());
+            Assertions.assertFalse(first.isLastBatch());
+            Assertions.assertEquals("data/a.csv", 
provider.getNextOffset(ONCE_PROPS, TVF_PROPS).getEndFile());
+            
provider.updateOffset(provider.deserializeOffset(first.toSerializedJson()));
+            Assertions.assertTrue(provider.hasMoreDataToConsume());
+            S3Offset second = provider.getNextOffset(ONCE_PROPS, TVF_PROPS);
+            Assertions.assertEquals("data/b.csv", second.getEndFile());
+            Assertions.assertTrue(second.isLastBatch());
+            Assertions.assertFalse(provider.hasReachedEnd());
+            // A concurrent metadata probe cannot change this task's 
completion decision.
+            provider.fetchRemoteMeta(TVF_PROPS);
+            Assertions.assertTrue(second.isLastBatch());
+            Assertions.assertFalse(provider.hasReachedEnd());
+            Assertions.assertEquals("data/b.csv", 
provider.getNextOffset(ONCE_PROPS, TVF_PROPS).getEndFile());
+            
provider.updateOffset(provider.deserializeOffset(second.toSerializedJson()));
+            Assertions.assertTrue(provider.hasReachedEnd());
+            Assertions.assertFalse(provider.hasMoreDataToConsume());
+
+            S3SourceOffsetProvider recovered = new 
S3SourceOffsetProvider(ONCE_PROPS);
+            recovered.restoreFromPersistInfo(provider.getPersistInfo());
+            Assertions.assertTrue(recovered.hasReachedEnd());
+            Assertions.assertFalse(recovered.hasMoreDataToConsume());
+
+            recovered.fetchRemoteMeta(TVF_PROPS);
+            Assertions.assertTrue(recovered.hasReachedEnd());
+
+            Mockito.when(fs.globListWithLimit(Mockito.any(Location.class), 
Mockito.eq("data/b.csv"),
+                    Mockito.eq(1L), 
Mockito.eq(1L))).thenReturn(page("data/c.csv", "data/c.csv"));
+            S3SourceOffsetProvider recoveredWithNewFile = new 
S3SourceOffsetProvider(ONCE_PROPS);
+            
recoveredWithNewFile.restoreFromPersistInfo(provider.getPersistInfo());
+            recoveredWithNewFile.fetchRemoteMeta(TVF_PROPS);
+            Assertions.assertTrue(recoveredWithNewFile.hasReachedEnd());
+            
Assertions.assertFalse(recoveredWithNewFile.hasMoreDataToConsume());
+
+            S3SourceOffsetProvider legacyWithNewFile = new 
S3SourceOffsetProvider();
+            
legacyWithNewFile.restoreFromPersistInfo("{\"endFile\":\"data/b.csv\"}");
+            legacyWithNewFile.fetchRemoteMeta(TVF_PROPS);
+            Assertions.assertFalse(legacyWithNewFile.hasReachedEnd());
+            Assertions.assertTrue(legacyWithNewFile.hasMoreDataToConsume());
+
+            S3SourceOffsetProvider lexical = new S3SourceOffsetProvider();
+            lexical.restoreFromPersistInfo(provider.getPersistInfo());
+            lexical.fetchRemoteMeta(TVF_PROPS);
+            Assertions.assertFalse(lexical.hasReachedEnd());
+            Assertions.assertTrue(lexical.hasMoreDataToConsume());
+        }
+    }
+
+    @Test
+    public void testEmptyAndListingErrorRemainDistinct() throws Exception {
+        FileSystem fs = Mockito.mock(FileSystem.class);
+        try (MockedStatic<FileSystemFactory> factory = 
Mockito.mockStatic(FileSystemFactory.class)) {
+            factory.when(() -> 
FileSystemFactory.getFileSystem(Mockito.any(StorageAdapter.class))).thenReturn(fs);
+            Mockito.when(fs.globListWithLimit(Mockito.any(Location.class), 
Mockito.isNull(),
+                    Mockito.anyLong(), Mockito.anyLong()))
+                    .thenThrow(new IOException("listing failed"))
+                    .thenReturn(new GlobListing(Collections.emptyList(), 
"bucket", "data/", ""));
+
+            S3SourceOffsetProvider once = new 
S3SourceOffsetProvider(ONCE_PROPS);
+            Assertions.assertThrows(IOException.class, () -> 
once.fetchRemoteMeta(TVF_PROPS));
+            Assertions.assertFalse(once.hasReachedEnd());
+            once.fetchRemoteMeta(TVF_PROPS);
+            Assertions.assertFalse(once.hasReachedEnd());
+            Assertions.assertTrue(once.hasMoreDataToConsume());
+            RuntimeException emptyError = 
Assertions.assertThrows(RuntimeException.class,
+                    () -> once.getNextOffset(ONCE_PROPS, TVF_PROPS));
+            Assertions.assertTrue(emptyError.getMessage().contains("No new 
files found in path:"));
+
+            S3SourceOffsetProvider lexical = new S3SourceOffsetProvider();
+            Assertions.assertThrows(RuntimeException.class, () -> 
lexical.getNextOffset(ONCE_PROPS, TVF_PROPS));
+        }
+    }
+
+    private static GlobListing page(String key, String maxFile) {
+        return new GlobListing(List.of(new 
FileEntry(Location.of("s3://bucket/" + key),
+                10, false, 0, null)), "bucket", "data/", maxFile);
+    }
+}
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/job/scheduler/StreamingTaskSchedulerTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/job/scheduler/StreamingTaskSchedulerTest.java
new file mode 100644
index 00000000000..d2a5676102c
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/job/scheduler/StreamingTaskSchedulerTest.java
@@ -0,0 +1,70 @@
+// 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.scheduler;
+
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.job.common.JobStatus;
+import org.apache.doris.job.extensions.insert.streaming.AbstractStreamingTask;
+import org.apache.doris.job.extensions.insert.streaming.StreamingInsertJob;
+import org.apache.doris.job.manager.JobManager;
+import org.apache.doris.job.manager.StreamingTaskManager;
+
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+public class StreamingTaskSchedulerTest {
+
+    @Test
+    public void testRecoveredSourceFinishesOrContinuesAfterProbe() throws 
Exception {
+        Env env = Mockito.mock(Env.class);
+        JobManager<?, ?> jobManager = Mockito.mock(JobManager.class);
+        StreamingTaskManager taskManager = 
Mockito.mock(StreamingTaskManager.class);
+        StreamingInsertJob exhaustedJob = 
Mockito.mock(StreamingInsertJob.class);
+        AbstractStreamingTask exhaustedTask = 
Mockito.mock(AbstractStreamingTask.class);
+        StreamingInsertJob continuedJob = 
Mockito.mock(StreamingInsertJob.class);
+        AbstractStreamingTask continuedTask = 
Mockito.mock(AbstractStreamingTask.class);
+
+        Mockito.when(env.getJobManager()).thenReturn(jobManager);
+        
Mockito.when(jobManager.getStreamingTaskManager()).thenReturn(taskManager);
+        Mockito.when(exhaustedTask.getJobId()).thenReturn(1L);
+        Mockito.doReturn(exhaustedJob).when(jobManager).getJob(1L);
+        Mockito.when(exhaustedJob.needScheduleTask()).thenReturn(true);
+        Mockito.when(exhaustedJob.hasMoreDataToConsume()).thenReturn(false);
+        Mockito.when(exhaustedJob.hasReachedEnd()).thenReturn(true);
+        Mockito.when(continuedTask.getJobId()).thenReturn(2L);
+        Mockito.doReturn(continuedJob).when(jobManager).getJob(2L);
+        Mockito.when(continuedJob.needScheduleTask()).thenReturn(true);
+        Mockito.when(continuedJob.hasMoreDataToConsume()).thenReturn(true);
+
+        try (MockedStatic<Env> envMockedStatic = 
Mockito.mockStatic(Env.class)) {
+            envMockedStatic.when(Env::getCurrentEnv).thenReturn(env);
+            StreamingTaskScheduler scheduler = new StreamingTaskScheduler();
+
+            Deencapsulation.invoke(scheduler, "scheduleOneTask", 
exhaustedTask);
+            Mockito.verify(exhaustedJob).tryFinishJob();
+            Mockito.verify(exhaustedTask, Mockito.never()).execute();
+
+            Deencapsulation.invoke(scheduler, "scheduleOneTask", 
continuedTask);
+            Mockito.verify(taskManager).addRunningTask(continuedTask);
+            Mockito.verify(continuedTask).execute();
+            Mockito.verify(continuedJob, 
Mockito.never()).updateJobStatus(JobStatus.FINISHED);
+        }
+    }
+}
diff --git 
a/regression-test/suites/job_p0/streaming_job/test_streaming_insert_job_s3_once.groovy
 
b/regression-test/suites/job_p0/streaming_job/test_streaming_insert_job_s3_once.groovy
new file mode 100644
index 00000000000..799312f3de9
--- /dev/null
+++ 
b/regression-test/suites/job_p0/streaming_job/test_streaming_insert_job_s3_once.groovy
@@ -0,0 +1,107 @@
+// 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.awaitility.Awaitility
+
+import static java.util.concurrent.TimeUnit.SECONDS
+
+suite("test_streaming_insert_job_s3_once") {
+    sql """DROP JOB IF EXISTS WHERE jobname = 
'test_streaming_insert_job_s3_once'"""
+    sql """DROP TABLE IF EXISTS test_streaming_insert_job_s3_once_tbl FORCE"""
+
+    sql """
+        CREATE TABLE test_streaming_insert_job_s3_once_tbl (
+            `c1` INT NULL,
+            `c2` STRING NULL,
+            `c3` INT NULL
+        ) ENGINE=OLAP
+        DUPLICATE KEY(`c1`)
+        DISTRIBUTED BY HASH(`c1`) BUCKETS 3
+        PROPERTIES ("replication_allocation" = "tag.location.default: 1")
+    """
+
+    test {
+        sql """
+            CREATE JOB test_streaming_insert_job_s3_once
+            PROPERTIES (
+                "s3.ingestion_mode" = "ONCE",
+                "offset" = '{"fileName":"regression/load/data/example_0.csv"}'
+            )
+            ON STREAMING DO INSERT INTO test_streaming_insert_job_s3_once_tbl
+            SELECT * FROM S3 (
+                "uri" = 
"s3://${s3BucketName}/regression/load/data/example_[0-1].csv",
+                "format" = "csv",
+                "provider" = "${getS3Provider()}",
+                "column_separator" = ",",
+                "s3.endpoint" = "${getS3Endpoint()}",
+                "s3.region" = "${getS3Region()}",
+                "s3.access_key" = "${getS3AK()}",
+                "s3.secret_key" = "${getS3SK()}"
+            )
+        """
+        exception "offset is not supported when s3.ingestion_mode is ONCE"
+    }
+
+    sql """
+        CREATE JOB test_streaming_insert_job_s3_once
+        PROPERTIES (
+            "s3.ingestion_mode" = "ONCE",
+            "s3.max_batch_files" = "1"
+        )
+        ON STREAMING DO INSERT INTO test_streaming_insert_job_s3_once_tbl
+        SELECT * FROM S3 (
+            "uri" = 
"s3://${s3BucketName}/regression/load/data/example_[0-1].csv",
+            "format" = "csv",
+            "provider" = "${getS3Provider()}",
+            "column_separator" = ",",
+            "s3.endpoint" = "${getS3Endpoint()}",
+            "s3.region" = "${getS3Region()}",
+            "s3.access_key" = "${getS3AK()}",
+            "s3.secret_key" = "${getS3SK()}"
+        )
+    """
+
+    try {
+        Awaitility.await().atMost(300, SECONDS)
+                .pollInterval(1, SECONDS).until {
+                    def job = sql """
+                        SELECT Status, SucceedTaskCount
+                        FROM jobs("type"="insert")
+                        WHERE Name = 'test_streaming_insert_job_s3_once'
+                          AND ExecuteType = 'STREAMING'
+                    """
+                    def rows = sql """SELECT COUNT(*) FROM 
test_streaming_insert_job_s3_once_tbl"""
+                    log.info("S3 ONCE job: ${job}, row count: ${rows}")
+                    job.size() == 1
+                            && job.get(0).get(0) == "FINISHED"
+                            && job.get(0).get(1).toString() == "2"
+                            && rows.get(0).get(0).toString() == "20"
+                }
+    } catch (Exception ex) {
+        def showJob = sql """
+            SELECT * FROM jobs("type"="insert")
+            WHERE Name = 'test_streaming_insert_job_s3_once'
+        """
+        def showTask = sql """
+            SELECT * FROM tasks("type"="insert")
+            WHERE JobName = 'test_streaming_insert_job_s3_once'
+        """
+        log.info("show job: " + showJob)
+        log.info("show task: " + showTask)
+        throw ex
+    }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to