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

JNSimba pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/branch-4.1 by this push:
     new 9e432943f3e branch-4.1: [fix](streaming-job) Stabilize lag and TVF 
pause-resume tests #68168 (#68400)
9e432943f3e is described below

commit 9e432943f3e5f1ecc2da41a55165b710183d58fd
Author: github-actions[bot] 
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Wed Sep 23 17:04:48 2026 +0800

    branch-4.1: [fix](streaming-job) Stabilize lag and TVF pause-resume tests 
#68168 (#68400)
    
    Cherry-picked from #68168
    
    Co-authored-by: wudi <[email protected]>
---
 .../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 e07caa6edd1..520ad9c1005 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