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 46afe021abf [improvement](streamingjob) change streaming job source
log lag in bytes (#66409)
46afe021abf is described below
commit 46afe021abfd53da12da29259f6ce949ae661c18
Author: wudi <[email protected]>
AuthorDate: Thu Aug 27 10:52:35 2026 +0800
[improvement](streamingjob) change streaming job source log lag in bytes
(#66409)
### What problem does this PR solve?
Issue Number: N/A
Related PR: https://github.com/apache/doris/pull/62269
Problem Summary:
Streaming job lag was calculated from the current time and the timestamp
of the last consumed event. The value kept increasing when the source
was idle, so it could not distinguish source inactivity from an actual
synchronization backlog.
This change reports source-log backlog in bytes. MySQL-compatible
sources calculate the remaining bytes from the committed binlog file and
position to the current end offset. PostgreSQL calculates the WAL byte
difference from the replication slot confirmed position to the current
WAL position. Snapshot-only jobs report `-1`, and transient lag
calculation failures retain the latest successful observation.
The `jobs("type"="insert")` output renames `Lag` to `LagBytes` and adds
`LastSourceEventTimestamp`. The new field is the Unix timestamp in
seconds of the latest source binlog or WAL event recorded in the job's
committed offset. MySQL-compatible sources use `ts_sec`, while
PostgreSQL converts `ts_usec` to seconds. It is empty before a source
event timestamp is available. The existing formatted
`LastTaskSuccessTime` remains unchanged.
The per-job Prometheus metrics now expose:
- `doris_fe_streaming_job_per_job_lag_bytes`
- `doris_fe_streaming_job_per_job_last_source_event_timestamp_seconds`
- `doris_fe_streaming_job_per_job_last_task_success_time_seconds`
The two timestamp metrics are Unix seconds. A value of `0` means that
the corresponding source event or successful task has not been observed
yet.
---
.../job/cdc/request/FetchEndOffsetRequest.java | 43 ++++
.../job/cdc/response/FetchEndOffsetResult.java | 32 +++
.../doris/job/extensions/insert/InsertJob.java | 4 +-
.../insert/streaming/StreamingInsertJob.java | 31 +--
.../doris/job/offset/SourceOffsetProvider.java | 23 +-
.../job/offset/jdbc/JdbcSourceOffsetProvider.java | 130 +++++++----
.../java/org/apache/doris/metric/MetricRepo.java | 44 +++-
.../streaming/StreamingInsertJobLagTest.java | 76 +++++++
.../jdbc/JdbcSourceOffsetProviderLagTest.java | 243 +++++++++++++++++++++
.../java/org/apache/doris/metric/MetricsTest.java | 54 +++++
.../cdcclient/controller/ClientController.java | 10 +-
.../cdcclient/source/reader/SourceReader.java | 6 +-
.../reader/mysql/MySqlBinlogLagCalculator.java | 77 +++++++
.../source/reader/mysql/MySqlSourceReader.java | 63 +++++-
.../reader/postgres/PostgresSourceReader.java | 89 ++++++--
.../cdcclient/itcase/CdcClientWriteHarness.java | 11 +
.../cdcclient/itcase/MySqlVersionSmokeITCase.java | 1 +
.../itcase/PostgresVersionSmokeITCase.java | 1 +
.../reader/mysql/MySqlBinlogLagCalculatorTest.java | 165 ++++++++++++++
.../cdc/test_streaming_mysql_job_lag.groovy | 24 +-
.../cdc/test_streaming_mysql_job_metrics.groovy | 33 ++-
.../cdc/test_streaming_oceanbase_job.groovy | 12 +
.../cdc/test_streaming_postgres_job_lag.groovy | 18 +-
23 files changed, 1083 insertions(+), 107 deletions(-)
diff --git
a/fe/fe-common/src/main/java/org/apache/doris/job/cdc/request/FetchEndOffsetRequest.java
b/fe/fe-common/src/main/java/org/apache/doris/job/cdc/request/FetchEndOffsetRequest.java
new file mode 100644
index 00000000000..579d148caa1
--- /dev/null
+++
b/fe/fe-common/src/main/java/org/apache/doris/job/cdc/request/FetchEndOffsetRequest.java
@@ -0,0 +1,43 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.job.cdc.request;
+
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+import lombok.NoArgsConstructor;
+
+import java.util.Collections;
+import java.util.Map;
+
+@Data
+@NoArgsConstructor
+@EqualsAndHashCode(callSuper = true)
+public class FetchEndOffsetRequest extends JobBaseConfig {
+ private Map<String, String> referenceOffset;
+
+ public FetchEndOffsetRequest(
+ String jobId,
+ String dataSource,
+ Map<String, String> config,
+ String frontendAddress,
+ Map<String, String> referenceOffset) {
+ super(jobId, dataSource, config, frontendAddress);
+ // Its presence tells the CDC client that this FE accepts the
lag-aware response.
+ this.referenceOffset = referenceOffset == null ?
Collections.emptyMap() : referenceOffset;
+ }
+}
diff --git
a/fe/fe-common/src/main/java/org/apache/doris/job/cdc/response/FetchEndOffsetResult.java
b/fe/fe-common/src/main/java/org/apache/doris/job/cdc/response/FetchEndOffsetResult.java
new file mode 100644
index 00000000000..660c4661959
--- /dev/null
+++
b/fe/fe-common/src/main/java/org/apache/doris/job/cdc/response/FetchEndOffsetResult.java
@@ -0,0 +1,32 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.job.cdc.response;
+
+import lombok.AllArgsConstructor;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+import java.util.Map;
+
+@Data
+@AllArgsConstructor
+@NoArgsConstructor
+public class FetchEndOffsetResult {
+ private Map<String, String> endOffset;
+ private long lagBytes = -1;
+}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/InsertJob.java
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/InsertJob.java
index 3a2a930b878..490562503dd 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/InsertJob.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/InsertJob.java
@@ -105,7 +105,8 @@ public class InsertJob extends AbstractJob<InsertTask,
Map<Object, Object>> impl
.add(new Column("LoadStatistic", ScalarType.createStringType()))
.add(new Column("ErrorMsg", ScalarType.createStringType()))
.add(new Column("JobRuntimeMsg", ScalarType.createStringType()))
- .add(new Column("Lag", ScalarType.createStringType()))
+ .add(new Column("LagBytes", ScalarType.createStringType()))
+ .add(new Column("LastSourceEventTimestamp",
ScalarType.createStringType()))
.add(new Column("LastTaskSuccessTime",
ScalarType.createStringType()))
.build();
@@ -572,6 +573,7 @@ public class InsertJob extends AbstractJob<InsertTask,
Map<Object, Object>> impl
trow.addToColumnValue(new TCell().setStringVal(failMsg == null ?
FeConstants.null_string : failMsg.getMsg()));
trow.addToColumnValue(new
TCell().setStringVal(FeConstants.null_string));
trow.addToColumnValue(new
TCell().setStringVal(FeConstants.null_string));
+ trow.addToColumnValue(new
TCell().setStringVal(FeConstants.null_string));
trow.addToColumnValue(new TCell().setStringVal(lastTaskSuccessTime > 0
? TimeUtils.longToTimeString(lastTaskSuccessTime) :
FeConstants.null_string));
return trow;
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
index affb6d6c599..6b9c2234012 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
@@ -1000,21 +1000,19 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
}
public String getLag() {
- return offsetProvider != null ? offsetProvider.getLag() : "";
+ return offsetProvider != null ? offsetProvider.getLag() : "-1";
}
- // Numeric lag for metrics. Returns -1 when lag is not applicable (S3,
snapshot phase)
- // or unparseable, so dashboards can filter N/A jobs via lag >= 0.
- public long getLagSeconds() {
- String lagStr = getLag();
- if (lagStr == null || lagStr.isEmpty()) {
- return -1L;
- }
- try {
- return Long.parseLong(lagStr);
- } catch (NumberFormatException e) {
- return -1L;
- }
+ public long getLagBytes() {
+ return offsetProvider != null ? offsetProvider.getLagBytes() : -1;
+ }
+
+ public long getLastSourceEventTimestampSeconds() {
+ return offsetProvider != null ?
offsetProvider.getLastSourceEventTimestampSeconds() : 0;
+ }
+
+ public long getLastTaskSuccessTimeSeconds() {
+ return lastTaskSuccessTime / 1000L;
}
/**
@@ -1077,6 +1075,7 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
if (StringUtils.isNotEmpty(inputStreamProps.getOffsetProperty())) {
Offset offset =
validateOffset(inputStreamProps.getOffsetProperty());
this.offsetProvider.updateOffset(offset);
+ this.offsetProvider.resetLag();
this.offsetProviderPersist = offsetProvider.getPersistInfo();
log.info("modifyPropertiesInternal: offset updated to {}, job {}",
inputStreamProps.getOffsetProperty(), getJobId());
@@ -1174,8 +1173,10 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
? "" : GsonUtils.GSON.toJson(failureReason)));
trow.addToColumnValue(new TCell().setStringVal(jobRuntimeMsg == null
? "" : jobRuntimeMsg));
- trow.addToColumnValue(new TCell().setStringVal(
- offsetProvider != null ? offsetProvider.getLag() : ""));
+ trow.addToColumnValue(new TCell().setStringVal(getLag()));
+ long lastSourceEventTimestampSeconds =
getLastSourceEventTimestampSeconds();
+ trow.addToColumnValue(new
TCell().setStringVal(lastSourceEventTimestampSeconds > 0
+ ? String.valueOf(lastSourceEventTimestampSeconds) : ""));
trow.addToColumnValue(new TCell().setStringVal(lastTaskSuccessTime > 0
? TimeUtils.longToTimeString(lastTaskSuccessTime) : ""));
return trow;
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/SourceOffsetProvider.java
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/SourceOffsetProvider.java
index 58c83eac632..be70fdf83bb 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/SourceOffsetProvider.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/SourceOffsetProvider.java
@@ -206,15 +206,22 @@ public interface SourceOffsetProvider {
return 0;
}
- /**
- * Get the lag of the data source in seconds.
- * For CDC sources, lag = (now - last consumed event timestamp) in seconds.
- *
- * @return lag in seconds as string, empty string if not applicable
- */
+ /** Get the latest successfully observed source-log lag in bytes, or -1
before any observation. */
+ default long getLagBytes() {
+ return -1;
+ }
+
+ /** Get the source event timestamp at the committed offset as Unix
seconds, or 0 if unavailable. */
+ default long getLastSourceEventTimestampSeconds() {
+ return 0;
+ }
+
+ /** Discard a lag value that was calculated from an offset explicitly
replaced by the user. */
+ default void resetLag() {}
+
+ /** Get source lag as a numeric string for SHOW output. */
default String getLag() {
- return "";
+ return String.valueOf(getLagBytes());
}
}
-
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProvider.java
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProvider.java
index 4eb5fad0890..64752144853 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProvider.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProvider.java
@@ -24,8 +24,10 @@ import org.apache.doris.httpv2.entity.ResponseBody;
import org.apache.doris.httpv2.rest.RestApiStatusCode;
import org.apache.doris.job.cdc.DataSourceConfigKeys;
import org.apache.doris.job.cdc.request.CompareOffsetRequest;
+import org.apache.doris.job.cdc.request.FetchEndOffsetRequest;
import org.apache.doris.job.cdc.request.FetchTableSplitsRequest;
import org.apache.doris.job.cdc.request.JobBaseConfig;
+import org.apache.doris.job.cdc.response.FetchEndOffsetResult;
import org.apache.doris.job.cdc.split.AbstractSourceSplit;
import org.apache.doris.job.cdc.split.BinlogSplit;
import org.apache.doris.job.cdc.split.SnapshotSplit;
@@ -106,6 +108,8 @@ public class JdbcSourceOffsetProvider implements
SourceOffsetProvider {
volatile boolean hasMoreData = true;
+ transient volatile long lagBytes = -1;
+
transient volatile String cloudCluster;
// Route fetchEndOffset/compareOffset to the bound BE (synced from job,
not persisted).
@@ -304,8 +308,13 @@ public class JdbcSourceOffsetProvider implements
SourceOffsetProvider {
@Override
public void fetchRemoteMeta(Map<String, String> properties) throws
Exception {
Backend backend = StreamingJobUtils.selectBackend(cloudCluster,
boundBackendId);
- JobBaseConfig requestParams =
- new JobBaseConfig(getJobId().toString(), sourceType.name(),
sourceProperties, getFrontendAddress());
+ FetchEndOffsetRequest requestParams =
+ new FetchEndOffsetRequest(
+ getJobId().toString(),
+ sourceType.name(),
+ sourceProperties,
+ getFrontendAddress(),
+ getLagReferenceOffset());
InternalService.PRequestCdcClientRequest request =
InternalService.PRequestCdcClientRequest.newBuilder()
.setApi("/api/fetchEndOffset")
.setParams(new Gson().toJson(requestParams)).build();
@@ -322,14 +331,17 @@ public class JdbcSourceOffsetProvider implements
SourceOffsetProvider {
"Failed to get end offset from backend," +
result.getStatus().getErrorMsgs(0) + ", response: "
+ result.getResponse());
}
- Map<String, String> newEndOffset = parseCdcResponseData(
- result.getResponse(), new TypeReference<Map<String,
String>>() {});
+ FetchEndOffsetResult fetchResult =
parseFetchEndOffsetResponse(result.getResponse());
+ Map<String, String> newEndOffset = fetchResult.getEndOffset();
synchronized (splitsLock) {
// null→value also counts as a change: upstream may have
advanced while fetch was blocked.
if (endBinlogOffset == null ||
!endBinlogOffset.equals(newEndOffset)) {
hasMoreData = true;
}
endBinlogOffset = newEndOffset;
+ if (!isSnapshotOnlyMode()) {
+ updateLagBytes(fetchResult.getLagBytes());
+ }
}
} catch (TimeoutException te) {
log.warn("cdc_client RPC timeout api=/api/fetchEndOffset jobId={}
backend={}:{} timeout_sec={}",
@@ -342,6 +354,30 @@ public class JdbcSourceOffsetProvider implements
SourceOffsetProvider {
}
}
+ Map<String, String> getLagReferenceOffset() {
+ if (isSnapshotOnlyMode()) {
+ return null;
+ }
+ synchronized (splitsLock) {
+ if (currentOffset != null && !currentOffset.snapshotSplit()) {
+ BinlogSplit binlogSplit = (BinlogSplit)
currentOffset.getSplits().get(0);
+ if (MapUtils.isNotEmpty(binlogSplit.getStartingOffset())) {
+ return new HashMap<>(binlogSplit.getStartingOffset());
+ }
+ }
+ if (sourceType == DataSourceType.POSTGRES) {
+ // PostgreSQL can use the replication slot's confirmed flush
LSN during snapshot.
+ return null;
+ }
+ return finishedSplits.stream()
+ .map(SnapshotSplit::getHighWatermark)
+ .filter(MapUtils::isNotEmpty)
+ .findFirst()
+ .map(HashMap::new)
+ .orElse(null);
+ }
+ }
+
@Override
public boolean hasMoreDataToConsume() {
if (currentOffset == null) {
@@ -945,6 +981,23 @@ public class JdbcSourceOffsetProvider implements
SourceOffsetProvider {
}
}
+ FetchEndOffsetResult parseFetchEndOffsetResponse(String response) throws
JobException {
+ JsonNode data = parseCdcResponseData(response, new
TypeReference<JsonNode>() {});
+ if (data == null) {
+ throw new JobException(response);
+ }
+ try {
+ if (data.has("endOffset")) {
+ return objectMapper.convertValue(data,
FetchEndOffsetResult.class);
+ }
+ Map<String, String> endOffset = objectMapper.convertValue(
+ data, new TypeReference<Map<String, String>>() {});
+ return new FetchEndOffsetResult(endOffset, -1);
+ } catch (IllegalArgumentException exception) {
+ throw new JobException(response);
+ }
+ }
+
protected boolean checkNeedSplitChunks(Map<String, String>
sourceProperties) {
String startMode = sourceProperties.get(DataSourceConfigKeys.OFFSET);
if (startMode == null) {
@@ -960,45 +1013,46 @@ public class JdbcSourceOffsetProvider implements
SourceOffsetProvider {
}
@Override
- public String getLag() {
- if (currentOffset == null || currentOffset.snapshotSplit()) {
- return "";
- }
- // Source is idle (last task consumed no data), report zero lag
- if (!hasMoreData) {
- return "0";
- }
- BinlogSplit binlogSplit = (BinlogSplit)
currentOffset.getSplits().get(0);
- Map<String, String> offsetMap = binlogSplit.getStartingOffset();
- if (MapUtils.isEmpty(offsetMap)) {
- return "";
- }
- long eventTimeMs = extractEventTimeMs(offsetMap);
- if (eventTimeMs <= 0) {
- return "0";
- }
- long lagSec = (System.currentTimeMillis() - eventTimeMs) / 1000;
- return String.valueOf(Math.max(lagSec, 0));
+ public long getLagBytes() {
+ return lagBytes;
}
- /**
- * Extract event timestamp in milliseconds from binlog offset map.
- * MySQL: ts_sec (seconds), PostgreSQL: ts_usec (microseconds).
- */
- protected long extractEventTimeMs(Map<String, String> offsetMap) {
- try {
- String tsSec = offsetMap.get("ts_sec");
- if (tsSec != null) {
- return Long.parseLong(tsSec) * 1000;
+ @Override
+ public long getLastSourceEventTimestampSeconds() {
+ synchronized (splitsLock) {
+ if (currentOffset == null || currentOffset.snapshotSplit()) {
+ return 0;
}
- String tsUsec = offsetMap.get("ts_usec");
- if (tsUsec != null) {
- return Long.parseLong(tsUsec) / 1000;
+ BinlogSplit binlogSplit = (BinlogSplit)
currentOffset.getSplits().get(0);
+ Map<String, String> offsetMap = binlogSplit.getStartingOffset();
+ if (MapUtils.isEmpty(offsetMap)) {
+ return 0;
}
- } catch (NumberFormatException e) {
- log.warn("Failed to parse event timestamp from offset: {}",
offsetMap, e);
+ try {
+ String timestampSeconds = offsetMap.get("ts_sec");
+ if (timestampSeconds != null) {
+ return Math.max(Long.parseLong(timestampSeconds), 0);
+ }
+ String timestampMicros = offsetMap.get("ts_usec");
+ if (timestampMicros != null) {
+ return Math.max(Long.parseLong(timestampMicros) /
1_000_000, 0);
+ }
+ } catch (NumberFormatException e) {
+ log.warn("Failed to parse source event timestamp from offset:
{}", offsetMap, e);
+ }
+ return 0;
+ }
+ }
+
+ @Override
+ public void resetLag() {
+ lagBytes = -1;
+ }
+
+ void updateLagBytes(long fetchedLagBytes) {
+ if (fetchedLagBytes >= 0) {
+ lagBytes = fetchedLagBytes;
}
- return -1;
}
@Override
diff --git a/fe/fe-core/src/main/java/org/apache/doris/metric/MetricRepo.java
b/fe/fe-core/src/main/java/org/apache/doris/metric/MetricRepo.java
index 0bc8b133a1b..3c1aeb6a037 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/metric/MetricRepo.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/metric/MetricRepo.java
@@ -100,7 +100,11 @@ public final class MetricRepo {
public static final String STREAMING_JOB_PER_JOB_FILTERED_ROWS =
"streaming_job_per_job_filtered_rows";
public static final String STREAMING_JOB_PER_JOB_SUCCEED_TASK_COUNT =
"streaming_job_per_job_succeed_task_count";
public static final String STREAMING_JOB_PER_JOB_FAILED_TASK_COUNT =
"streaming_job_per_job_failed_task_count";
- public static final String STREAMING_JOB_PER_JOB_LAG =
"streaming_job_per_job_lag";
+ public static final String STREAMING_JOB_PER_JOB_LAG_BYTES =
"streaming_job_per_job_lag_bytes";
+ public static final String
STREAMING_JOB_PER_JOB_LAST_SOURCE_EVENT_TIMESTAMP_SECONDS =
+ "streaming_job_per_job_last_source_event_timestamp_seconds";
+ public static final String
STREAMING_JOB_PER_JOB_LAST_TASK_SUCCESS_TIME_SECONDS =
+ "streaming_job_per_job_last_task_success_time_seconds";
public static final String ROUTINE_LOAD_PER_JOB_TOTAL_ROWS =
"routine_load_per_job_total_rows";
public static final String ROUTINE_LOAD_PER_JOB_ERROR_ROWS =
"routine_load_per_job_error_rows";
public static final String ROUTINE_LOAD_PER_JOB_RECEIVED_BYTES =
"routine_load_per_job_received_bytes";
@@ -1378,7 +1382,9 @@ public final class MetricRepo {
DORIS_METRIC_REGISTER.removeMetrics(STREAMING_JOB_PER_JOB_FILTERED_ROWS);
DORIS_METRIC_REGISTER.removeMetrics(STREAMING_JOB_PER_JOB_SUCCEED_TASK_COUNT);
DORIS_METRIC_REGISTER.removeMetrics(STREAMING_JOB_PER_JOB_FAILED_TASK_COUNT);
- DORIS_METRIC_REGISTER.removeMetrics(STREAMING_JOB_PER_JOB_LAG);
+ DORIS_METRIC_REGISTER.removeMetrics(STREAMING_JOB_PER_JOB_LAG_BYTES);
+
DORIS_METRIC_REGISTER.removeMetrics(STREAMING_JOB_PER_JOB_LAST_SOURCE_EVENT_TIMESTAMP_SECONDS);
+
DORIS_METRIC_REGISTER.removeMetrics(STREAMING_JOB_PER_JOB_LAST_TASK_SUCCESS_TIME_SECONDS);
try {
List<org.apache.doris.job.base.AbstractJob> jobs =
@@ -1463,16 +1469,44 @@ public final class MetricRepo {
DORIS_METRIC_REGISTER.addMetrics(failedTaskCount);
GaugeMetric<Long> lag = new GaugeMetric<Long>(
- STREAMING_JOB_PER_JOB_LAG, MetricUnit.SECONDS,
- "per job lag in seconds of streaming job, -1 means
N/A") {
+ STREAMING_JOB_PER_JOB_LAG_BYTES, MetricUnit.BYTES,
+ "latest successfully observed source log lag in bytes,
-1 means no valid observation") {
@Override
public Long getValue() {
- return sJob.getLagSeconds();
+ return sJob.getLagBytes();
}
};
lag.addLabel(new MetricLabel("job_id", jobId))
.addLabel(new MetricLabel("job_name", jobName));
DORIS_METRIC_REGISTER.addMetrics(lag);
+
+ long lastSourceEventTimestampSeconds =
sJob.getLastSourceEventTimestampSeconds();
+ GaugeMetric<Long> lastSourceEventTimestamp = new
GaugeMetric<Long>(
+
STREAMING_JOB_PER_JOB_LAST_SOURCE_EVENT_TIMESTAMP_SECONDS, MetricUnit.SECONDS,
+ "timestamp of the latest source binlog or WAL event
recorded in the job's committed offset"
+ + " as Unix seconds, 0 means unavailable") {
+ @Override
+ public Long getValue() {
+ return lastSourceEventTimestampSeconds;
+ }
+ };
+ lastSourceEventTimestamp.addLabel(new MetricLabel("job_id",
jobId))
+ .addLabel(new MetricLabel("job_name", jobName));
+ DORIS_METRIC_REGISTER.addMetrics(lastSourceEventTimestamp);
+
+ long lastTaskSuccessTimeSeconds =
sJob.getLastTaskSuccessTimeSeconds();
+ GaugeMetric<Long> lastTaskSuccessTime = new GaugeMetric<Long>(
+ STREAMING_JOB_PER_JOB_LAST_TASK_SUCCESS_TIME_SECONDS,
MetricUnit.SECONDS,
+ "timestamp of the latest successful task completion as
Unix seconds,"
+ + " 0 means no successful task") {
+ @Override
+ public Long getValue() {
+ return lastTaskSuccessTimeSeconds;
+ }
+ };
+ lastTaskSuccessTime.addLabel(new MetricLabel("job_id", jobId))
+ .addLabel(new MetricLabel("job_name", jobName));
+ DORIS_METRIC_REGISTER.addMetrics(lastTaskSuccessTime);
}
} catch (Throwable t) {
LOG.warn("failed to update streaming job per-job metrics", t);
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobLagTest.java
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobLagTest.java
new file mode 100644
index 00000000000..e1a64931c8b
--- /dev/null
+++
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobLagTest.java
@@ -0,0 +1,76 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.job.extensions.insert.streaming;
+
+import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.job.base.JobExecutionConfiguration;
+import org.apache.doris.job.base.TimerDefinition;
+import org.apache.doris.job.cdc.split.BinlogSplit;
+import org.apache.doris.job.offset.jdbc.JdbcOffset;
+import org.apache.doris.job.offset.jdbc.JdbcSourceOffsetProvider;
+
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+
+public class StreamingInsertJobLagTest {
+
+ @Test
+ public void testLastSourceEventTimestampUsesOffsetProvider() {
+ StreamingInsertJob job =
Deencapsulation.newInstance(StreamingInsertJob.class);
+ Assert.assertEquals(0L, job.getLastSourceEventTimestampSeconds());
+
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ Map<String, String> committedOffset = new HashMap<>();
+ committedOffset.put("file", "mysql-bin.000001");
+ committedOffset.put("pos", "100");
+ committedOffset.put("ts_sec", "1787039821");
+ provider.setCurrentOffset(new JdbcOffset(Collections.singletonList(new
BinlogSplit(committedOffset))));
+ Deencapsulation.setField(job, "offsetProvider", provider);
+
+ Assert.assertEquals(1787039821L,
job.getLastSourceEventTimestampSeconds());
+ job.setLastTaskSuccessTime(1787039821123L);
+ Assert.assertEquals(1787039821L, job.getLastTaskSuccessTimeSeconds());
+ }
+
+ @Test
+ public void testExplicitOffsetChangeInvalidatesLastObservedLag() throws
Exception {
+ StreamingInsertJob job =
Deencapsulation.newInstance(StreamingInsertJob.class);
+ Map<String, String> properties = new HashMap<>();
+ properties.put(StreamingJobProperties.MAX_INTERVAL_SECOND_PROPERTY,
"10");
+ Deencapsulation.setField(job, "properties", properties);
+ Deencapsulation.setField(job, "jobProperties", new
StreamingJobProperties(properties));
+
+ JobExecutionConfiguration configuration = new
JobExecutionConfiguration();
+ configuration.setTimerDefinition(new TimerDefinition());
+ Deencapsulation.setField(job, "jobConfig", configuration);
+
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ provider.setLagBytes(4096);
+ Deencapsulation.setField(job, "offsetProvider", provider);
+
+ Map<String, String> alterProperties = new HashMap<>();
+ alterProperties.put(StreamingJobProperties.OFFSET_PROPERTY,
"{\"lsn\":\"200\"}");
+ Deencapsulation.invoke(job, "modifyPropertiesInternal",
alterProperties);
+
+ Assert.assertEquals(-1, provider.getLagBytes());
+ }
+}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderLagTest.java
b/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderLagTest.java
new file mode 100644
index 00000000000..fbcf94af12a
--- /dev/null
+++
b/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderLagTest.java
@@ -0,0 +1,243 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.job.offset.jdbc;
+
+import org.apache.doris.job.cdc.DataSourceConfigKeys;
+import org.apache.doris.job.cdc.request.FetchEndOffsetRequest;
+import org.apache.doris.job.cdc.response.FetchEndOffsetResult;
+import org.apache.doris.job.cdc.split.BinlogSplit;
+import org.apache.doris.job.cdc.split.SnapshotSplit;
+import org.apache.doris.job.common.DataSourceType;
+import org.apache.doris.job.exception.JobException;
+
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+
+public class JdbcSourceOffsetProviderLagTest {
+
+ @Test
+ public void testPostgresSnapshotDoesNotSendFeReferenceOffset() {
+ JdbcSourceOffsetProvider provider = provider(DataSourceType.POSTGRES,
DataSourceConfigKeys.OFFSET_INITIAL);
+ provider.finishedSplits.add(snapshotSplit("split-1", offset("lsn",
"300")));
+ provider.finishedSplits.add(snapshotSplit("split-2", offset("lsn",
"100")));
+
+ Assert.assertNull(provider.getLagReferenceOffset());
+ }
+
+ @Test
+ public void testPostgresIncrementalPhaseUsesCommittedOffset() {
+ JdbcSourceOffsetProvider provider = provider(DataSourceType.POSTGRES,
DataSourceConfigKeys.OFFSET_INITIAL);
+ Map<String, String> committedOffset = offset("lsn", "700");
+ provider.currentOffset =
+ new JdbcOffset(Collections.singletonList(new
BinlogSplit(committedOffset)));
+
+ Assert.assertEquals(committedOffset, provider.getLagReferenceOffset());
+ }
+
+ @Test
+ public void testInitialSnapshotUsesFirstCommittedMysqlHighWatermark() {
+ JdbcSourceOffsetProvider provider = provider(DataSourceType.MYSQL,
DataSourceConfigKeys.OFFSET_INITIAL);
+ provider.finishedSplits.add(snapshotSplit("split-1",
mysqlOffset("mysql-bin.000010", 100)));
+ provider.finishedSplits.add(snapshotSplit("split-2",
mysqlOffset("mysql-bin.000009", 900)));
+ provider.finishedSplits.add(snapshotSplit("split-3",
mysqlOffset("mysql-bin.000010", 50)));
+
+ Assert.assertEquals(mysqlOffset("mysql-bin.000010", 100),
provider.getLagReferenceOffset());
+ }
+
+ @Test
+ public void testIncrementalPhaseUsesCommittedBinlogOffset() {
+ JdbcSourceOffsetProvider provider = provider(DataSourceType.MYSQL,
DataSourceConfigKeys.OFFSET_INITIAL);
+ provider.finishedSplits.add(snapshotSplit("split-1",
mysqlOffset("mysql-bin.000001", 100)));
+ Map<String, String> committedOffset = mysqlOffset("mysql-bin.000002",
250);
+ provider.currentOffset =
+ new JdbcOffset(Collections.singletonList(new
BinlogSplit(committedOffset)));
+
+ Assert.assertEquals(committedOffset, provider.getLagReferenceOffset());
+ }
+
+ @Test
+ public void
testRestoredSnapshotToBinlogTransitionUsesSnapshotHighWatermark() {
+ JdbcSourceOffsetProvider provider = provider(DataSourceType.MYSQL,
DataSourceConfigKeys.OFFSET_INITIAL);
+ provider.finishedSplits.add(snapshotSplit("split-1",
mysqlOffset("mysql-bin.000003", 300)));
+ provider.finishedSplits.add(snapshotSplit("split-2",
mysqlOffset("mysql-bin.000001", 100)));
+ provider.currentOffset =
+ new JdbcOffset(Collections.singletonList(new BinlogSplit()));
+
+ Assert.assertEquals(mysqlOffset("mysql-bin.000003", 300),
provider.getLagReferenceOffset());
+ }
+
+ @Test
+ public void testSnapshotOnlyDoesNotExposeSourceLogLag() {
+ JdbcSourceOffsetProvider provider = provider(DataSourceType.POSTGRES,
DataSourceConfigKeys.OFFSET_SNAPSHOT);
+ provider.finishedSplits.add(snapshotSplit("split-1", offset("lsn",
"100")));
+
+ Assert.assertNull(provider.getLagReferenceOffset());
+ Assert.assertEquals("-1", provider.getLag());
+ }
+
+ @Test
+ public void testLagIsAlwaysNumeric() {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+
+ Assert.assertEquals("-1", provider.getLag());
+ provider.setLagBytes(4096);
+ Assert.assertEquals("4096", provider.getLag());
+ Assert.assertEquals(4096, provider.getLagBytes());
+ }
+
+ @Test
+ public void testMysqlLastSourceEventTimestampUsesCommittedOffsetSeconds() {
+ JdbcSourceOffsetProvider provider = provider(DataSourceType.MYSQL,
DataSourceConfigKeys.OFFSET_INITIAL);
+ Map<String, String> committedOffset = mysqlOffset("mysql-bin.000002",
250);
+ committedOffset.put("ts_sec", "1787039821");
+ provider.currentOffset =
+ new JdbcOffset(Collections.singletonList(new
BinlogSplit(committedOffset)));
+
+ Assert.assertEquals(1787039821L,
provider.getLastSourceEventTimestampSeconds());
+ }
+
+ @Test
+ public void
testPostgresLastSourceEventTimestampConvertsCommittedOffsetMicrosToSeconds() {
+ JdbcSourceOffsetProvider provider = provider(DataSourceType.POSTGRES,
DataSourceConfigKeys.OFFSET_INITIAL);
+ Map<String, String> committedOffset = new HashMap<>();
+ committedOffset.put("lsn", "700");
+ committedOffset.put("ts_usec", "1787039821987654");
+ provider.currentOffset =
+ new JdbcOffset(Collections.singletonList(new
BinlogSplit(committedOffset)));
+
+ Assert.assertEquals(1787039821L,
provider.getLastSourceEventTimestampSeconds());
+ }
+
+ @Test
+ public void testPostgresInitialOffsetTimestampIsUnavailable() {
+ JdbcSourceOffsetProvider provider = provider(DataSourceType.POSTGRES,
DataSourceConfigKeys.OFFSET_LATEST);
+ Map<String, String> committedOffset = new HashMap<>();
+ committedOffset.put("lsn", "0");
+ committedOffset.put("ts_usec", String.valueOf(Long.MIN_VALUE));
+ provider.currentOffset =
+ new JdbcOffset(Collections.singletonList(new
BinlogSplit(committedOffset)));
+
+ Assert.assertEquals(0L, provider.getLastSourceEventTimestampSeconds());
+ }
+
+ @Test
+ public void
testLastSourceEventTimestampUnavailableBeforeCommittedBinlogTimestamp() {
+ JdbcSourceOffsetProvider provider = provider(DataSourceType.MYSQL,
DataSourceConfigKeys.OFFSET_INITIAL);
+ provider.currentOffset = new JdbcOffset(Collections.singletonList(
+ snapshotSplit("split-1", mysqlOffset("mysql-bin.000001",
100))));
+
+ Assert.assertEquals(0L, provider.getLastSourceEventTimestampSeconds());
+
+ provider.currentOffset = new JdbcOffset(Collections.singletonList(
+ new BinlogSplit(mysqlOffset("mysql-bin.000002", 250))));
+ Assert.assertEquals(0L, provider.getLastSourceEventTimestampSeconds());
+ }
+
+ @Test
+ public void testUnavailableLagDoesNotOverwriteLastSuccessfulValue() {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ provider.setLagBytes(4096);
+
+ provider.updateLagBytes(-1);
+
+ Assert.assertEquals(4096, provider.getLagBytes());
+ }
+
+ @Test
+ public void testSuccessfulLagReplacesLastSuccessfulValue() {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ provider.setLagBytes(4096);
+
+ provider.updateLagBytes(2048);
+
+ Assert.assertEquals(2048, provider.getLagBytes());
+ }
+
+ @Test
+ public void testParseFetchEndOffsetResponse() throws JobException {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ String response = "{\"code\":0,\"msg\":\"Success\",\"data\":{"
+ + "\"endOffset\":{\"lsn\":\"200\"},\"lagBytes\":4096}}";
+
+ FetchEndOffsetResult result =
provider.parseFetchEndOffsetResponse(response);
+
+ Assert.assertEquals(offset("lsn", "200"), result.getEndOffset());
+ Assert.assertEquals(4096, result.getLagBytes());
+ }
+
+ @Test
+ public void testParseLegacyFetchEndOffsetResponse() throws JobException {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ String response =
"{\"code\":0,\"msg\":\"Success\",\"data\":{\"lsn\":\"200\"}}";
+
+ FetchEndOffsetResult result =
provider.parseFetchEndOffsetResponse(response);
+
+ Assert.assertEquals(offset("lsn", "200"), result.getEndOffset());
+ Assert.assertEquals(-1, result.getLagBytes());
+ }
+
+ @Test
+ public void testParseFetchEndOffsetResponseWithoutLag() throws
JobException {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ String response = "{\"code\":0,\"msg\":\"Success\",\"data\":{"
+ + "\"endOffset\":{\"lsn\":\"200\"}}}";
+
+ FetchEndOffsetResult result =
provider.parseFetchEndOffsetResponse(response);
+
+ Assert.assertEquals(offset("lsn", "200"), result.getEndOffset());
+ Assert.assertEquals(-1, result.getLagBytes());
+ }
+
+ @Test
+ public void
testFetchEndOffsetRequestUsesEmptyReferenceOffsetAsCapabilitySignal() {
+ FetchEndOffsetRequest request =
+ new FetchEndOffsetRequest("123", "POSTGRES",
Collections.emptyMap(), null, null);
+
+ Assert.assertNotNull(request.getReferenceOffset());
+ Assert.assertTrue(request.getReferenceOffset().isEmpty());
+ }
+
+ private static JdbcSourceOffsetProvider provider(DataSourceType type,
String startupMode) {
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ provider.setSourceType(type);
+ provider.setJobId(123L);
+ provider.setSourceProperties(
+ Collections.singletonMap(DataSourceConfigKeys.OFFSET,
startupMode));
+ return provider;
+ }
+
+ private static SnapshotSplit snapshotSplit(String splitId, Map<String,
String> highWatermark) {
+ return new SnapshotSplit(splitId, "db.table", Arrays.asList("id"),
null, null, highWatermark);
+ }
+
+ private static Map<String, String> mysqlOffset(String file, long position)
{
+ Map<String, String> offset = new HashMap<>();
+ offset.put("file", file);
+ offset.put("pos", String.valueOf(position));
+ return offset;
+ }
+
+ private static Map<String, String> offset(String key, String value) {
+ return Collections.singletonMap(key, value);
+ }
+}
diff --git a/fe/fe-core/src/test/java/org/apache/doris/metric/MetricsTest.java
b/fe/fe-core/src/test/java/org/apache/doris/metric/MetricsTest.java
index 698de8f104b..36d272f6984 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/metric/MetricsTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/metric/MetricsTest.java
@@ -23,7 +23,15 @@ import org.apache.doris.cloud.JobWarmUpStats;
import org.apache.doris.common.Config;
import org.apache.doris.common.FeConstants;
import org.apache.doris.common.Pair;
+import org.apache.doris.common.jmockit.Deencapsulation;
import org.apache.doris.common.util.JsonUtil;
+import org.apache.doris.ha.FrontendNodeType;
+import org.apache.doris.job.base.AbstractJob;
+import org.apache.doris.job.cdc.split.BinlogSplit;
+import org.apache.doris.job.extensions.insert.streaming.StreamingInsertJob;
+import org.apache.doris.job.manager.JobManager;
+import org.apache.doris.job.offset.jdbc.JdbcOffset;
+import org.apache.doris.job.offset.jdbc.JdbcSourceOffsetProvider;
import org.apache.doris.metric.Metric.MetricUnit;
import org.apache.doris.monitor.jvm.JvmService;
import org.apache.doris.monitor.jvm.JvmStats;
@@ -40,7 +48,10 @@ import org.junit.Test;
import java.lang.management.GarbageCollectorMXBean;
import java.lang.management.ManagementFactory;
import java.util.Collections;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
@Slf4j
@@ -120,6 +131,49 @@ public class MetricsTest {
}
}
+ @Test
+ public void testStreamingJobTimeAndLagMetrics() {
+ StreamingInsertJob job =
Deencapsulation.newInstance(StreamingInsertJob.class);
+ job.setJobId(1787039821000L);
+ job.setJobName("streaming_metric_job");
+ job.setLastTaskSuccessTime(1787039821123L);
+
+ JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider();
+ provider.setLagBytes(4096);
+ Map<String, String> committedOffset = new HashMap<>();
+ committedOffset.put("file", "mysql-bin.000001");
+ committedOffset.put("pos", "100");
+ committedOffset.put("ts_sec", "1787039800");
+ provider.setCurrentOffset(
+ new JdbcOffset(Collections.singletonList(new
BinlogSplit(committedOffset))));
+ Deencapsulation.setField(job, "offsetProvider", provider);
+
+ Env env = Env.getCurrentEnv();
+ FrontendNodeType originalFeType = Deencapsulation.getField(env,
"feType");
+ Deencapsulation.setField(env, "feType", FrontendNodeType.MASTER);
+ JobManager jobManager = env.getJobManager();
+ ConcurrentHashMap<Long, AbstractJob> jobMap =
Deencapsulation.getField(jobManager, "jobMap");
+ jobMap.put(job.getJobId(), job);
+ try {
+ MetricRepo.updateStreamingJobPerJobMetrics();
+ String metricResult = getPrometheusMetrics();
+
+
Assert.assertTrue(metricResult.contains("doris_fe_streaming_job_per_job_lag_bytes"
+ + "{job_id=\"1787039821000\",
job_name=\"streaming_metric_job\"} 4096"));
+ Assert.assertTrue(metricResult.contains(
+
"doris_fe_streaming_job_per_job_last_source_event_timestamp_seconds"
+ + "{job_id=\"1787039821000\",
job_name=\"streaming_metric_job\"} 1787039800"));
+ Assert.assertTrue(metricResult.contains(
+
"doris_fe_streaming_job_per_job_last_task_success_time_seconds"
+ + "{job_id=\"1787039821000\",
job_name=\"streaming_metric_job\"} 1787039821"));
+
Assert.assertFalse(metricResult.contains("doris_fe_streaming_job_per_job_lag{"));
+ } finally {
+ jobMap.remove(job.getJobId());
+ MetricRepo.updateStreamingJobPerJobMetrics();
+ Deencapsulation.setField(env, "feType", originalFeType);
+ }
+ }
+
@Test
public void testUserQueryMetrics() {
MetricRepo.USER_COUNTER_QUERY_ALL.getOrAdd("test_user").increase(1L);
diff --git
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/controller/ClientController.java
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/controller/ClientController.java
index cc0bac0665b..3f220ef68f5 100644
---
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/controller/ClientController.java
+++
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/controller/ClientController.java
@@ -22,10 +22,12 @@ import org.apache.doris.cdcclient.model.rest.RestResponse;
import org.apache.doris.cdcclient.service.PipelineCoordinator;
import org.apache.doris.cdcclient.source.reader.SourceReader;
import org.apache.doris.job.cdc.request.CompareOffsetRequest;
+import org.apache.doris.job.cdc.request.FetchEndOffsetRequest;
import org.apache.doris.job.cdc.request.FetchRecordRequest;
import org.apache.doris.job.cdc.request.FetchTableSplitsRequest;
import org.apache.doris.job.cdc.request.JobBaseConfig;
import org.apache.doris.job.cdc.request.WriteRecordRequest;
+import org.apache.doris.job.cdc.response.FetchEndOffsetResult;
import org.apache.commons.lang3.exception.ExceptionUtils;
@@ -103,12 +105,16 @@ public class ClientController {
/** Fetch lastest end meta */
@RequestMapping(path = "/api/fetchEndOffset", method = RequestMethod.POST)
- public Object fetchEndOffset(@RequestBody JobBaseConfig jobConfig) {
+ public Object fetchEndOffset(@RequestBody FetchEndOffsetRequest jobConfig)
{
LOG.info("Fetching end offset for job {}", jobConfig.getJobId());
try {
SourceReader reader = Env.getCurrentEnv().getMetaReader(jobConfig);
Env.getCurrentEnv().keepAlive(jobConfig.getJobId());
- return RestResponse.success(reader.getEndOffset(jobConfig));
+ FetchEndOffsetResult result = reader.fetchEndOffset(jobConfig);
+ // Requests from older FEs do not contain referenceOffset and
expect the legacy
+ // response.
+ return RestResponse.success(
+ jobConfig.getReferenceOffset() == null ?
result.getEndOffset() : result);
} catch (Exception ex) {
LOG.error("Failed to fetch end offset, jobId={}",
jobConfig.getJobId(), ex);
return
RestResponse.internalError(ExceptionUtils.getRootCauseMessage(ex));
diff --git
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/SourceReader.java
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/SourceReader.java
index 70577764ed4..95c6d4ff448 100644
---
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/SourceReader.java
+++
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/SourceReader.java
@@ -20,9 +20,11 @@ package org.apache.doris.cdcclient.source.reader;
import org.apache.doris.cdcclient.source.deserialize.DeserializeResult;
import org.apache.doris.cdcclient.source.factory.DataSource;
import org.apache.doris.job.cdc.request.CompareOffsetRequest;
+import org.apache.doris.job.cdc.request.FetchEndOffsetRequest;
import org.apache.doris.job.cdc.request.FetchTableSplitsRequest;
import org.apache.doris.job.cdc.request.JobBaseConfig;
import org.apache.doris.job.cdc.request.JobBaseRecordRequest;
+import org.apache.doris.job.cdc.response.FetchEndOffsetResult;
import org.apache.doris.job.cdc.split.AbstractSourceSplit;
import org.apache.flink.api.connector.source.SourceSplit;
@@ -77,8 +79,8 @@ public interface SourceReader {
/** Finish reading all split records */
void finishSplitRecords();
- /** Get the end offset for the job */
- Map<String, String> getEndOffset(JobBaseConfig jobConfig);
+ /** Get the end offset and latest source-log lag. */
+ FetchEndOffsetResult fetchEndOffset(FetchEndOffsetRequest request);
/** Compare the offsets */
int compareOffset(CompareOffsetRequest compareOffsetRequest);
diff --git
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlBinlogLagCalculator.java
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlBinlogLagCalculator.java
new file mode 100644
index 00000000000..732866ea46d
--- /dev/null
+++
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlBinlogLagCalculator.java
@@ -0,0 +1,77 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.cdcclient.source.reader.mysql;
+
+import java.util.List;
+import java.util.Map;
+
+final class MySqlBinlogLagCalculator {
+ private static final String FILE_KEY = "file";
+ private static final String POSITION_KEY = "pos";
+
+ private MySqlBinlogLagCalculator() {}
+
+ static long calculate(
+ Map<String, String> referenceOffset,
+ Map<String, String> endOffset,
+ List<BinlogFile> binlogFiles) {
+ if (referenceOffset == null || endOffset == null || binlogFiles ==
null) {
+ return -1;
+ }
+ String referenceFile = referenceOffset.get(FILE_KEY);
+ String endFile = endOffset.get(FILE_KEY);
+ String referencePositionValue = referenceOffset.get(POSITION_KEY);
+ if (referenceFile == null || referencePositionValue == null) {
+ // GTID-only startup offsets have no byte position until the
reader advances.
+ return -1;
+ }
+ long referencePosition = Long.parseLong(referencePositionValue);
+ long endPosition = Long.parseLong(endOffset.get(POSITION_KEY));
+ int referenceIndex = indexOf(binlogFiles, referenceFile);
+ int endIndex = indexOf(binlogFiles, endFile);
+ if (referenceIndex < 0
+ || endIndex < referenceIndex
+ || referencePosition < 0
+ || endPosition < 0) {
+ return -1;
+ }
+ if (referenceIndex == endIndex) {
+ return endPosition >= referencePosition ? endPosition -
referencePosition : -1;
+ }
+
+ long lag = binlogFiles.get(referenceIndex).size() - referencePosition;
+ if (lag < 0) {
+ return -1;
+ }
+ for (int i = referenceIndex + 1; i < endIndex; i++) {
+ lag = Math.addExact(lag, binlogFiles.get(i).size());
+ }
+ return Math.addExact(lag, endPosition);
+ }
+
+ private static int indexOf(List<BinlogFile> binlogFiles, String name) {
+ for (int i = 0; i < binlogFiles.size(); i++) {
+ if (binlogFiles.get(i).name().equals(name)) {
+ return i;
+ }
+ }
+ return -1;
+ }
+
+ record BinlogFile(String name, long size) {}
+}
diff --git
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java
index 3fd87c0335e..19380dfe421 100644
---
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java
+++
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java
@@ -28,9 +28,11 @@ import org.apache.doris.cdcclient.utils.ConfigUtil;
import org.apache.doris.cdcclient.utils.SmallFileMgr;
import org.apache.doris.job.cdc.DataSourceConfigKeys;
import org.apache.doris.job.cdc.request.CompareOffsetRequest;
+import org.apache.doris.job.cdc.request.FetchEndOffsetRequest;
import org.apache.doris.job.cdc.request.FetchTableSplitsRequest;
import org.apache.doris.job.cdc.request.JobBaseConfig;
import org.apache.doris.job.cdc.request.JobBaseRecordRequest;
+import org.apache.doris.job.cdc.response.FetchEndOffsetResult;
import org.apache.doris.job.cdc.split.AbstractSourceSplit;
import org.apache.doris.job.cdc.split.BinlogSplit;
import org.apache.doris.job.cdc.split.SnapshotSplit;
@@ -1174,13 +1176,62 @@ public class MySqlSourceReader extends
AbstractCdcSourceReader {
}
@Override
- public Map<String, String> getEndOffset(JobBaseConfig jobConfig) {
- MySqlSourceConfig sourceConfig = getSourceConfig(jobConfig);
+ public FetchEndOffsetResult fetchEndOffset(FetchEndOffsetRequest request) {
+ MySqlSourceConfig sourceConfig = getSourceConfig(request);
try (MySqlConnection jdbc =
DebeziumUtils.createMySqlConnection(sourceConfig)) {
- BinlogOffset binlogOffset =
DebeziumUtils.currentBinlogOffset(jdbc);
- return binlogOffset.getOffset();
- } catch (SQLException ex) {
- throw new RuntimeException(ex);
+ Map<String, String> endOffset =
DebeziumUtils.currentBinlogOffset(jdbc).getOffset();
+ long lagBytes;
+ try {
+ lagBytes = calculateLagBytes(request, endOffset, jdbc);
+ } catch (Exception exception) {
+ lagBytes = -1;
+ LOG.warn(
+ "Failed to calculate source log lag for job {}",
+ request.getJobId(),
+ exception);
+ }
+ return new FetchEndOffsetResult(endOffset, lagBytes);
+ } catch (SQLException exception) {
+ throw new RuntimeException(exception);
+ }
+ }
+
+ private long calculateLagBytes(
+ FetchEndOffsetRequest request, Map<String, String> endOffset,
MySqlConnection jdbc)
+ throws SQLException {
+ if (MapUtils.isEmpty(request.getReferenceOffset())) {
+ return -1;
+ }
+ try (Statement statement = jdbc.connection().createStatement();
+ ResultSet resultSet = statement.executeQuery("SHOW BINARY
LOGS")) {
+ List<MySqlBinlogLagCalculator.BinlogFile> binlogFiles = new
ArrayList<>();
+ while (resultSet.next()) {
+ binlogFiles.add(
+ new MySqlBinlogLagCalculator.BinlogFile(
+ resultSet.getString(1), resultSet.getLong(2)));
+ }
+ String referenceFile = request.getReferenceOffset().get("file");
+ String endFile = endOffset.get("file");
+ if (referenceFile != null && endFile != null) {
+ boolean referenceFileExists =
+ binlogFiles.stream().anyMatch(file ->
file.name().equals(referenceFile));
+ boolean endFileExists =
+ binlogFiles.stream().anyMatch(file ->
file.name().equals(endFile));
+ if (!referenceFileExists || !endFileExists) {
+ LOG.warn(
+ "Cannot calculate source log lag because a binlog
file is unavailable,"
+ + " jobId={}, referenceFile={}
(available={}),"
+ + " endFile={} (available={})",
+ request.getJobId(),
+ referenceFile,
+ referenceFileExists,
+ endFile,
+ endFileExists);
+ return -1;
+ }
+ }
+ return MySqlBinlogLagCalculator.calculate(
+ request.getReferenceOffset(), endOffset, binlogFiles);
}
}
diff --git
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/postgres/PostgresSourceReader.java
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/postgres/PostgresSourceReader.java
index 330f461510b..e52502e5576 100644
---
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/postgres/PostgresSourceReader.java
+++
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/postgres/PostgresSourceReader.java
@@ -26,8 +26,10 @@ import org.apache.doris.cdcclient.utils.ConfigUtil;
import org.apache.doris.cdcclient.utils.SmallFileMgr;
import org.apache.doris.job.cdc.DataSourceConfigKeys;
import org.apache.doris.job.cdc.request.CompareOffsetRequest;
+import org.apache.doris.job.cdc.request.FetchEndOffsetRequest;
import org.apache.doris.job.cdc.request.JobBaseConfig;
import org.apache.doris.job.cdc.request.JobBaseRecordRequest;
+import org.apache.doris.job.cdc.response.FetchEndOffsetResult;
import org.apache.commons.lang3.StringUtils;
import org.apache.commons.lang3.exception.ExceptionUtils;
@@ -55,6 +57,7 @@ import
org.apache.flink.cdc.connectors.postgres.source.utils.PostgresTypeUtils;
import
org.apache.flink.cdc.connectors.postgres.source.utils.TableDiscoveryUtils;
import org.apache.flink.table.types.DataType;
+import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;
@@ -466,28 +469,74 @@ public class PostgresSourceReader extends
JdbcIncrementalSourceReader {
}
}
- /**
- * Why not call dialect.displayCurrentOffset(sourceConfig) ? The
underlying system calls
- * `txid_current()` to advance the WAL log. Here, it's just a query;
retrieving the LSN is
- * sufficient because `PostgresOffset.compare` only compares the LSN.
- */
@Override
- public Map<String, String> getEndOffset(JobBaseConfig jobConfig) {
- PostgresSourceConfig sourceConfig = getSourceConfig(jobConfig);
- try {
- PostgresDialect dialect = new PostgresDialect(sourceConfig);
- try (JdbcConnection jdbcConnection =
dialect.openJdbcConnection(sourceConfig)) {
- PostgresConnection pgConnection = (PostgresConnection)
jdbcConnection;
- Long lsn = pgConnection.currentXLogLocation();
- Map<String, String> offsetMap = new HashMap<>();
- offsetMap.put(SourceInfo.LSN_KEY, lsn.toString());
- offsetMap.put(
- SourceInfo.TIMESTAMP_USEC_KEY,
-
String.valueOf(Conversions.toEpochMicros(Instant.MIN)));
- return offsetMap;
+ public FetchEndOffsetResult fetchEndOffset(FetchEndOffsetRequest request) {
+ PostgresSourceConfig sourceConfig = getSourceConfig(request);
+ PostgresDialect dialect = new PostgresDialect(sourceConfig);
+ String slotName = dialect.getSlotName();
+ try (JdbcConnection jdbcConnection =
dialect.openJdbcConnection(sourceConfig)) {
+ PostgresConnection pgConnection = (PostgresConnection)
jdbcConnection;
+ // displayCurrentOffset() calls txid_current() and advances WAL;
reading the current LSN
+ // is sufficient because PostgresOffset.compare() only compares
LSN.
+ Long lsn = pgConnection.currentXLogLocation();
+ Map<String, String> endOffset = new HashMap<>();
+ endOffset.put(SourceInfo.LSN_KEY, lsn.toString());
+ endOffset.put(
+ SourceInfo.TIMESTAMP_USEC_KEY,
+ String.valueOf(Conversions.toEpochMicros(Instant.MIN)));
+ long lagBytes;
+ try {
+ lagBytes = calculateLagBytes(request, lsn, slotName,
jdbcConnection);
+ } catch (Exception exception) {
+ lagBytes = -1;
+ LOG.warn(
+ "Failed to calculate source log lag for job {}",
+ request.getJobId(),
+ exception);
+ }
+ return new FetchEndOffsetResult(endOffset, lagBytes);
+ } catch (Exception exception) {
+ throw new RuntimeException(exception);
+ }
+ }
+
+ private long calculateLagBytes(
+ FetchEndOffsetRequest request,
+ long endOffset,
+ String slotName,
+ JdbcConnection jdbcConnection)
+ throws SQLException {
+ try (PreparedStatement statement =
+ jdbcConnection
+ .connection()
+ .prepareStatement(
+ "SELECT pg_wal_lsn_diff("
+ + "?::pg_lsn,"
+ + " GREATEST(confirmed_flush_lsn,"
+ + " COALESCE(?::pg_lsn,
confirmed_flush_lsn)))::bigint"
+ + " FROM pg_replication_slots"
+ + " WHERE slot_name = ?")) {
+ String currentOffset = null;
+ Map<String, String> referenceOffset = request.getReferenceOffset();
+ if (referenceOffset != null &&
referenceOffset.get(SourceInfo.LSN_KEY) != null) {
+ currentOffset =
+
Lsn.valueOf(Long.parseLong(referenceOffset.get(SourceInfo.LSN_KEY)))
+ .asString();
+ }
+ statement.setString(1, Lsn.valueOf(endOffset).asString());
+ statement.setString(2, currentOffset);
+ statement.setString(3, slotName);
+ try (ResultSet resultSet = statement.executeQuery()) {
+ if (!resultSet.next()) {
+ throw new SQLException("Replication slot not found: " +
slotName);
+ }
+ long lagBytes = resultSet.getLong(1);
+ if (resultSet.wasNull()) {
+ throw new SQLException(
+ "Replication slot has no confirmed flush LSN: " +
slotName);
+ }
+ return lagBytes >= 0 ? lagBytes : -1;
}
- } catch (Exception ex) {
- throw new RuntimeException(ex);
}
}
diff --git
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/CdcClientWriteHarness.java
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/CdcClientWriteHarness.java
index 07e5b40a811..6a93bca540a 100644
---
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/CdcClientWriteHarness.java
+++
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/CdcClientWriteHarness.java
@@ -22,9 +22,11 @@ import
org.apache.doris.cdcclient.service.PipelineCoordinator;
import org.apache.doris.cdcclient.source.reader.AbstractCdcSourceReader;
import org.apache.doris.cdcclient.source.reader.SourceReader;
import org.apache.doris.job.cdc.DataSourceConfigKeys;
+import org.apache.doris.job.cdc.request.FetchEndOffsetRequest;
import org.apache.doris.job.cdc.request.FetchTableSplitsRequest;
import org.apache.doris.job.cdc.request.JobBaseConfig;
import org.apache.doris.job.cdc.request.WriteRecordRequest;
+import org.apache.doris.job.cdc.response.FetchEndOffsetResult;
import org.apache.doris.job.cdc.split.AbstractSourceSplit;
import org.apache.doris.job.cdc.split.BinlogSplit;
import org.apache.doris.job.cdc.split.SnapshotSplit;
@@ -473,6 +475,15 @@ final class CdcClientWriteHarness implements AutoCloseable
{
return lastTableSchemas;
}
+ long sourceLogLagBytes() throws Exception {
+ Map<String, String> referenceOffset = committedBinlogOffset();
+ FetchEndOffsetRequest request =
+ new FetchEndOffsetRequest(jobId, dataSource, config, null,
referenceOffset);
+ SourceReader reader = openReader();
+ FetchEndOffsetResult result = reader.fetchEndOffset(request);
+ return result.getLagBytes();
+ }
+
@Override
public void close() {
SourceReader reader = Env.getCurrentEnv().getReaderIfPresent(jobId);
diff --git
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/MySqlVersionSmokeITCase.java
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/MySqlVersionSmokeITCase.java
index 66ea089b674..85a67fbc60c 100644
---
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/MySqlVersionSmokeITCase.java
+++
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/MySqlVersionSmokeITCase.java
@@ -125,6 +125,7 @@ class MySqlVersionSmokeITCase {
assertThat(byId.get(3).get(Constants.DORIS_DELETE_SIGN).asInt()).isZero();
assertThat(byId.get(1).get("age").asInt()).isEqualTo(31);
assertThat(byId.get(2).get(Constants.DORIS_DELETE_SIGN).asInt()).isEqualTo(1);
+
assertThat(harness.sourceLogLagBytes()).isGreaterThanOrEqualTo(0);
} finally {
Env.getCurrentEnv().close(jobId);
}
diff --git
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/PostgresVersionSmokeITCase.java
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/PostgresVersionSmokeITCase.java
index c1c1f111143..cb6cbb4615f 100644
---
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/PostgresVersionSmokeITCase.java
+++
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/PostgresVersionSmokeITCase.java
@@ -111,6 +111,7 @@ class PostgresVersionSmokeITCase {
assertThat(byId.get(3).get(Constants.DORIS_DELETE_SIGN).asInt()).isZero();
assertThat(byId.get(1).get("age").asInt()).isEqualTo(31);
assertThat(byId.get(2).get(Constants.DORIS_DELETE_SIGN).asInt()).isEqualTo(1);
+
assertThat(harness.sourceLogLagBytes()).isGreaterThanOrEqualTo(0);
} finally {
Env.getCurrentEnv().close(jobId);
}
diff --git
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlBinlogLagCalculatorTest.java
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlBinlogLagCalculatorTest.java
new file mode 100644
index 00000000000..bf0a12c2e74
--- /dev/null
+++
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlBinlogLagCalculatorTest.java
@@ -0,0 +1,165 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.cdcclient.source.reader.mysql;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+class MySqlBinlogLagCalculatorTest {
+
+ @Test
+ void sameFileLagUsesPositionsInsteadOfActiveFileSize() {
+ Map<String, String> reference = offset("mysql-bin.000003", 1000);
+ Map<String, String> end = offset("mysql-bin.000003", 1250);
+ List<MySqlBinlogLagCalculator.BinlogFile> files =
+ Collections.singletonList(
+ new
MySqlBinlogLagCalculator.BinlogFile("mysql-bin.000003", 5000));
+
+ assertThat(MySqlBinlogLagCalculator.calculate(reference, end,
files)).isEqualTo(250);
+ }
+
+ @Test
+ void crossFileLagIncludesTailIntermediateFilesAndHeadPosition() {
+ Map<String, String> reference = offset("custom-prefix.000001", 1500);
+ Map<String, String> end = offset("custom-prefix.000003", 700);
+ List<MySqlBinlogLagCalculator.BinlogFile> files =
+ Arrays.asList(
+ new
MySqlBinlogLagCalculator.BinlogFile("custom-prefix.000001", 2000),
+ new
MySqlBinlogLagCalculator.BinlogFile("custom-prefix.000002", 3000),
+ new
MySqlBinlogLagCalculator.BinlogFile("custom-prefix.000003", 5000));
+
+ assertThat(MySqlBinlogLagCalculator.calculate(reference, end,
files)).isEqualTo(4200);
+ }
+
+ @Test
+ void equalOffsetsReportCaughtUp() {
+ Map<String, String> offset = offset("mysql-bin.000003", 1250);
+
+ assertThat(
+ MySqlBinlogLagCalculator.calculate(
+ offset,
+ offset,
+ Collections.singletonList(
+ new
MySqlBinlogLagCalculator.BinlogFile(
+ "mysql-bin.000003", 5000))))
+ .isZero();
+ }
+
+ @Test
+ void purgedReferenceFileReportsUnavailable() {
+ assertThat(
+ MySqlBinlogLagCalculator.calculate(
+ offset("mysql-bin.000001", 1500),
+ offset("mysql-bin.000003", 700),
+ Arrays.asList(
+ new
MySqlBinlogLagCalculator.BinlogFile(
+ "mysql-bin.000002", 3000),
+ new
MySqlBinlogLagCalculator.BinlogFile(
+ "mysql-bin.000003", 5000))))
+ .isEqualTo(-1);
+ }
+
+ @Test
+ void missingEndFileReportsUnavailable() {
+ assertThat(
+ MySqlBinlogLagCalculator.calculate(
+ offset("mysql-bin.000001", 1500),
+ offset("mysql-bin.000003", 700),
+ Arrays.asList(
+ new
MySqlBinlogLagCalculator.BinlogFile(
+ "mysql-bin.000001", 2000),
+ new
MySqlBinlogLagCalculator.BinlogFile(
+ "mysql-bin.000002", 3000))))
+ .isEqualTo(-1);
+ }
+
+ @Test
+ void referenceAheadOfHeadReportsUnavailable() {
+ assertThat(
+ MySqlBinlogLagCalculator.calculate(
+ offset("mysql-bin.000003", 1500),
+ offset("mysql-bin.000003", 1250),
+ Collections.singletonList(
+ new
MySqlBinlogLagCalculator.BinlogFile(
+ "mysql-bin.000003", 5000))))
+ .isEqualTo(-1);
+ }
+
+ @Test
+ void gtidOnlyReferenceReportsUnavailableUntilFilePositionIsCommitted() {
+ Map<String, String> reference =
+ Collections.singletonMap("gtids",
"24bc7850-2c16-11ef-a0c9-0242ac120002:1-9");
+
+ assertThat(
+ MySqlBinlogLagCalculator.calculate(
+ reference,
+ offset("mysql-bin.000003", 1250),
+ Collections.singletonList(
+ new
MySqlBinlogLagCalculator.BinlogFile(
+ "mysql-bin.000003", 5000))))
+ .isEqualTo(-1);
+ }
+
+ @Test
+ void overflowIsReportedToTheCaller() {
+ assertThatThrownBy(
+ () ->
+ MySqlBinlogLagCalculator.calculate(
+ offset("mysql-bin.000001", 0),
+ offset("mysql-bin.000003", 1),
+ Arrays.asList(
+ new
MySqlBinlogLagCalculator.BinlogFile(
+ "mysql-bin.000001",
Long.MAX_VALUE),
+ new
MySqlBinlogLagCalculator.BinlogFile(
+ "mysql-bin.000002", 1),
+ new
MySqlBinlogLagCalculator.BinlogFile(
+ "mysql-bin.000003",
1))))
+ .isInstanceOf(ArithmeticException.class);
+ }
+
+ @Test
+ void malformedPositionIsReportedToTheCaller() {
+ Map<String, String> reference = offset("mysql-bin.000003", 1000);
+ reference.put("pos", "not-a-number");
+
+ assertThatThrownBy(
+ () ->
+ MySqlBinlogLagCalculator.calculate(
+ reference,
+ offset("mysql-bin.000003", 1250),
+ Collections.singletonList(
+ new
MySqlBinlogLagCalculator.BinlogFile(
+ "mysql-bin.000003",
5000))))
+ .isInstanceOf(NumberFormatException.class);
+ }
+
+ private static Map<String, String> offset(String file, long position) {
+ Map<String, String> offset = new HashMap<>();
+ offset.put("file", file);
+ offset.put("pos", String.valueOf(position));
+ return offset;
+ }
+}
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 08b4405bcc8..ac17b4a82a3 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
@@ -74,14 +74,34 @@ suite("test_streaming_mysql_job_lag",
// wait for binlog data consumed and lag is available
Awaitility.await().atMost(300, SECONDS)
.pollInterval(1, SECONDS).until({
- def jobInfo = sql """ select SucceedTaskCount, Lag
from jobs("type"="insert") where Name = '${jobName}' and
ExecuteType='STREAMING' """
+ 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
}
def lagValue = jobInfo[0][1] as String
+ def sourceEventTime = jobInfo[0][2] as String
log.info("lag value: " + lagValue)
- return lagValue != null && lagValue != "" &&
lagValue.isNumber()
+ return lagValue != null && lagValue != ""
+ && lagValue.isLong() &&
Long.parseLong(lagValue) >= 0
+ && sourceEventTime != null &&
sourceEventTime.isLong()
+ && Long.parseLong(sourceEventTime) > 0
+ })
+
+ sql "PAUSE JOB where jobname = '${jobName}'"
+ Awaitility.await().atMost(30, SECONDS)
+ .pollInterval(1, SECONDS).until({
+ def jobInfo = sql """ select Status, LagBytes,
LastSourceEventTimestamp from jobs("type"="insert") where Name = '${jobName}'
and ExecuteType='STREAMING' """
+ if (jobInfo.size() != 1 || jobInfo[0][0] != "PAUSED") {
+ return false
+ }
+ def lagValue = jobInfo[0][1] as String
+ def sourceEventTime = jobInfo[0][2] as String
+ log.info("paused lag value: " + lagValue)
+ return lagValue != null && lagValue != ""
+ && lagValue.isLong() &&
Long.parseLong(lagValue) >= 0
+ && sourceEventTime != null &&
sourceEventTime.isLong()
+ && Long.parseLong(sourceEventTime) > 0
})
} catch (Exception ex) {
def showjob = sql """select * from jobs("type"="insert") where
Name='${jobName}'"""
diff --git
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_metrics.groovy
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_metrics.groovy
index a780a16aaaf..a67207a8baa 100644
---
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_metrics.groovy
+++
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_metrics.groovy
@@ -201,21 +201,42 @@ suite("test_streaming_mysql_job_metrics",
metricCount++
}
- def perJobLag = result.find {
- it.tags?.metric ==
"doris_fe_streaming_job_per_job_lag" &&
+ def perJobLagBytes = result.find {
+ it.tags?.metric ==
"doris_fe_streaming_job_per_job_lag_bytes" &&
it.tags?.job_name == "${jobName}"
}
- if (perJobLag != null) {
- log.info("per-job lag: ${perJobLag}".toString())
+ if (perJobLagBytes != null && perJobLagBytes.value != null
+ && new
BigDecimal(perJobLagBytes.value.toString()).signum() >= 0) {
+ log.info("per-job lag_bytes:
${perJobLagBytes}".toString())
metricCount++
}
+ def perJobLastSourceEventTimestamp = result.find {
+ it.tags?.metric ==
+
"doris_fe_streaming_job_per_job_last_source_event_timestamp_seconds" &&
+ it.tags?.job_name == "${jobName}"
+ }
+ if (perJobLastSourceEventTimestamp != null
+ && perJobLastSourceEventTimestamp.value != null
+ && new
BigDecimal(perJobLastSourceEventTimestamp.value.toString()).signum() >= 0) {
+ log.info("per-job last_source_event_timestamp:
${perJobLastSourceEventTimestamp}".toString())
+ metricCount++
+ }
+ def perJobLastTaskSuccessTime = result.find {
+ it.tags?.metric ==
"doris_fe_streaming_job_per_job_last_task_success_time_seconds" &&
+ it.tags?.job_name == "${jobName}"
+ }
+ if (perJobLastTaskSuccessTime != null &&
perJobLastTaskSuccessTime.value != null
+ && new
BigDecimal(perJobLastTaskSuccessTime.value.toString()).signum() > 0) {
+ log.info("per-job last_task_success_time:
${perJobLastTaskSuccessTime}".toString())
+ metricCount++
+ }
}
}
- // 9 streaming_job_* counters + 1 doris_fe_job RUNNING gauge + 6
per-job metrics
- if (metricCount >= 16) {
+ // 9 streaming_job_* counters + 1 doris_fe_job RUNNING gauge + 8
per-job metrics
+ if (metricCount >= 18) {
break
}
diff --git
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_oceanbase_job.groovy
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_oceanbase_job.groovy
index 416c8ddc016..ea492cbc91e 100644
---
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_oceanbase_job.groovy
+++
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_oceanbase_job.groovy
@@ -138,6 +138,18 @@ suite("test_streaming_oceanbase_job",
def status = sql """SELECT Status FROM jobs("type"="insert") WHERE
Name='${jobName}'"""
assert status.size() == 1 && status[0][0] == "RUNNING"
+
+ Awaitility.await().atMost(30, SECONDS).pollInterval(1, SECONDS).until({
+ def jobInfo = sql """SELECT LagBytes FROM jobs("type"="insert")
WHERE Name='${jobName}'"""
+ if (jobInfo.size() != 1) {
+ return false
+ }
+ def lagValue = jobInfo[0][0] as String
+ log.info("OceanBase lag value: " + lagValue)
+ return lagValue != null && lagValue != ""
+ && lagValue.isLong() && Long.parseLong(lagValue) >= 0
+ })
+
sql """DROP JOB IF EXISTS 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 57abd35db2a..d66dc10f50b 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,6 +69,14 @@ 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)"""
@@ -77,14 +85,20 @@ suite("test_streaming_postgres_job_lag",
// wait for binlog data consumed and lag is available
Awaitility.await().atMost(300, SECONDS)
.pollInterval(1, SECONDS).until({
- def jobInfo = sql """ select SucceedTaskCount, Lag
from jobs("type"="insert") where Name = '${jobName}' and
ExecuteType='STREAMING' """
+ 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
}
def lagValue = jobInfo[0][1] as String
+ def sourceEventTime = jobInfo[0][2] as String
log.info("lag value: " + lagValue)
- return lagValue != null && lagValue != "" &&
lagValue.isNumber()
+ return lagValue != null && lagValue != ""
+ && lagValue.isLong() &&
Long.parseLong(lagValue) >= 0
+ && sourceEventTime != null &&
sourceEventTime.isLong()
+ && Long.parseLong(sourceEventTime) > 0
})
} 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]