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]