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 776d8758b0c [fix](streaming-job) Stabilize lag and TVF pause-resume 
tests (#68168)
776d8758b0c is described below

commit 776d8758b0c15395a516bb7cd136d7f635a55f8e
Author: wudi <[email protected]>
AuthorDate: Wed Sep 23 10:03:37 2026 +0800

    [fix](streaming-job) Stabilize lag and TVF pause-resume tests (#68168)
    
    ### What problem does this PR solve?
    
    Problem Summary:
    
    This PR addresses three unstable streaming job regression cases:
    
    1. `test_streaming_postgres_job_lag`
    2. `test_streaming_mysql_job_lag`
    3. `test_streaming_job_cdc_stream_postgres_pause_resume`
    
    The first two lag cases assumed that a source write immediately after
    job creation would always be captured by an `offset=latest` reader. If
    the reader initialized after that write, no event remained to initialize
    the lag metrics. They now generate source updates only while lag metrics
    are unavailable, crossing the reader initialization window without fixed
    sleeps or production behavior changes.
    
    The third case exposed a production reader lifecycle race. After pause
    and resume, a canceled job-driven `cdc_stream` TVF request could overlap
    the successor task. Both requests reused the job-scoped `SourceReader`,
    allowing the old request to poll or close the successor's fetcher and
    causing a row inserted during pause to be missed. Each job-driven TVF
    request now claims a fresh reader instance, displaced tasks stop reading
    or publishing offsets, and request cleanup releases only its captured
    reader while preserving source-side resources such as the PostgreSQL
    replication slot. Standalone TVF cleanup and FROM-TO reader reuse remain
    unchanged.
---
 .../insert/streaming/StreamingInsertTask.java      |  2 +
 .../offset/jdbc/JdbcTvfSourceOffsetProvider.java   |  5 ++
 .../CdcStreamTableValuedFunction.java              |  1 +
 .../streaming/StreamingInsertTaskAuditTest.java    | 13 ++++
 .../org/apache/doris/cdcclient/common/Env.java     | 14 ++--
 .../cdcclient/service/PipelineCoordinator.java     | 70 ++++++++++++-----
 .../org/apache/doris/cdcclient/common/EnvTest.java | 88 ++++++++++++++++++++++
 .../cdc/test_streaming_mysql_job_lag.groovy        | 25 +++---
 .../cdc/test_streaming_postgres_job_lag.groovy     | 34 ++++-----
 9 files changed, 202 insertions(+), 50 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTask.java
 
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTask.java
index da74592287c..f6589d730c2 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTask.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTask.java
@@ -42,6 +42,7 @@ import org.apache.doris.qe.AuditLogHelper;
 import org.apache.doris.qe.ConnectContext;
 import org.apache.doris.qe.QueryState;
 import org.apache.doris.qe.StmtExecutor;
+import org.apache.doris.tablefunction.CdcStreamTableValuedFunction;
 import org.apache.doris.tablefunction.S3TableValuedFunction;
 import org.apache.doris.thrift.TCell;
 import org.apache.doris.thrift.TRow;
@@ -92,6 +93,7 @@ public class StreamingInsertTask extends 
AbstractStreamingTask {
         this.originTvfProps = originTvfProps;
         this.cloudCluster = cloudCluster;
         this.auditEnabled = 
S3TableValuedFunction.NAME.equalsIgnoreCase(offsetProvider.getSourceType());
+        this.noRetry = 
CdcStreamTableValuedFunction.NAME.equalsIgnoreCase(offsetProvider.getSourceType());
     }
 
     @Override
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 e7324015d93..9d0c88bbdcc 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
@@ -96,6 +96,11 @@ public class JdbcTvfSourceOffsetProvider extends 
JdbcSourceOffsetProvider {
         super();
     }
 
+    @Override
+    public String getSourceType() {
+        return CdcStreamTableValuedFunction.NAME;
+    }
+
     /** Initializes provider state from TVF properties; called every schedule 
tick. */
     @Override
     public void ensureInitialized(Long jobId, Map<String, String> 
originTvfProps) throws JobException {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/tablefunction/CdcStreamTableValuedFunction.java
 
b/fe/fe-core/src/main/java/org/apache/doris/tablefunction/CdcStreamTableValuedFunction.java
index e4323cef82a..f6df6595cd8 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/tablefunction/CdcStreamTableValuedFunction.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/tablefunction/CdcStreamTableValuedFunction.java
@@ -45,6 +45,7 @@ import java.util.Map;
 import java.util.UUID;
 
 public class CdcStreamTableValuedFunction extends 
ExternalFileTableValuedFunction {
+    public static final String NAME = "cdc_stream";
     private static final ObjectMapper objectMapper = new ObjectMapper();
     private static final String URI = 
"http://127.0.0.1:CDC_CLIENT_PORT/api/fetchRecordStream";;
     private static final String ENABLE_CDC_CLIENT_KEY = "enable_cdc_client";
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTaskAuditTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTaskAuditTest.java
index 07cbbfce52b..0a9e917bdfc 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTaskAuditTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTaskAuditTest.java
@@ -67,6 +67,19 @@ public class StreamingInsertTaskAuditTest {
             + "\"s3.secret_key\" = \"private-value\", "
             + "\"enclose\" = \"\\\"\")";
 
+    @Test
+    public void testCdcTaskDisablesInPlaceRetry() {
+        StreamingInsertTask cdcTask = new StreamingInsertTask(
+                1L, 2L, "", new JdbcTvfSourceOffsetProvider(), "test_db", null,
+                Collections.emptyMap(), UserIdentity.ROOT, null);
+        StreamingInsertTask s3Task = new StreamingInsertTask(
+                1L, 3L, "", new S3SourceOffsetProvider(), "test_db", null,
+                Collections.emptyMap(), UserIdentity.ROOT, null);
+
+        Assertions.assertTrue(cdcTask.isNoRetry());
+        Assertions.assertFalse(s3Task.isNoRetry());
+    }
+
     @Test
     public void testS3RunSubmitsAuditEvent() throws Exception {
         AuditEvent auditEvent = runS3Task(null);
diff --git 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/common/Env.java
 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/common/Env.java
index f6d4e5d1243..e93243126ba 100644
--- 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/common/Env.java
+++ 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/common/Env.java
@@ -21,6 +21,7 @@ import org.apache.doris.cdcclient.source.factory.DataSource;
 import org.apache.doris.cdcclient.source.factory.SourceReaderFactory;
 import org.apache.doris.cdcclient.source.reader.AbstractCdcSourceReader;
 import org.apache.doris.cdcclient.source.reader.SourceReader;
+import org.apache.doris.job.cdc.request.FetchRecordRequest;
 import org.apache.doris.job.cdc.request.JobBaseConfig;
 import org.apache.doris.job.cdc.request.WriteRecordRequest;
 
@@ -140,12 +141,15 @@ public class Env {
         try {
             JobContext context = jobContexts.get(jobId);
             if (context != null
-                    && jobConfig instanceof WriteRecordRequest
-                    && ((WriteRecordRequest) jobConfig).isRebuildReader()) {
-                // FE declared the previous task abnormal: swap in a fresh 
reader instance so the
-                // old task's thread can never reach the new fetcher.
+                    && (jobConfig instanceof FetchRecordRequest
+                            || (jobConfig instanceof WriteRecordRequest
+                                    && ((WriteRecordRequest) 
jobConfig).isRebuildReader()))) {
+                // Swap in a fresh reader instance so the old task's thread
+                // can never reach the new fetcher.
                 LOG.info(
-                        "Rebuild reader for job {} on FE request, discard 
current instance", jobId);
+                        "Rebuild reader for job {} task {}, discard current 
instance",
+                        jobId,
+                        taskId);
                 jobContexts.remove(jobId);
                 staleReader = context.reader;
                 staleConfig = context.jobConfig != null ? context.jobConfig : 
jobConfig;
diff --git 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/service/PipelineCoordinator.java
 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/service/PipelineCoordinator.java
index 9c133e1f172..c3ea478e7a9 100644
--- 
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/service/PipelineCoordinator.java
+++ 
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/service/PipelineCoordinator.java
@@ -114,7 +114,6 @@ public class PipelineCoordinator {
     /** return data for http_file_reader */
     public StreamingResponseBody fetchRecordStream(FetchRecordRequest 
fetchReq) throws Exception {
         SourceReader sourceReader;
-        SplitReadResult readResult;
         try {
             LOG.info(
                     "Fetch record request with meta {}, jobId={}, taskId={}",
@@ -128,9 +127,23 @@ public class PipelineCoordinator {
                 LOG.info("Generated meta for job {}: {}", fetchReq.getJobId(), 
meta);
             }
 
-            sourceReader = Env.getCurrentEnv().getReader(fetchReq, 
!isLong(fetchReq.getJobId()));
+            sourceReader =
+                    isJobDrivenTvf(fetchReq.getJobId())
+                            ? Env.getCurrentEnv().getReaderAndClaim(fetchReq, 
fetchReq.getTaskId())
+                            : Env.getCurrentEnv().getReader(fetchReq, true);
+        } catch (Exception ex) {
+            throw new CommonException(ex);
+        }
+
+        SplitReadResult readResult;
+        try {
             readResult = sourceReader.prepareAndSubmitSplit(fetchReq);
         } catch (Exception ex) {
+            try {
+                closeTvfReader(fetchReq, sourceReader);
+            } catch (Exception cleanupEx) {
+                ex.addSuppressed(cleanupEx);
+            }
             throw new CommonException(ex);
         }
 
@@ -144,10 +157,27 @@ public class PipelineCoordinator {
                         fetchReq.getTaskId(),
                         ex);
                 throw new StreamException(ex);
+            } finally {
+                closeTvfReader(fetchReq, sourceReader);
             }
         };
     }
 
+    private void closeTvfReader(FetchRecordRequest request, SourceReader 
sourceReader) {
+        Env env = Env.getCurrentEnv();
+        if (isJobDrivenTvf(request.getJobId())) {
+            env.detachReaderIfOwner(request.getJobId(), request.getTaskId());
+            // Release only this request's instance and keep the PG slot for 
the next task.
+            sourceReader.release(request);
+        } else {
+            try {
+                sourceReader.close(request);
+            } finally {
+                env.close(request.getJobId());
+            }
+        }
+    }
+
     private void buildStreamRecords(
             SourceReader sourceReader,
             FetchRecordRequest fetchRecord,
@@ -170,6 +200,7 @@ public class PipelineCoordinator {
                     fetchRecord.getTaskId(),
                     isSnapshotSplit);
             while (!shouldStop) {
+                checkTvfReaderOwner(fetchRecord);
                 Iterator<SourceRecord> recordIterator = 
sourceReader.pollRecords();
                 if (!recordIterator.hasNext()) {
                     Thread.sleep(100);
@@ -233,29 +264,29 @@ public class PipelineCoordinator {
         }
 
         List<Map<String, String>> offsetMeta = extractOffsetMeta(sourceReader, 
readResult);
+        checkTvfReaderOwner(fetchRecord);
         if (StringUtils.isNotEmpty(fetchRecord.getTaskId())) {
             taskOffsetCache.put(fetchRecord.getTaskId(), offsetMeta);
         }
-        // Convention: standalone TVF uses a UUID jobId; job-driven TVF will 
use a numeric Long
-        // jobId (set via rewriteTvfParams). When the job-driven path is 
implemented,
-        // rewriteTvfParams must inject the job's Long jobId into the TVF 
properties
-        // so that generateParams() can read it, keeping isLong() correct.
-        // TODO: replace isLong() with an explicit field in FetchRecordRequest
-        // once the job-driven TVF path is fully implemented.
-        if (!isLong(fetchRecord.getJobId())) {
-            // TVF requires closing the window after each execution,
-            // while PG requires dropping the slot.
-            sourceReader.close(fetchRecord);
-            // Clean up the job context so it does not accumulate in 
Env.jobContexts.
-            // Each TVF call uses a fresh UUID job ID, so without this the map 
grows unboundedly.
-            Env.getCurrentEnv().close(fetchRecord.getJobId());
+    }
+
+    private void checkTvfReaderOwner(FetchRecordRequest request) {
+        if (isJobDrivenTvf(request.getJobId())
+                && !Env.getCurrentEnv().isOwner(request.getJobId(), 
request.getTaskId())) {
+            throw new IllegalStateException(
+                    String.format(
+                            "TVF reader released or replaced for job %s task 
%s",
+                            request.getJobId(), request.getTaskId()));
         }
     }
 
-    private boolean isLong(String s) {
-        if (s == null || s.isEmpty()) return false;
+    // Convention: standalone TVF uses a UUID jobId; job-driven TVF uses the 
job's numeric Long
+    // jobId, injected into the TVF properties by rewriteTvfParams and read by 
generateParams().
+    // TODO: replace this jobId-based check with an explicit field in 
FetchRecordRequest.
+    private boolean isJobDrivenTvf(String jobId) {
+        if (jobId == null || jobId.isEmpty()) return false;
         try {
-            Long.parseLong(s);
+            Long.parseLong(jobId);
             return true;
         } catch (NumberFormatException e) {
             return false;
@@ -288,7 +319,8 @@ public class PipelineCoordinator {
     public RecordWithMeta fetchRecords(FetchRecordRequest fetchRecordRequest) 
throws Exception {
         SourceReader sourceReader =
                 Env.getCurrentEnv()
-                        .getReader(fetchRecordRequest, 
!isLong(fetchRecordRequest.getJobId()));
+                        .getReader(
+                                fetchRecordRequest, 
!isJobDrivenTvf(fetchRecordRequest.getJobId()));
         SplitReadResult readResult = 
sourceReader.prepareAndSubmitSplit(fetchRecordRequest);
         return buildRecordResponse(sourceReader, fetchRecordRequest, 
readResult);
     }
diff --git 
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/common/EnvTest.java
 
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/common/EnvTest.java
index 685beb29f8a..0bdfcd8378e 100644
--- 
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/common/EnvTest.java
+++ 
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/common/EnvTest.java
@@ -17,12 +17,91 @@
 
 package org.apache.doris.cdcclient.common;
 
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
 import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+import org.apache.doris.cdcclient.exception.CommonException;
+import org.apache.doris.cdcclient.service.PipelineCoordinator;
+import org.apache.doris.cdcclient.source.reader.SourceReader;
+import org.apache.doris.job.cdc.request.FetchRecordRequest;
+import org.apache.doris.job.cdc.request.WriteRecordRequest;
 
 import org.junit.jupiter.api.Test;
 
+import java.lang.reflect.Method;
+import java.util.Collections;
+
 class EnvTest {
 
+    @Test
+    void tvfRequestReplacesReaderOwnedByPreviousTask() throws Exception {
+        Env env = Env.getCurrentEnv();
+        PipelineCoordinator coordinator = new PipelineCoordinator(1);
+        Method closeTvfReader =
+                PipelineCoordinator.class.getDeclaredMethod(
+                        "closeTvfReader", FetchRecordRequest.class, 
SourceReader.class);
+        closeTvfReader.setAccessible(true);
+        FetchRecordRequest firstRequest = tvfRequest("68168001", "first");
+        FetchRecordRequest secondRequest = tvfRequest("68168001", "second");
+        SourceReader first = env.getReaderAndClaim(firstRequest, 
firstRequest.getTaskId());
+        try {
+            SourceReader second = env.getReaderAndClaim(secondRequest, 
secondRequest.getTaskId());
+            assertNotSame(first, second);
+            closeTvfReader.invoke(coordinator, firstRequest, first);
+            assertSame(second, 
env.getReaderIfPresent(secondRequest.getJobId()));
+            closeTvfReader.invoke(coordinator, secondRequest, second);
+            assertNull(env.getReaderIfPresent(secondRequest.getJobId()));
+        } finally {
+            SourceReader reader = 
env.getReaderIfPresent(secondRequest.getJobId());
+            if (reader != null) {
+                reader.release(secondRequest);
+            }
+            env.close(secondRequest.getJobId());
+        }
+    }
+
+    @Test
+    void failedTvfPreparationRemovesClaimedReader() {
+        Env env = Env.getCurrentEnv();
+        FetchRecordRequest request = tvfRequest("68168002", "task");
+        try {
+            CommonException exception =
+                    assertThrows(
+                            CommonException.class,
+                            () -> new 
PipelineCoordinator(1).fetchRecordStream(request));
+            assertEquals("miss meta offset", 
exception.getCause().getMessage());
+            assertNull(env.getReaderIfPresent(request.getJobId()));
+        } finally {
+            SourceReader reader = env.getReaderIfPresent(request.getJobId());
+            if (reader != null) {
+                reader.release(request);
+            }
+            env.close(request.getJobId());
+        }
+    }
+
+    @Test
+    void fromToStillReusesReaderUnlessRebuildRequested() {
+        Env env = Env.getCurrentEnv();
+        WriteRecordRequest request = new WriteRecordRequest();
+        request.setJobId("from-to-reader-reuse");
+        request.setDataSource("POSTGRES");
+        request.setConfig(Collections.emptyMap());
+        SourceReader first = env.getReaderAndClaim(request, "first");
+        try {
+            assertSame(first, env.getReaderAndClaim(request, "second"));
+            request.setRebuildReader(true);
+            assertNotSame(first, env.getReaderAndClaim(request, "third"));
+        } finally {
+            env.getReaderIfPresent(request.getJobId()).release(request);
+            first.release(request);
+            env.close(request.getJobId());
+        }
+    }
+
     @Test
     void getReaderIfPresentReturnsNullForUnknownJob() {
         // An off-target releaseReader RPC must be a no-op, never create a 
reader -> peek returns null.
@@ -34,4 +113,13 @@ class EnvTest {
         // Stale release for an unknown job (no lock/context) must be a no-op.
         assertNull(Env.getCurrentEnv().detachReaderIfOwner("no-such-job-id", 
"t1"));
     }
+
+    private FetchRecordRequest tvfRequest(String jobId, String taskId) {
+        FetchRecordRequest request = new FetchRecordRequest();
+        request.setJobId(jobId);
+        request.setTaskId(taskId);
+        request.setDataSource("POSTGRES");
+        request.setConfig(Collections.emptyMap());
+        return request;
+    }
 }
diff --git 
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_lag.groovy
 
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_lag.groovy
index ac17b4a82a3..a29d8f13fd4 100644
--- 
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_lag.groovy
+++ 
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_lag.groovy
@@ -76,16 +76,23 @@ suite("test_streaming_mysql_job_lag",
                     .pollInterval(1, SECONDS).until({
                         def jobInfo = sql """ select SucceedTaskCount, 
LagBytes, LastSourceEventTimestamp from jobs("type"="insert") where Name = 
'${jobName}' and ExecuteType='STREAMING' """
                         log.info("jobInfo: " + jobInfo)
-                        if (jobInfo.size() != 1 || 
Integer.parseInt(jobInfo[0][0] as String) < 1) {
-                            return false
+                        if (jobInfo.size() == 1 && 
Integer.parseInt(jobInfo[0][0] as String) >= 1) {
+                            def lagValue = jobInfo[0][1] as String
+                            def sourceEventTime = jobInfo[0][2] as String
+                            log.info("lag value: " + lagValue)
+                            if (lagValue != null && lagValue != ""
+                                    && lagValue.isLong() && 
Long.parseLong(lagValue) >= 0
+                                    && sourceEventTime != null && 
sourceEventTime.isLong()
+                                    && Long.parseLong(sourceEventTime) > 0) {
+                                return true
+                            }
                         }
-                        def lagValue = jobInfo[0][1] as String
-                        def sourceEventTime = jobInfo[0][2] as String
-                        log.info("lag value: " + lagValue)
-                        return lagValue != null && lagValue != ""
-                                && lagValue.isLong() && 
Long.parseLong(lagValue) >= 0
-                                && sourceEventTime != null && 
sourceEventTime.isLong()
-                                && Long.parseLong(sourceEventTime) > 0
+                        // Keep generating binlog events until the 
latest-offset reader is ready.
+                        connect("root", "123456", 
"jdbc:mysql://${externalEnvIp}:${mysql_port}") {
+                            sql """UPDATE ${mysqlDb}.${mysqlTable} SET age = 
age + 1
+                                    WHERE name = 'Alice'"""
+                        }
+                        return false
                     })
 
             sql "PAUSE JOB where jobname = '${jobName}'"
diff --git 
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_postgres_job_lag.groovy
 
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_postgres_job_lag.groovy
index d66dc10f50b..6f12da8c4b3 100644
--- 
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_postgres_job_lag.groovy
+++ 
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_postgres_job_lag.groovy
@@ -69,14 +69,6 @@ suite("test_streaming_postgres_job_lag",
             """
 
         try {
-            // Wait until the offset=latest baseline is committed before 
writing incremental data.
-            Awaitility.await().atMost(120, SECONDS)
-                    .pollInterval(1, SECONDS).until({
-                        def jobInfo = sql """ select SucceedTaskCount from 
jobs("type"="insert")
-                                where Name = '${jobName}' and 
ExecuteType='STREAMING' """
-                        return jobInfo.size() == 1 && 
Integer.parseInt(jobInfo[0][0] as String) >= 1
-                    })
-
             // insert incremental data to trigger WAL consumption
             connect("${pgUser}", "${pgPassword}", 
"jdbc:postgresql://${externalEnvIp}:${pg_port}/${pgDB}") {
                 sql """INSERT INTO ${pgDB}.${pgSchema}.${pgTable} (name, age) 
VALUES ('Bob', 20)"""
@@ -89,16 +81,24 @@ suite("test_streaming_postgres_job_lag",
                                 from jobs("type"="insert")
                                 where Name = '${jobName}' and 
ExecuteType='STREAMING' """
                         log.info("jobInfo: " + jobInfo)
-                        if (jobInfo.size() != 1 || 
Integer.parseInt(jobInfo[0][0] as String) < 1) {
-                            return false
+                        if (jobInfo.size() == 1 && 
Integer.parseInt(jobInfo[0][0] as String) >= 1) {
+                            def lagValue = jobInfo[0][1] as String
+                            def sourceEventTime = jobInfo[0][2] as String
+                            log.info("lag value: " + lagValue)
+                            if (lagValue != null && lagValue != ""
+                                    && lagValue.isLong() && 
Long.parseLong(lagValue) >= 0
+                                    && sourceEventTime != null && 
sourceEventTime.isLong()
+                                    && Long.parseLong(sourceEventTime) > 0) {
+                                return true
+                            }
+                        }
+                        // Keep generating WAL until the latest-offset reader 
is ready.
+                        connect("${pgUser}", "${pgPassword}",
+                                
"jdbc:postgresql://${externalEnvIp}:${pg_port}/${pgDB}") {
+                            sql """UPDATE ${pgDB}.${pgSchema}.${pgTable} SET 
age = age + 1
+                                    WHERE name = 'Alice'"""
                         }
-                        def lagValue = jobInfo[0][1] as String
-                        def sourceEventTime = jobInfo[0][2] as String
-                        log.info("lag value: " + lagValue)
-                        return lagValue != null && lagValue != ""
-                                && lagValue.isLong() && 
Long.parseLong(lagValue) >= 0
-                                && sourceEventTime != null && 
sourceEventTime.isLong()
-                                && Long.parseLong(sourceEventTime) > 0
+                        return false
                     })
         } catch (Exception ex) {
             def showjob = sql """select * from jobs("type"="insert") where 
Name='${jobName}'"""


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

Reply via email to