diff --git a/fluss-server/src/test/java/org/apache/fluss/server/testutils/FlussClusterExtension.java b/fluss-server/src/test/java/org/apache/fluss/server/testutils/FlussClusterExtension.java index bc833e9c02b..0cc615482ab 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/testutils/FlussClusterExtension.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/testutils/FlussClusterExtension.java @@ -943,8 +943,21 @@ public void waitUntilPartitionsDropped(TablePath tablePath, List dropped "Fail to wait partitions dropped"); } + /** Wait until the assigned replica is ready as the local leader and return its server id. */ public int waitAndGetLeader(TableBucket tb) { - return waitLeaderAndIsrReady(tb).leader(); + ZooKeeperClient zkClient = getZooKeeperClient(); + return waitValue( + () -> { + Optional leaderAndIsrOpt = zkClient.getLeaderAndIsr(tb); + if (!leaderAndIsrOpt.isPresent()) { + return Optional.empty(); + } + + int leader = leaderAndIsrOpt.get().leader(); + return getReplica(tb, leader, true).map(ignored -> leader); + }, + Duration.ofMinutes(1), + "Fail to wait leader replica ready for " + tb); } public int waitAndGetStandby(TableBucket tb) {