This is an automated email from the ASF dual-hosted git repository.
renqs pushed a commit to branch release-3.1
in repository https://gitbox.apache.org/repos/asf/flink-cdc.git
The following commit(s) were added to refs/heads/release-3.1 by this push:
new aa2ffd9c2 [FLINK-35128][cdc-connector][cdc-base] Re-calculate the
starting changelog offset after the new table added (#3230) (#3257)
aa2ffd9c2 is described below
commit aa2ffd9c2747cde21d09305c2852debb8937b1fb
Author: Hongshun Wang <[email protected]>
AuthorDate: Fri Apr 26 11:29:41 2024 +0800
[FLINK-35128][cdc-connector][cdc-base] Re-calculate the starting changelog
offset after the new table added (#3230) (#3257)
---
.../cdc/connectors/base/source/meta/split/StreamSplit.java | 10 +++++++++-
1 file changed, 9 insertions(+), 1 deletion(-)
diff --git
a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/meta/split/StreamSplit.java
b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/meta/split/StreamSplit.java
index f4143364a..7b0b048c4 100644
---
a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/meta/split/StreamSplit.java
+++
b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/meta/split/StreamSplit.java
@@ -163,10 +163,18 @@ public class StreamSplit extends SourceSplitBase {
// -------------------------------------------------------------------
public static StreamSplit appendFinishedSplitInfos(
StreamSplit streamSplit, List<FinishedSnapshotSplitInfo>
splitInfos) {
+ // re-calculate the starting changelog offset after the new table added
+ Offset startingOffset = streamSplit.getStartingOffset();
+ for (FinishedSnapshotSplitInfo splitInfo : splitInfos) {
+ if (splitInfo.getHighWatermark().isBefore(startingOffset)) {
+ startingOffset = splitInfo.getHighWatermark();
+ }
+ }
splitInfos.addAll(streamSplit.getFinishedSnapshotSplitInfos());
+
return new StreamSplit(
streamSplit.splitId,
- streamSplit.getStartingOffset(),
+ startingOffset,
streamSplit.getEndingOffset(),
splitInfos,
streamSplit.getTableSchemas(),