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]