Skip to content

Commit 51e43b6

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 c57d2c8 commit 51e43b6

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
@@ -772,33 +772,33 @@ public CompletedSnapshot triggerAndWaitSnapshot(TableBucket tableBucket) {
772772
}
773773

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

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

0 commit comments

Comments
 (0)