Skip to content

Commit 1c61b9c

Browse files
[flink][test] fix Flink 2.x open() compatibility and flaky snapshot trigger
- ScanAndCleanFunction: change open(Configuration) to open(OpenContext) for Flink 2.x compatibility (Flink 2.x removed the Configuration overload) - FlussClusterExtension: triggerSnapshot() returns null on no-op instead of failing hard when snapshot ID does not advance (initSnapshot skips when logOffset <= lastSnapshotOffset) - triggerAndWaitSnapshots() silently skips null buckets (original behavior)
1 parent 663fc5f commit 1c61b9c

2 files changed

Lines changed: 20 additions & 18 deletions

File tree

fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/action/orphan/job/ScanAndCleanFunction.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -77,7 +77,9 @@ public ScanAndCleanFunction(long deleteRateLimitPerSecond, Map<String, String> e
7777
}
7878

7979
@Override
80-
public void open(org.apache.flink.configuration.Configuration parameters) {
80+
public void open(org.apache.flink.api.common.functions.OpenContext openContext)
81+
throws Exception {
82+
super.open(openContext);
8183
if (!extraConfigs.isEmpty()) {
8284
FileSystem.initialize(Configuration.fromMap(extraConfigs), null);
8385
}

fluss-server/src/test/java/org/apache/fluss/server/testutils/FlussClusterExtension.java

Lines changed: 17 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -771,33 +771,33 @@ public CompletedSnapshot triggerAndWaitSnapshot(TableBucket tableBucket) {
771771
}
772772

773773
private Long triggerSnapshot(TableBucket tableBucket) {
774-
Long snapshotId = null;
775-
Long nextSnapshotId = null;
776774
for (TabletServer ts : tabletServers.values()) {
777775
ReplicaManager.HostedReplica replica = ts.getReplicaManager().getReplica(tableBucket);
778776
if (replica instanceof ReplicaManager.OnlineReplica) {
779777
Replica r = ((ReplicaManager.OnlineReplica) replica).getReplica();
780778
PeriodicSnapshotManager kvSnapshotManager = r.getKvSnapshotManager();
781779
if (r.isLeader() && kvSnapshotManager != null) {
782-
snapshotId = kvSnapshotManager.currentSnapshotId();
780+
long snapshotId = kvSnapshotManager.currentSnapshotId();
783781
kvSnapshotManager.triggerSnapshot();
784-
nextSnapshotId = kvSnapshotManager.currentSnapshotId();
785-
break;
782+
// triggerSnapshot() submits to guardedExecutor asynchronously.
783+
// Wait for the counter to advance; if it does not, initSnapshot()
784+
// determined there is no new data (logOffset <= lastSnapshotOffset)
785+
// and the trigger was a legitimate no-op — return null.
786+
try {
787+
waitUntil(
788+
() -> kvSnapshotManager.currentSnapshotId() > snapshotId,
789+
Duration.ofSeconds(3),
790+
Duration.ofMillis(50),
791+
"");
792+
} catch (AssertionError e) {
793+
return null;
794+
}
795+
return snapshotId;
786796
}
787797
}
788798
}
789-
790-
if (snapshotId != null) {
791-
if (nextSnapshotId > snapshotId) {
792-
// only there is a new snapshot triggered, we return the snapshot id
793-
return snapshotId;
794-
} else {
795-
return null;
796-
}
797-
} else {
798-
fail("No KV snapshot manager found for table bucket " + tableBucket);
799-
return null;
800-
}
799+
fail("No KV snapshot manager found for table bucket " + tableBucket);
800+
return null;
801801
}
802802

803803
public CompletedSnapshot waitUntilSnapshotFinished(TableBucket tableBucket, long snapshotId) {

0 commit comments

Comments
 (0)