diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java index 247b610aee0..02e47aec526 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java @@ -32,6 +32,9 @@ import java.util.TreeSet; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Semaphore; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; + import org.apache.bookkeeper.client.BKException; import org.apache.bookkeeper.conf.ServerConfiguration; import org.apache.bookkeeper.meta.LedgerManager; @@ -139,8 +142,9 @@ public void gc(GarbageCleaner garbageCleaner) { } // Iterate over all the ledger on the metadata store - long zkOpTimeout = this.conf.getZkTimeout() * 2; - LedgerRangeIterator ledgerRangeIterator = ledgerManager.getLedgerRanges(zkOpTimeout); + long zkOpTimeoutMs = this.conf.getZkTimeout() * 2; + LedgerRangeIterator ledgerRangeIterator = ledgerManager + .getLedgerRanges(zkOpTimeoutMs); Set ledgersInMetadata = null; long start; long end = -1; @@ -167,13 +171,20 @@ public void gc(GarbageCleaner garbageCleaner) { if (verifyMetadataOnGc) { int rc = BKException.Code.OK; try { - result(ledgerManager.readLedgerMetadata(bkLid)); - } catch (BKException e) { - rc = e.getCode(); + result(ledgerManager.readLedgerMetadata(bkLid), zkOpTimeoutMs, TimeUnit.MILLISECONDS); + } catch (BKException | TimeoutException e) { + if (e instanceof BKException) { + rc = ((BKException) e).getCode(); + } else { + LOG.warn("Time-out while fetching metadata for Ledger {} : {}.", bkLid, + e.getMessage()); + + continue; + } } if (rc != BKException.Code.NoSuchLedgerExistsOnMetadataServerException) { LOG.warn("Ledger {} Missing in metadata list, but ledgerManager returned rc: {}.", - bkLid, rc); + bkLid, rc); continue; } } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/CleanupLedgerManager.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/CleanupLedgerManager.java index 3c7ddc51d30..a3ce432961f 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/CleanupLedgerManager.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/CleanupLedgerManager.java @@ -210,13 +210,13 @@ public void processResult(int rc, String path, Object ctx) { } @Override - public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutSec) { + public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutMs) { closeLock.readLock().lock(); try { if (closed) { return new ClosedLedgerRangeIterator(); } - return underlying.getLedgerRanges(zkOpTimeoutSec); + return underlying.getLedgerRanges(zkOpTimeoutMs); } finally { closeLock.readLock().unlock(); } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/FlatLedgerManager.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/FlatLedgerManager.java index 8c7f05f9f7f..7c89326ffe7 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/FlatLedgerManager.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/FlatLedgerManager.java @@ -89,7 +89,7 @@ public void asyncProcessLedgers(final Processor processor, } @Override - public LedgerRangeIterator getLedgerRanges(long zkOpTimeOutSec) { + public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutMs) { return new LedgerRangeIterator() { // single iterator, can visit only one time boolean nextCalled = false; @@ -103,7 +103,7 @@ private synchronized void preload() throws IOException { try { zkActiveLedgers = ledgerListToSet( - ZkUtils.getChildrenInSingleNode(zk, ledgerRootPath, zkOpTimeOutSec), + ZkUtils.getChildrenInSingleNode(zk, ledgerRootPath, zkOpTimeoutMs), ledgerRootPath); nextRange = new LedgerRange(zkActiveLedgers); } catch (KeeperException.NoNodeException e) { diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/HierarchicalLedgerManager.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/HierarchicalLedgerManager.java index 079fcfacd25..9cd6aed92e3 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/HierarchicalLedgerManager.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/HierarchicalLedgerManager.java @@ -87,9 +87,9 @@ protected long getLedgerId(String ledgerPath) throws IOException { } @Override - public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutSec) { - LedgerRangeIterator legacyLedgerRangeIterator = legacyLM.getLedgerRanges(zkOpTimeoutSec); - LedgerRangeIterator longLedgerRangeIterator = longLM.getLedgerRanges(zkOpTimeoutSec); + public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutMs) { + LedgerRangeIterator legacyLedgerRangeIterator = legacyLM.getLedgerRanges(zkOpTimeoutMs); + LedgerRangeIterator longLedgerRangeIterator = longLM.getLedgerRanges(zkOpTimeoutMs); return new HierarchicalLedgerRangeIterator(legacyLedgerRangeIterator, longLedgerRangeIterator); } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/LedgerManager.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/LedgerManager.java index 15da14b083b..56c447d538c 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/LedgerManager.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/LedgerManager.java @@ -149,12 +149,12 @@ void asyncProcessLedgers(Processor processor, AsyncCallback.VoidCallback f /** * Loop to scan a range of metadata from metadata storage. * - * @param zkOpTimeOutSec + * @param zkOpTimeOutMs * Iterator considers timeout while fetching ledger-range from * zk. * @return will return a iterator of the Ranges */ - LedgerRangeIterator getLedgerRanges(long zkOpTimeOutSec); + LedgerRangeIterator getLedgerRanges(long zkOpTimeOutMs); /** * Used to represent the Ledgers range returned from the diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/LegacyHierarchicalLedgerManager.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/LegacyHierarchicalLedgerManager.java index 1a63407e7ad..9828719ba5a 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/LegacyHierarchicalLedgerManager.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/LegacyHierarchicalLedgerManager.java @@ -153,8 +153,8 @@ protected String getLedgerParentNodeRegex() { } @Override - public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutSec) { - return new LegacyHierarchicalLedgerRangeIterator(zkOpTimeoutSec); + public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutMs) { + return new LegacyHierarchicalLedgerRangeIterator(zkOpTimeoutMs); } /** @@ -166,10 +166,10 @@ private class LegacyHierarchicalLedgerRangeIterator implements LedgerRangeIterat private String curL1Nodes = ""; private boolean iteratorDone = false; private LedgerRange nextRange = null; - private final long zkOpTimeoutSec; + private final long zkOpTimeoutMs; - public LegacyHierarchicalLedgerRangeIterator(long zkOpTimeoutSec) { - this.zkOpTimeoutSec = zkOpTimeoutSec; + public LegacyHierarchicalLedgerRangeIterator(long zkOpTimeoutMs) { + this.zkOpTimeoutMs = zkOpTimeoutMs; } /** @@ -266,7 +266,7 @@ LedgerRange getLedgerRangeByLevel(final String level1, final String level2) String nodePath = nodeBuilder.toString(); List ledgerNodes = null; try { - ledgerNodes = ZkUtils.getChildrenInSingleNode(zk, nodePath, zkOpTimeoutSec); + ledgerNodes = ZkUtils.getChildrenInSingleNode(zk, nodePath, zkOpTimeoutMs); } catch (KeeperException.NoNodeException e) { /* If the node doesn't exist, we must have raced with a recursive node removal, just * return an empty list. */ diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/LongHierarchicalLedgerManager.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/LongHierarchicalLedgerManager.java index 95f8e48bde8..7ec2f5af1b8 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/LongHierarchicalLedgerManager.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/LongHierarchicalLedgerManager.java @@ -139,8 +139,8 @@ public void process(String lNode, VoidCallback cb) { } @Override - public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutSec) { - return new LongHierarchicalLedgerRangeIterator(zkOpTimeoutSec); + public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutMs) { + return new LongHierarchicalLedgerRangeIterator(zkOpTimeoutMs); } @@ -149,7 +149,7 @@ public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutSec) { */ private class LongHierarchicalLedgerRangeIterator implements LedgerRangeIterator { LedgerRangeIterator rootIterator; - final long zkOpTimeoutSec; + final long zkOpTimeoutMs; /** * Returns all children with path as a parent. If path is non-existent, @@ -163,7 +163,7 @@ private class LongHierarchicalLedgerRangeIterator implements LedgerRangeIterator */ List getChildrenAt(String path) throws IOException { try { - List children = ZkUtils.getChildrenInSingleNode(zk, path, zkOpTimeoutSec); + List children = ZkUtils.getChildrenInSingleNode(zk, path, zkOpTimeoutMs); Collections.sort(children); return children; } catch (KeeperException.NoNodeException e) { @@ -285,8 +285,8 @@ public LedgerRange next() throws IOException { } } - private LongHierarchicalLedgerRangeIterator(long zkOpTimeoutSec) { - this.zkOpTimeoutSec = zkOpTimeoutSec; + private LongHierarchicalLedgerRangeIterator(long zkOpTimeoutMs) { + this.zkOpTimeoutMs = zkOpTimeoutMs; } private void bootstrap() throws IOException { diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/MSLedgerManagerFactory.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/MSLedgerManagerFactory.java index 201a0a3edb9..9b21c876e34 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/MSLedgerManagerFactory.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/meta/MSLedgerManagerFactory.java @@ -642,7 +642,7 @@ public LedgerRange next() throws IOException { } @Override - public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutSec) { + public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutMs) { return new MSLedgerRangeIterator(); } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/util/ZkUtils.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/util/ZkUtils.java index 23d0b421c3e..cc0612fc77c 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/util/ZkUtils.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/util/ZkUtils.java @@ -25,7 +25,6 @@ import java.io.IOException; import java.util.List; import java.util.concurrent.CountDownLatch; -import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import org.apache.bookkeeper.conf.AbstractConfiguration; @@ -222,7 +221,7 @@ private static class GetChildrenCtx { * @throws InterruptedException * @throws IOException */ - public static List getChildrenInSingleNode(final ZooKeeper zk, final String node, long timeOutSec) + public static List getChildrenInSingleNode(final ZooKeeper zk, final String node, long zkOpTimeoutMs) throws InterruptedException, IOException, KeeperException.NoNodeException { final GetChildrenCtx ctx = new GetChildrenCtx(); getChildrenInSingleNode(zk, node, new GenericCallback>() { @@ -240,13 +239,20 @@ public void operationComplete(int rc, List ledgers) { }); synchronized (ctx) { + long startTime = System.currentTimeMillis(); while (!ctx.done) { try { - ctx.wait(timeOutSec > 0 ? TimeUnit.SECONDS.toMillis(timeOutSec) : 0); + ctx.wait(zkOpTimeoutMs > 0 ? zkOpTimeoutMs : 0); } catch (InterruptedException e) { ctx.rc = Code.OPERATIONTIMEOUT.intValue(); ctx.done = true; } + // timeout the process if get-children response not received + // zkOpTimeoutMs. + if (zkOpTimeoutMs > 0 && (System.currentTimeMillis() - startTime) >= zkOpTimeoutMs) { + ctx.rc = Code.OPERATIONTIMEOUT.intValue(); + ctx.done = true; + } } } if (Code.NONODE.intValue() == ctx.rc) { diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java index 0d6e5d64775..adef116475c 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java @@ -970,7 +970,7 @@ void unsupported() { } @Override - public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutSec) { + public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutMs) { final AtomicBoolean hasnext = new AtomicBoolean(true); return new LedgerManager.LedgerRangeIterator() { @Override diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/ParallelLedgerRecoveryTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/ParallelLedgerRecoveryTest.java index e6c3c222928..177cbe38df3 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/ParallelLedgerRecoveryTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/ParallelLedgerRecoveryTest.java @@ -112,8 +112,8 @@ public CompletableFuture> readLedgerMetadata(long ledg } @Override - public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutSec) { - return lm.getLedgerRanges(zkOpTimeoutSec); + public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutMs) { + return lm.getLedgerRanges(zkOpTimeoutMs); } @Override diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/meta/MockLedgerManager.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/meta/MockLedgerManager.java index 0f818b66c9e..952837be37e 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/meta/MockLedgerManager.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/meta/MockLedgerManager.java @@ -188,7 +188,7 @@ public void asyncProcessLedgers(Processor processor, AsyncCallback.VoidCal } @Override - public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutSec) { + public LedgerRangeIterator getLedgerRanges(long zkOpTimeoutMs) { return null; } diff --git a/metadata-drivers/etcd/src/main/java/org/apache/bookkeeper/metadata/etcd/EtcdLedgerManager.java b/metadata-drivers/etcd/src/main/java/org/apache/bookkeeper/metadata/etcd/EtcdLedgerManager.java index a40a4bc383b..f91f2f7f1f4 100644 --- a/metadata-drivers/etcd/src/main/java/org/apache/bookkeeper/metadata/etcd/EtcdLedgerManager.java +++ b/metadata-drivers/etcd/src/main/java/org/apache/bookkeeper/metadata/etcd/EtcdLedgerManager.java @@ -420,7 +420,7 @@ private void processLedgers(KeyStream ks, } @Override - public LedgerRangeIterator getLedgerRanges(long opTimeOutSec) { + public LedgerRangeIterator getLedgerRanges(long opTimeOutMs) { KeyStream ks = new KeyStream<>( kvClient, ByteSequence.fromString(EtcdUtils.getLedgerKey(scope, 0L)),