Skip to content

Commit 48bfd04

Browse files
branch-4.1: [Improve](job) update flink cdc version for streaming job #62212 (#62430)
Cherry-picked from #62212 Co-authored-by: wudi <wudi@selectdb.com>
1 parent 6820b0e commit 48bfd04

File tree

3 files changed

+4
-63
lines changed

3 files changed

+4
-63
lines changed

fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProvider.java

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -431,9 +431,7 @@ public void replayIfNeed(StreamingInsertJob job) throws JobException {
431431
protected List<SnapshotSplit> recalculateRemainingSplits(
432432
Map<String, Map<String, Map<String, String>>> chunkHighWatermarkMap,
433433
Map<String, List<SnapshotSplit>> snapshotSplits) {
434-
if (this.finishedSplits == null) {
435-
this.finishedSplits = new ArrayList<>();
436-
}
434+
this.finishedSplits = new ArrayList<>();
437435
for (Map.Entry<String, Map<String, Map<String, String>>> entry : chunkHighWatermarkMap.entrySet()) {
438436
String tableId = entry.getKey();
439437
Map<String, Map<String, String>> splitIdToHighWatermark = entry.getValue();

fs_brokers/cdc_client/pom.xml

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,7 @@ under the License.
6363
<properties>
6464
<spring.version>3.5.7</spring.version>
6565
<flink.version>1.20.3</flink.version>
66+
<flink.cdc.version>3.6.0-1.20</flink.cdc.version>
6667
<slf4j.version>1.7.36</slf4j.version>
6768
<doris.version>1.2-SNAPSHOT</doris.version>
6869
<spotless.version>2.43.0</spotless.version>
@@ -118,12 +119,12 @@ under the License.
118119
<dependency>
119120
<groupId>org.apache.flink</groupId>
120121
<artifactId>flink-connector-mysql-cdc</artifactId>
121-
<version>3.5.0</version>
122+
<version>${flink.cdc.version}</version>
122123
</dependency>
123124
<dependency>
124125
<groupId>org.apache.flink</groupId>
125126
<artifactId>flink-connector-postgres-cdc</artifactId>
126-
<version>3.5.0</version>
127+
<version>${flink.cdc.version}</version>
127128
</dependency>
128129
<dependency>
129130
<groupId>org.apache.flink</groupId>

fs_brokers/cdc_client/src/main/java/org/apache/flink/cdc/connectors/postgres/source/PostgresConnectionPoolFactory.java

Lines changed: 0 additions & 58 deletions
This file was deleted.

0 commit comments

Comments
 (0)