From a0bc3d6f86b27df978857959b8c4a09d39d36ff9 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Mon, 7 Oct 2019 15:48:24 -0700 Subject: [PATCH 1/5] Switch ManagedLedger to use MetadataStore interface --- managed-ledger/pom.xml | 6 + .../mledger/ManagedLedgerException.java | 16 +- .../bookkeeper/mledger/ManagedLedgerInfo.java | 4 +- .../mledger/impl/ManagedCursorImpl.java | 2 +- .../impl/ManagedLedgerFactoryImpl.java | 6 +- .../mledger/impl/ManagedLedgerImpl.java | 2 +- .../impl/ManagedLedgerOfflineBacklog.java | 7 +- .../bookkeeper/mledger/impl/MetaStore.java | 16 +- .../mledger/impl/MetaStoreImpl.java | 286 ++++++++++++ .../mledger/impl/MetaStoreImplZookeeper.java | 422 ------------------ .../impl/ReadOnlyManagedLedgerImpl.java | 2 +- .../mledger/impl/ManagedCursorTest.java | 8 +- .../mledger/impl/ManagedLedgerErrorsTest.java | 2 +- .../mledger/impl/ManagedLedgerTest.java | 5 +- ...keeperTest.java => MetaStoreImplTest.java} | 75 +--- 15 files changed, 342 insertions(+), 517 deletions(-) create mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/MetaStoreImpl.java delete mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/MetaStoreImplZookeeper.java rename managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/{MetaStoreImplZookeeperTest.java => MetaStoreImplTest.java} (70%) diff --git a/managed-ledger/pom.xml b/managed-ledger/pom.xml index 06d057263317b..b883c8661d9c4 100644 --- a/managed-ledger/pom.xml +++ b/managed-ledger/pom.xml @@ -60,6 +60,12 @@ pulsar-common ${project.version} + + + org.apache.pulsar + pulsar-metadata + ${project.version} + com.google.guava diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java index 698acfd33e74d..6e6e6a2b848c7 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java @@ -39,8 +39,12 @@ public static ManagedLedgerException getManagedLedgerException(Throwable e) { } public static class MetaStoreException extends ManagedLedgerException { - public MetaStoreException(Exception e) { - super(e); + public MetaStoreException(Throwable t) { + super(t); + } + + public MetaStoreException(String msg) { + super(msg); } } @@ -48,12 +52,20 @@ public static class BadVersionException extends MetaStoreException { public BadVersionException(Exception e) { super(e); } + + public BadVersionException(String msg) { + super(msg); + } } public static class MetadataNotFoundException extends MetaStoreException { public MetadataNotFoundException(Exception e) { super(e); } + + public MetadataNotFoundException(String msg) { + super(msg); + } } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerInfo.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerInfo.java index ff3f6b505a57a..c03da89921e7c 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerInfo.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerInfo.java @@ -24,7 +24,7 @@ @SuppressWarnings("checkstyle:javadoctype") public class ManagedLedgerInfo { /** Z-Node version. */ - public int version; + public long version; public String creationDate; public String modificationDate; @@ -42,7 +42,7 @@ public static class LedgerInfo { public static class CursorInfo { /** Z-Node version. */ - public int version; + public long version; public String creationDate; public String modificationDate; diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index d76553877fbef..681d35534186e 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -83,7 +83,6 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.PositionBound; import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback; -import org.apache.bookkeeper.mledger.impl.MetaStore.Stat; import org.apache.bookkeeper.mledger.proto.MLDataFormats; import org.apache.bookkeeper.mledger.proto.MLDataFormats.LongProperty; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedCursorInfo; @@ -93,6 +92,7 @@ import org.apache.pulsar.common.util.collections.ConcurrentOpenLongPairRangeSet; import org.apache.pulsar.common.util.collections.LongPairRangeSet; import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPairConsumer; +import org.apache.pulsar.metadata.api.Stat; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index f8abcebafda4f..85164642f1b1e 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -68,7 +68,6 @@ import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.ManagedLedgerInitializeLedgerCallback; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.State; import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback; -import org.apache.bookkeeper.mledger.impl.MetaStore.Stat; import org.apache.bookkeeper.mledger.proto.MLDataFormats; import org.apache.bookkeeper.mledger.proto.MLDataFormats.LongProperty; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedCursorInfo; @@ -76,7 +75,8 @@ import org.apache.bookkeeper.mledger.util.Futures; import org.apache.bookkeeper.zookeeper.ZooKeeperClient; import org.apache.pulsar.common.util.DateFormatter; -import org.apache.zookeeper.KeeperException; +import org.apache.pulsar.metadata.api.Stat; +import org.apache.pulsar.metadata.impl.zookeeper.ZKMetadataStore; import org.apache.zookeeper.ZooKeeper; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -160,7 +160,7 @@ private ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPoli this.bookkeeperFactory = bookKeeperGroupFactory; this.isBookkeeperManaged = isBookkeeperManaged; this.zookeeper = isBookkeeperManaged ? zooKeeper : null; - this.store = new MetaStoreImplZookeeper(zooKeeper, orderedExecutor); + this.store = new MetaStoreImpl(new ZKMetadataStore(zooKeeper), orderedExecutor); this.config = config; this.mbean = new ManagedLedgerFactoryMBeanImpl(this); this.entryCacheManager = new EntryCacheManager(this); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java index 3a803864f5ecd..3b0495e6deb05 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java @@ -106,7 +106,6 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.VoidCallback; import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback; -import org.apache.bookkeeper.mledger.impl.MetaStore.Stat; import org.apache.bookkeeper.mledger.offload.OffloadUtils; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo.LedgerInfo; @@ -117,6 +116,7 @@ import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.InitialPosition; import org.apache.pulsar.common.util.collections.ConcurrentLongHashMap; +import org.apache.pulsar.metadata.api.Stat; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerOfflineBacklog.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerOfflineBacklog.java index 382f457688b4a..d8b49ad959b62 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerOfflineBacklog.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerOfflineBacklog.java @@ -41,6 +41,7 @@ import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.PersistentOfflineTopicStats; import org.apache.pulsar.common.util.collections.ConcurrentOpenHashMap; +import org.apache.pulsar.metadata.api.Stat; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -146,7 +147,7 @@ private void readLedgerMeta(final ManagedLedgerFactoryImpl factory, final TopicN store.getManagedLedgerInfo(managedLedgerName, false /* createIfMissing */, new MetaStore.MetaStoreCallback() { @Override - public void operationComplete(MLDataFormats.ManagedLedgerInfo mlInfo, MetaStore.Stat version) { + public void operationComplete(MLDataFormats.ManagedLedgerInfo mlInfo, Stat stat) { for (MLDataFormats.ManagedLedgerInfo.LedgerInfo ls : mlInfo.getLedgerInfoList()) { ledgers.put(ls.getLedgerId(), ls); } @@ -230,7 +231,7 @@ private void calculateCursorBacklogs(final ManagedLedgerFactoryImpl factory, fin store.getCursors(managedLedgerName, new MetaStore.MetaStoreCallback>() { @Override - public void operationComplete(List cursors, MetaStore.Stat v) { + public void operationComplete(List cursors, Stat v) { // Load existing cursors if (log.isDebugEnabled()) { log.debug("[{}] Found {} cursors", managedLedgerName, cursors.size()); @@ -336,7 +337,7 @@ public void readComplete(int rc, LedgerHandle lh, Enumeration seq, new MetaStore.MetaStoreCallback() { @Override public void operationComplete(MLDataFormats.ManagedCursorInfo info, - MetaStore.Stat version) { + Stat stat) { long cursorLedgerId = info.getCursorsLedgerId(); if (log.isDebugEnabled()) { log.debug("[{}] Cursor {} meta-data read ledger id {}", managedLedgerName, diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/MetaStore.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/MetaStore.java index 72ef566a9cc7d..8b99203c91a57 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/MetaStore.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/MetaStore.java @@ -22,19 +22,13 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.MetaStoreException; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedCursorInfo; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo; +import org.apache.pulsar.metadata.api.Stat; /** * Interface that describes the operations that the ManagedLedger need to do on the metadata store. */ public interface MetaStore { - @SuppressWarnings("checkstyle:javadoctype") - interface Stat { - int getVersion(); - long getCreationTimestamp(); - long getModificationTimestamp(); - } - @SuppressWarnings("checkstyle:javadoctype") interface UpdateLedgersIdsCallback { void updateLedgersIdsComplete(MetaStoreException status, Stat stat); @@ -64,12 +58,12 @@ interface MetaStoreCallback { * the name of the ManagedLedger * @param mlInfo * managed ledger info - * @param version + * @param stat * version object associated with current state * @param callback * callback object */ - void asyncUpdateLedgerIds(String ledgerName, ManagedLedgerInfo mlInfo, Stat version, + void asyncUpdateLedgerIds(String ledgerName, ManagedLedgerInfo mlInfo, Stat stat, MetaStoreCallback callback); /** @@ -98,12 +92,12 @@ void asyncUpdateLedgerIds(String ledgerName, ManagedLedgerInfo mlInfo, Stat vers * the name of the ManagedLedger * @param cursorName * @param info - * @param version + * @param stat * @param callback * the callback * @throws MetaStoreException */ - void asyncUpdateCursorInfo(String ledgerName, String cursorName, ManagedCursorInfo info, Stat version, + void asyncUpdateCursorInfo(String ledgerName, String cursorName, ManagedCursorInfo info, Stat stat, MetaStoreCallback callback); /** diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/MetaStoreImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/MetaStoreImpl.java new file mode 100644 index 0000000000000..8f229fdb478e9 --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/MetaStoreImpl.java @@ -0,0 +1,286 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.bookkeeper.mledger.impl; + +import com.google.protobuf.InvalidProtocolBufferException; + +import java.util.ArrayList; +import java.util.List; +import java.util.Optional; +import java.util.concurrent.CompletionException; + +import lombok.extern.slf4j.Slf4j; + +import org.apache.bookkeeper.common.util.OrderedExecutor; +import org.apache.bookkeeper.mledger.ManagedLedgerException; +import org.apache.bookkeeper.mledger.ManagedLedgerException.MetaStoreException; +import org.apache.bookkeeper.mledger.ManagedLedgerException.MetadataNotFoundException; +import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedCursorInfo; +import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo; +import org.apache.bookkeeper.util.SafeRunnable; +import org.apache.pulsar.metadata.api.MetadataStore; +import org.apache.pulsar.metadata.api.MetadataStoreException; +import org.apache.pulsar.metadata.api.Stat; + +@Slf4j +public class MetaStoreImpl implements MetaStore { + + private static final String BASE_NODE = "/managed-ledgers"; + private static final String PREFIX = BASE_NODE + "/"; + + private final MetadataStore store; + private final OrderedExecutor executor; + + public MetaStoreImpl(MetadataStore store, OrderedExecutor executor) { + this.store = store; + this.executor = executor; + } + + @Override + public void getManagedLedgerInfo(String ledgerName, boolean createIfMissing, + MetaStoreCallback callback) { + // Try to get the content or create an empty node + String path = PREFIX + ledgerName; + store.get(path) + .thenAcceptAsync(optResult -> { + if (optResult.isPresent()) { + ManagedLedgerInfo info; + try { + info = ManagedLedgerInfo.parseFrom(optResult.get().getValue()); + info = updateMLInfoTimestamp(info); + callback.operationComplete(info, optResult.get().getStat()); + } catch (InvalidProtocolBufferException e) { + callback.operationFailed(getException(e)); + } + } else { + // Z-node doesn't exist + if (createIfMissing) { + log.info("Creating '{}'", path); + + store.put(path, new byte[0], Optional.of(-1L)) + .thenAccept(stat -> { + ManagedLedgerInfo info = ManagedLedgerInfo.getDefaultInstance(); + callback.operationComplete(info, stat); + }).exceptionally(ex -> { + callback.operationFailed(getException(ex)); + return null; + }); + } else { + // Tried to open a managed ledger but it doesn't exist and we shouldn't creating it at this + // point + callback.operationFailed(new MetadataNotFoundException("Managed ledger not found")); + } + } + }, executor.chooseThread(ledgerName)) + .exceptionally(ex -> { + executor.executeOrdered(ledgerName, SafeRunnable.safeRun(() -> { + callback.operationFailed(getException(ex)); + })); + return null; + }); + } + + @Override + public void asyncUpdateLedgerIds(String ledgerName, ManagedLedgerInfo mlInfo, Stat stat, + MetaStoreCallback callback) { + if (log.isDebugEnabled()) { + log.debug("[{}] Updating metadata version={} with content={}", ledgerName, stat, mlInfo); + } + + byte[] serializedMlInfo = mlInfo.toByteArray(); // Binary format + String path = PREFIX + ledgerName; + store.put(path, serializedMlInfo, Optional.of(stat.getVersion())) + .thenAcceptAsync(newVersion -> { + callback.operationComplete(null, newVersion); + }, executor.chooseThread(ledgerName)) + .exceptionally(ex -> { + executor.executeOrdered(ledgerName, SafeRunnable.safeRun(() -> { + callback.operationFailed(getException(ex)); + })); + return null; + }); + } + + @Override + public void getCursors(String ledgerName, MetaStoreCallback> callback) { + if (log.isDebugEnabled()) { + log.debug("[{}] Get cursors list", ledgerName); + } + + String path = PREFIX + ledgerName; + store.getChildren(path) + .thenAcceptAsync(cursors -> { + callback.operationComplete(cursors, null); + }, executor.chooseThread(ledgerName)) + .exceptionally(ex -> { + executor.executeOrdered(ledgerName, SafeRunnable.safeRun(() -> { + callback.operationFailed(getException(ex)); + })); + return null; + }); + } + + @Override + public void asyncGetCursorInfo(String ledgerName, String cursorName, + MetaStoreCallback callback) { + String path = PREFIX + ledgerName + "/" + cursorName; + if (log.isDebugEnabled()) { + log.debug("Reading from {}", path); + } + + store.get(path) + .thenAcceptAsync(optRes -> { + if (optRes.isPresent()) { + try { + ManagedCursorInfo info = ManagedCursorInfo.parseFrom(optRes.get().getValue()); + callback.operationComplete(info, optRes.get().getStat()); + } catch (InvalidProtocolBufferException e) { + callback.operationFailed(getException(e)); + } + } else { + callback.operationFailed(new MetadataNotFoundException("Cursor metadata not found")); + } + }, executor.chooseThread(ledgerName)) + .exceptionally(ex -> { + executor.executeOrdered(ledgerName, SafeRunnable.safeRun(() -> { + callback.operationFailed(getException(ex)); + })); + return null; + }); + } + + @Override + public void asyncUpdateCursorInfo(String ledgerName, String cursorName, ManagedCursorInfo info, Stat stat, + MetaStoreCallback callback) { + log.info("[{}] [{}] Updating cursor info ledgerId={} mark-delete={}:{}", ledgerName, cursorName, + info.getCursorsLedgerId(), info.getMarkDeleteLedgerId(), info.getMarkDeleteEntryId()); + + String path = PREFIX + ledgerName + "/" + cursorName; + byte[] content = info.toByteArray(); // Binary format + + long expectedVersion; + + if (stat != null) { + expectedVersion = stat.getVersion(); + if (log.isDebugEnabled()) { + log.debug("[{}] Creating consumer {} on meta-data store with {}", ledgerName, cursorName, info); + } + } else { + expectedVersion = -1; + if (log.isDebugEnabled()) { + log.debug("[{}] Updating consumer {} on meta-data store with {}", ledgerName, cursorName, info); + } + } + + store.put(path, content, Optional.of(expectedVersion)) + .thenAcceptAsync(optStat -> { + callback.operationComplete(null, optStat); + }, executor.chooseThread(ledgerName)) + .exceptionally(ex -> { + executor.executeOrdered(ledgerName, SafeRunnable.safeRun(() -> { + callback.operationFailed(getException(ex)); + })); + return null; + }); + } + + @Override + public void asyncRemoveCursor(String ledgerName, String cursorName, MetaStoreCallback callback) { + String path = PREFIX + ledgerName + "/" + cursorName; + log.info("[{}] Remove consumer={}", ledgerName, cursorName); + + store.delete(path, Optional.empty()) + .thenAcceptAsync(v -> { + if (log.isDebugEnabled()) { + log.debug("[{}] [{}] cursor delete done", ledgerName, cursorName); + } + callback.operationComplete(null, null); + }, executor.chooseThread(ledgerName)) + .exceptionally(ex -> { + executor.executeOrdered(ledgerName, SafeRunnable.safeRun(() -> { + callback.operationFailed(getException(ex)); + })); + return null; + }); + } + + @Override + public void removeManagedLedger(String ledgerName, MetaStoreCallback callback) { + log.info("[{}] Remove ManagedLedger", ledgerName); + + String path = PREFIX + ledgerName; + store.delete(path, Optional.empty()) + .thenAcceptAsync(v -> { + if (log.isDebugEnabled()) { + log.debug("[{}] managed ledger delete done", ledgerName); + } + callback.operationComplete(null, null); + }, executor.chooseThread(ledgerName)) + .exceptionally(ex -> { + executor.executeOrdered(ledgerName, SafeRunnable.safeRun(() -> { + callback.operationFailed(getException(ex)); + })); + return null; + }); + } + + @Override + public Iterable getManagedLedgers() throws MetaStoreException { + try { + return store.getChildren(BASE_NODE).join(); + } catch (CompletionException e) { + throw getException(e); + } + } + + // + // update timestamp if missing or 0 + // 3 cases - timestamp does not exist for ledgers serialized before + // - timestamp is 0 for a ledger in recovery + // - ledger has timestamp which is the normal case now + + private static ManagedLedgerInfo updateMLInfoTimestamp(ManagedLedgerInfo info) { + List infoList = new ArrayList<>(info.getLedgerInfoCount()); + long currentTime = System.currentTimeMillis(); + + for (ManagedLedgerInfo.LedgerInfo ledgerInfo : info.getLedgerInfoList()) { + if (!ledgerInfo.hasTimestamp() || ledgerInfo.getTimestamp() == 0) { + ManagedLedgerInfo.LedgerInfo.Builder singleInfoBuilder = ledgerInfo.toBuilder(); + singleInfoBuilder.setTimestamp(currentTime); + infoList.add(singleInfoBuilder.build()); + } else { + infoList.add(ledgerInfo); + } + } + ManagedLedgerInfo.Builder mlInfo = ManagedLedgerInfo.newBuilder(); + mlInfo.addAllLedgerInfo(infoList); + if (info.hasTerminatedPosition()) { + mlInfo.setTerminatedPosition(info.getTerminatedPosition()); + } + return mlInfo.build(); + } + + private static MetaStoreException getException(Throwable t) { + if (t.getCause() instanceof MetadataStoreException.BadVersionException) { + return new ManagedLedgerException.BadVersionException(t.getMessage()); + } else { + return new MetaStoreException(t); + } + } +} diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/MetaStoreImplZookeeper.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/MetaStoreImplZookeeper.java deleted file mode 100644 index 2e69614f65000..0000000000000 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/MetaStoreImplZookeeper.java +++ /dev/null @@ -1,422 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - */ -package org.apache.bookkeeper.mledger.impl; - -import static org.apache.bookkeeper.mledger.util.SafeRun.safeRun; - -import com.google.common.base.Charsets; -import com.google.protobuf.InvalidProtocolBufferException; -import com.google.protobuf.TextFormat; -import com.google.protobuf.TextFormat.ParseException; - -import java.io.File; -import java.nio.charset.Charset; -import java.util.ArrayList; -import java.util.List; -import java.util.concurrent.ForkJoinPool; -import java.util.function.Consumer; - -import org.apache.bookkeeper.common.util.OrderedExecutor; -import org.apache.bookkeeper.mledger.ManagedLedgerException; -import org.apache.bookkeeper.mledger.ManagedLedgerException.BadVersionException; -import org.apache.bookkeeper.mledger.ManagedLedgerException.MetaStoreException; -import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedCursorInfo; -import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo; -import org.apache.zookeeper.AsyncCallback.StringCallback; -import org.apache.zookeeper.CreateMode; -import org.apache.zookeeper.KeeperException; -import org.apache.zookeeper.KeeperException.Code; -import org.apache.zookeeper.ZooDefs; -import org.apache.zookeeper.ZooKeeper; -import org.apache.zookeeper.data.ACL; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -@SuppressWarnings("checkstyle:javadoctype") -public class MetaStoreImplZookeeper implements MetaStore { - - private static final Charset Encoding = Charsets.UTF_8; - private static final List Acl = ZooDefs.Ids.OPEN_ACL_UNSAFE; - - private static final String prefixName = "/managed-ledgers"; - private static final String prefix = prefixName + "/"; - - private final ZooKeeper zk; - private final OrderedExecutor executor; - - private static class ZKStat implements Stat { - private final int version; - private final long creationTimestamp; - private final long modificationTimestamp; - - ZKStat(org.apache.zookeeper.data.Stat stat) { - this.version = stat.getVersion(); - this.creationTimestamp = stat.getCtime(); - this.modificationTimestamp = stat.getMtime(); - } - - ZKStat() { - this.version = 0; - this.creationTimestamp = System.currentTimeMillis(); - this.modificationTimestamp = System.currentTimeMillis(); - } - - @Override - public int getVersion() { - return version; - } - - @Override - public long getCreationTimestamp() { - return creationTimestamp; - } - - @Override - public long getModificationTimestamp() { - return modificationTimestamp; - } - } - - public MetaStoreImplZookeeper(ZooKeeper zk, OrderedExecutor executor) - throws Exception { - this.zk = zk; - this.executor = executor; - } - - // - // update timestamp if missing or 0 - // 3 cases - timestamp does not exist for ledgers serialized before - // - timestamp is 0 for a ledger in recovery - // - ledger has timestamp which is the normal case now - - private ManagedLedgerInfo updateMLInfoTimestamp(ManagedLedgerInfo info) { - List infoList = new ArrayList<>(info.getLedgerInfoCount()); - long currentTime = System.currentTimeMillis(); - - for (ManagedLedgerInfo.LedgerInfo ledgerInfo : info.getLedgerInfoList()) { - if (!ledgerInfo.hasTimestamp() || ledgerInfo.getTimestamp() == 0) { - ManagedLedgerInfo.LedgerInfo.Builder singleInfoBuilder = ledgerInfo.toBuilder(); - singleInfoBuilder.setTimestamp(currentTime); - infoList.add(singleInfoBuilder.build()); - } else { - infoList.add(ledgerInfo); - } - } - ManagedLedgerInfo.Builder mlInfo = ManagedLedgerInfo.newBuilder(); - mlInfo.addAllLedgerInfo(infoList); - if (info.hasTerminatedPosition()) { - mlInfo.setTerminatedPosition(info.getTerminatedPosition()); - } - return mlInfo.build(); - } - - @Override - public void getManagedLedgerInfo(final String ledgerName, boolean createIfMissing, - final MetaStoreCallback callback) { - // Try to get the content or create an empty node - zk.getData(prefix + ledgerName, false, - (rc, path, ctx, readData, stat) -> executor.executeOrdered(ledgerName, safeRun(() -> { - if (rc == Code.OK.intValue()) { - try { - ManagedLedgerInfo info = parseManagedLedgerInfo(readData); - info = updateMLInfoTimestamp(info); - callback.operationComplete(info, new ZKStat(stat)); - } catch (ParseException | InvalidProtocolBufferException e) { - callback.operationFailed(new MetaStoreException(e)); - } - } else if (rc == Code.NONODE.intValue()) { - // Z-node doesn't exist - if (createIfMissing) { - log.info("Creating '{}{}'", prefix, ledgerName); - - StringCallback createcb = (rc1, path1, ctx1, name) -> { - if (rc1 == Code.OK.intValue()) { - ManagedLedgerInfo info = ManagedLedgerInfo.getDefaultInstance(); - callback.operationComplete(info, new ZKStat()); - } else { - callback.operationFailed( - new MetaStoreException(KeeperException.create(Code.get(rc1)))); - } - }; - - asyncCreateFullPathOptimistic(prefixName, ledgerName, new byte[0], Acl, - CreateMode.PERSISTENT, createcb); - } else { - // Tried to open a managed ledger but it doesn't exist and we shouldn't creating it at this - // point - - callback.operationFailed(new ManagedLedgerException.MetadataNotFoundException( - KeeperException.create(Code.get(rc)))); - } - } else { - // Other ZK error - callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc)))); - } - })), null); - } - - @Override - public void asyncUpdateLedgerIds(String ledgerName, ManagedLedgerInfo mlInfo, Stat stat, - final MetaStoreCallback callback) { - - ZKStat zkStat = (ZKStat) stat; - if (log.isDebugEnabled()) { - log.debug("[{}] Updating metadata version={} with content={}", ledgerName, zkStat.version, mlInfo); - } - - byte[] serializedMlInfo = mlInfo.toByteArray(); // Binary format - - zk.setData(prefix + ledgerName, serializedMlInfo, zkStat.getVersion(), - (rc, path, zkCtx, stat1) -> executor.executeOrdered(ledgerName, safeRun(() -> { - if (log.isDebugEnabled()) { - log.debug("[{}] UpdateLedgersIdsCallback.processResult rc={} newVersion={}", ledgerName, - Code.get(rc), stat != null ? stat.getVersion() : "null"); - } - MetaStoreException status = null; - if (rc == Code.BADVERSION.intValue()) { - // Content has been modified on ZK since our last read - status = new BadVersionException(KeeperException.create(Code.get(rc))); - callback.operationFailed(status); - } else if (rc != Code.OK.intValue()) { - status = new MetaStoreException(KeeperException.create(Code.get(rc))); - callback.operationFailed(status); - } else { - callback.operationComplete(null, new ZKStat(stat1)); - } - })), null); - } - - @Override - public void getCursors(final String ledgerName, final MetaStoreCallback> callback) { - if (log.isDebugEnabled()) { - log.debug("[{}] Get cursors list", ledgerName); - } - zk.getChildren(prefix + ledgerName, false, - (rc, path, ctx, children, stat) -> executor.executeOrdered(ledgerName, safeRun(() -> { - if (log.isDebugEnabled()) { - log.debug("[{}] getConsumers complete rc={} children={}", ledgerName, Code.get(rc), children); - } - if (rc != Code.OK.intValue()) { - callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc)))); - return; - } - - if (log.isDebugEnabled()) { - log.debug("[{}] Get childrend completed version={}", ledgerName, stat.getVersion()); - } - callback.operationComplete(children, new ZKStat(stat)); - })), null); - } - - @Override - public void asyncGetCursorInfo(String ledgerName, String consumerName, - final MetaStoreCallback callback) { - String path = prefix + ledgerName + "/" + consumerName; - if (log.isDebugEnabled()) { - log.debug("Reading from {}", path); - } - - zk.getData(path, false, (rc, path1, ctx, data, stat) -> executor.executeOrdered(ledgerName, safeRun(() -> { - if (rc != Code.OK.intValue()) { - callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc)))); - } else { - try { - ManagedCursorInfo info = parseManagedCursorInfo(data); - callback.operationComplete(info, new ZKStat(stat)); - } catch (ParseException | InvalidProtocolBufferException e) { - callback.operationFailed(new MetaStoreException(e)); - } - } - })), null); - - if (log.isDebugEnabled()) { - log.debug("Reading from {} ok", path); - } - } - - @Override - public void asyncUpdateCursorInfo(final String ledgerName, final String cursorName, final ManagedCursorInfo info, - Stat stat, final MetaStoreCallback callback) { - log.info("[{}] [{}] Updating cursor info ledgerId={} mark-delete={}:{}", ledgerName, cursorName, - info.getCursorsLedgerId(), info.getMarkDeleteLedgerId(), info.getMarkDeleteEntryId()); - - String path = prefix + ledgerName + "/" + cursorName; - byte[] content = info.toByteArray(); // Binary format - - if (stat == null) { - if (log.isDebugEnabled()) { - log.debug("[{}] Creating consumer {} on meta-data store with {}", ledgerName, cursorName, info); - } - zk.create(path, content, Acl, CreateMode.PERSISTENT, - (rc, path1, ctx, name) -> executor.executeOrdered(ledgerName, safeRun(() -> { - if (rc != Code.OK.intValue()) { - log.warn("[{}] Error creating cosumer {} node on meta-data store with {}: ", ledgerName, - cursorName, info, Code.get(rc)); - callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc)))); - } else { - if (log.isDebugEnabled()) { - log.debug("[{}] Created consumer {} on meta-data store with {}", ledgerName, cursorName, - info); - } - callback.operationComplete(null, new ZKStat()); - } - })), null); - } else { - ZKStat zkStat = (ZKStat) stat; - if (log.isDebugEnabled()) { - log.debug("[{}] Updating consumer {} on meta-data store with {}", ledgerName, cursorName, info); - } - zk.setData(path, content, zkStat.getVersion(), - (rc, path1, ctx, stat1) -> executor.executeOrdered(ledgerName, safeRun(() -> { - if (rc == Code.BADVERSION.intValue()) { - callback.operationFailed(new BadVersionException(KeeperException.create(Code.get(rc)))); - } else if (rc != Code.OK.intValue()) { - callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc)))); - } else { - callback.operationComplete(null, new ZKStat(stat1)); - } - })), null); - } - } - - @Override - public void asyncRemoveCursor(final String ledgerName, final String consumerName, - final MetaStoreCallback callback) { - log.info("[{}] Remove consumer={}", ledgerName, consumerName); - zk.delete(prefix + ledgerName + "/" + consumerName, -1, - (rc, path, ctx) -> executor.executeOrdered(ledgerName, safeRun(() -> { - if (log.isDebugEnabled()) { - log.debug("[{}] [{}] zk delete done. rc={}", ledgerName, consumerName, Code.get(rc)); - } - if (rc == Code.OK.intValue()) { - callback.operationComplete(null, null); - } else { - callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc)))); - } - })), null); - } - - @Override - public void removeManagedLedger(String ledgerName, MetaStoreCallback callback) { - log.info("[{}] Remove ManagedLedger", ledgerName); - zk.delete(prefix + ledgerName, -1, (rc, path, ctx) -> executor.executeOrdered(ledgerName, safeRun(() -> { - if (log.isDebugEnabled()) { - log.debug("[{}] zk delete done. rc={}", ledgerName, Code.get(rc)); - } - if (rc == Code.OK.intValue()) { - callback.operationComplete(null, null); - } else { - callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc)))); - } - })), null); - } - - @Override - public Iterable getManagedLedgers() throws MetaStoreException { - try { - return zk.getChildren(prefixName, false); - } catch (Exception e) { - throw new MetaStoreException(e); - } - } - - private ManagedLedgerInfo parseManagedLedgerInfo(byte[] data) - throws ParseException, InvalidProtocolBufferException { - // First try binary format, then fallback to text - try { - return ManagedLedgerInfo.parseFrom(data); - } catch (InvalidProtocolBufferException e) { - // Fallback to parsing protobuf text format - ManagedLedgerInfo.Builder builder = ManagedLedgerInfo.newBuilder(); - TextFormat.merge(new String(data, Encoding), builder); - return builder.build(); - } - } - - private ManagedCursorInfo parseManagedCursorInfo(byte[] data) - throws ParseException, InvalidProtocolBufferException { - // First try binary format, then fallback to text - try { - return ManagedCursorInfo.parseFrom(data); - } catch (InvalidProtocolBufferException e) { - // Fallback to parsing protobuf text format - ManagedCursorInfo.Builder builder = ManagedCursorInfo.newBuilder(); - TextFormat.merge(new String(data, Encoding), builder); - return builder.build(); - } - - } - - void asyncCreateFullPathOptimistic( - final String basePath, final String nodePath, final byte[] data, - final List acl, final CreateMode createMode, final StringCallback callback) { - String fullPath = basePath + "/" + nodePath; - - zk.create(fullPath, data, acl, createMode, - (rc, path, ignoreCtx1, name) -> { - Runnable retry = () -> { - asyncCreateFullPathOptimistic(basePath, nodePath, data, - acl, createMode, callback); - }; - - Consumer complete = (finalrc) -> { - executor.executeOrdered(nodePath, safeRun(() -> { - callback.processResult(finalrc, path, null, name); - })); - }; - - if (rc != Code.NONODE.intValue()) { - complete.accept(rc); - return; - } - - // Since I got a nonode, it means that my parents don't exist - // create mode is persistent since ephemeral nodes can't be - // parents - String nodeParent = new File(nodePath).getParent(); - if (nodeParent == null) { - zk.exists(basePath, false, - (existsRc, existsPath, ignoreCtx2, stat) -> { - if (existsRc == Code.OK.intValue()) { - if (stat != null) { - retry.run(); - } else { - complete.accept(Code.NONODE.intValue()); - } - } else { - complete.accept(existsRc); - } - }, null); - } else { - nodeParent = nodeParent.replace("\\", "/"); - asyncCreateFullPathOptimistic( - basePath, nodeParent, new byte[0], acl, CreateMode.PERSISTENT, - (parentRc, parentPath, ignoreCtx3, parentName) -> { - if (parentRc == Code.OK.intValue() || parentRc == Code.NODEEXISTS.intValue()) { - retry.run(); - } else { - complete.accept(parentRc); - } - }); - } - }, null); - } - - private static final Logger log = LoggerFactory.getLogger(MetaStoreImplZookeeper.class); -} diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ReadOnlyManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ReadOnlyManagedLedgerImpl.java index 9721b15bd4280..621131752720a 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ReadOnlyManagedLedgerImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ReadOnlyManagedLedgerImpl.java @@ -34,9 +34,9 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.MetadataNotFoundException; import org.apache.bookkeeper.mledger.ReadOnlyCursor; import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback; -import org.apache.bookkeeper.mledger.impl.MetaStore.Stat; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo.LedgerInfo; +import org.apache.pulsar.metadata.api.Stat; @Slf4j public class ReadOnlyManagedLedgerImpl extends ManagedLedgerImpl { diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index b0db9264d8f91..c3a1acd42b326 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -77,10 +77,10 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.VoidCallback; import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback; -import org.apache.bookkeeper.mledger.impl.MetaStore.Stat; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedCursorInfo; import org.apache.bookkeeper.mledger.proto.MLDataFormats.PositionInfo; import org.apache.bookkeeper.test.MockedBookKeeperTestCase; +import org.apache.pulsar.metadata.api.Stat; import org.apache.zookeeper.KeeperException.Code; import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; @@ -92,7 +92,7 @@ public class ManagedCursorTest extends MockedBookKeeperTestCase { private static final Charset Encoding = Charsets.UTF_8; - + @DataProvider(name = "useOpenRangeSet") public static Object[][] useOpenRangeSet() { return new Object[][] { { Boolean.TRUE }, { Boolean.FALSE } }; @@ -2873,7 +2873,7 @@ public void testRecoverCursorAheadOfLastPosition() throws Exception { final long markDeleteLedgerId = 2L; final long markDeleteEntryId = -1L; - MetaStoreImplZookeeper mockMetaStore = mock(MetaStoreImplZookeeper.class); + MetaStore mockMetaStore = mock(MetaStore.class); doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) { ManagedCursorInfo info = ManagedCursorInfo.newBuilder().setCursorsLedgerId(cursorsLedgerId) @@ -2967,6 +2967,6 @@ public void deleteMessagesCheckhMarkDelete() throws Exception { assertEquals(c1.getMarkDeletedPosition(), positions[markDelete]); assertEquals(c1.getReadPosition(), positions[markDelete + 1]); } - + private static final Logger log = LoggerFactory.getLogger(ManagedCursorTest.class); } diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerErrorsTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerErrorsTest.java index 3d47a6f5697d7..416a58cc256ca 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerErrorsTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerErrorsTest.java @@ -343,7 +343,7 @@ public void recoverAfterZnodeVersionError() throws Exception { ledger.addEntry("entry".getBytes()); fail("should fail"); } catch (ManagedLedgerFencedException e) { - assertEquals(e.getCause().getCause().getClass(), org.apache.zookeeper.KeeperException.BadVersionException.class); + assertEquals(e.getCause().getClass(), ManagedLedgerException.BadVersionException.class); // ok } diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java index d94099bc68069..0fa8b34040678 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java @@ -95,7 +95,6 @@ import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback; -import org.apache.bookkeeper.mledger.impl.MetaStore.Stat; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo.LedgerInfo; import org.apache.bookkeeper.test.MockedBookKeeperTestCase; @@ -106,6 +105,8 @@ import org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata; import org.apache.pulsar.common.protocol.ByteBufPair; import org.apache.pulsar.common.util.protobuf.ByteBufCodedOutputStream; +import org.apache.pulsar.metadata.api.Stat; +import org.apache.pulsar.metadata.impl.zookeeper.ZKMetadataStore; import org.apache.zookeeper.CreateMode; import org.apache.zookeeper.KeeperException.Code; import org.apache.zookeeper.ZooDefs; @@ -1811,7 +1812,7 @@ public void testBackwardCompatiblityForMeta() throws Exception { ml.addEntry("msg2".getBytes()); ml.close(); - MetaStore store = new MetaStoreImplZookeeper(zkc, executor); + MetaStore store = new MetaStoreImpl(new ZKMetadataStore(zkc), executor); CountDownLatch l1 = new CountDownLatch(1); // obtain the ledger info diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/MetaStoreImplZookeeperTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/MetaStoreImplTest.java similarity index 70% rename from managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/MetaStoreImplZookeeperTest.java rename to managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/MetaStoreImplTest.java index 021e4c153dbc9..2337b4dd28b8b 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/MetaStoreImplZookeeperTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/MetaStoreImplTest.java @@ -18,32 +18,30 @@ */ package org.apache.bookkeeper.mledger.impl; -import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.fail; -import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; -import java.util.concurrent.ExecutionException; import java.util.concurrent.atomic.AtomicReference; + import org.apache.bookkeeper.mledger.ManagedLedgerException.MetaStoreException; import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback; -import org.apache.bookkeeper.mledger.impl.MetaStore.Stat; import org.apache.bookkeeper.mledger.proto.MLDataFormats; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedCursorInfo; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo; import org.apache.bookkeeper.test.MockedBookKeeperTestCase; +import org.apache.pulsar.metadata.api.Stat; +import org.apache.pulsar.metadata.impl.zookeeper.ZKMetadataStore; import org.apache.zookeeper.CreateMode; -import org.apache.zookeeper.KeeperException; import org.apache.zookeeper.KeeperException.Code; import org.apache.zookeeper.ZooDefs; import org.testng.annotations.Test; -public class MetaStoreImplZookeeperTest extends MockedBookKeeperTestCase { +public class MetaStoreImplTest extends MockedBookKeeperTestCase { @Test void getMLList() throws Exception { - MetaStore store = new MetaStoreImplZookeeper(zkc, executor); + MetaStore store = new MetaStoreImpl(new ZKMetadataStore(zkc), executor); zkc.failNow(Code.CONNECTIONLOSS); @@ -57,7 +55,7 @@ void getMLList() throws Exception { @Test void deleteNonExistingML() throws Exception { - MetaStore store = new MetaStoreImplZookeeper(zkc, executor); + MetaStore store = new MetaStoreImpl(new ZKMetadataStore(zkc), executor); AtomicReference exception = new AtomicReference<>(); CountDownLatch counter = new CountDownLatch(1); @@ -82,7 +80,7 @@ public void operationFailed(MetaStoreException e) { @Test(timeOut = 20000) void readMalformedML() throws Exception { - MetaStore store = new MetaStoreImplZookeeper(zkc, executor); + MetaStore store = new MetaStoreImpl(new ZKMetadataStore(zkc), executor); zkc.create("/managed-ledgers/my_test", "non-valid".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); @@ -105,7 +103,7 @@ public void operationComplete(ManagedLedgerInfo result, Stat version) { @Test(timeOut = 20000) void readMalformedCursorNode() throws Exception { - MetaStore store = new MetaStoreImplZookeeper(zkc, executor); + MetaStore store = new MetaStoreImpl(new ZKMetadataStore(zkc), executor); zkc.create("/managed-ledgers/my_test", "".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); zkc.create("/managed-ledgers/my_test/c1", "non-valid".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, @@ -129,7 +127,7 @@ public void operationComplete(ManagedCursorInfo result, Stat version) { @Test(timeOut = 20000) void failInCreatingMLnode() throws Exception { - MetaStore store = new MetaStoreImplZookeeper(zkc, executor); + MetaStore store = new MetaStoreImpl(new ZKMetadataStore(zkc), executor); final CountDownLatch latch = new CountDownLatch(1); @@ -151,7 +149,7 @@ public void operationComplete(ManagedLedgerInfo result, Stat version) { @Test(timeOut = 20000) void updatingCursorNode() throws Exception { - final MetaStore store = new MetaStoreImplZookeeper(zkc, executor); + MetaStore store = new MetaStoreImpl(new ZKMetadataStore(zkc), executor); zkc.create("/managed-ledgers/my_test", "".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); @@ -187,7 +185,7 @@ public void operationComplete(Void result, Stat version) { @Test(timeOut = 20000) void updatingMLNode() throws Exception { - final MetaStore store = new MetaStoreImplZookeeper(zkc, executor); + MetaStore store = new MetaStoreImpl(new ZKMetadataStore(zkc), executor); zkc.create("/managed-ledgers/my_test", "".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); @@ -218,55 +216,4 @@ public void operationComplete(Void result, Stat version) { latch.await(); } - - @Test(timeOut = 20000) - public void createOptimisticBaseNotExist() throws Exception { - CompletableFuture promise = new CompletableFuture<>(); - - MetaStoreImplZookeeper store = new MetaStoreImplZookeeper(zkc, executor); - store.asyncCreateFullPathOptimistic( - "/foo", "bar/zar/gar", new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT, - (rc, path, ctx, name) -> { - if (rc != KeeperException.Code.OK.intValue()) { - promise.completeExceptionally(KeeperException.create(rc)); - } else { - promise.complete(null); - } - }); - try { - promise.get(); - fail("should have failed"); - } catch (ExecutionException ee) { - assertEquals(ee.getCause().getClass(), KeeperException.NoNodeException.class); - } - } - - @Test(timeOut = 20000) - public void createOptimisticBaseExists() throws Exception { - MetaStoreImplZookeeper store = new MetaStoreImplZookeeper(zkc, executor); - zkc.create("/foo", new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); - CompletableFuture promise = new CompletableFuture<>(); - store.asyncCreateFullPathOptimistic( - "/foo", "bar/zar/gar", new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT, - (rc, path, ctx, name) -> { - if (rc != KeeperException.Code.OK.intValue()) { - promise.completeExceptionally(KeeperException.create(rc)); - } else { - promise.complete(null); - } - }); - promise.get(); - - CompletableFuture promise2 = new CompletableFuture<>(); - store.asyncCreateFullPathOptimistic( - "/foo", "blah", new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT, - (rc, path, ctx, name) -> { - if (rc != KeeperException.Code.OK.intValue()) { - promise2.completeExceptionally(KeeperException.create(rc)); - } else { - promise2.complete(null); - } - }); - promise2.get(); - } } From 4c7ab2a521e40fe2e9a5f62615714b572a719ee6 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Sat, 4 Jan 2020 18:25:15 -0800 Subject: [PATCH 2/5] Use different thread for ZK callbacks --- .../impl/zookeeper/ZKMetadataStore.java | 187 ++++++++++-------- 1 file changed, 102 insertions(+), 85 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/zookeeper/ZKMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/zookeeper/ZKMetadataStore.java index 58b1c107dce13..348b6f5b71e56 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/zookeeper/ZKMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/zookeeper/ZKMetadataStore.java @@ -18,13 +18,13 @@ */ package org.apache.pulsar.metadata.impl.zookeeper; -import com.google.common.annotations.VisibleForTesting; - import java.io.IOException; import java.util.Collections; import java.util.List; import java.util.Optional; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import org.apache.bookkeeper.util.ZkUtils; import org.apache.bookkeeper.zookeeper.BoundExponentialBackoffRetryPolicy; @@ -43,26 +43,32 @@ import org.apache.zookeeper.ZooDefs; import org.apache.zookeeper.ZooKeeper; +import com.google.common.annotations.VisibleForTesting; + +import io.netty.util.concurrent.DefaultThreadFactory; + public class ZKMetadataStore implements MetadataStore { private final ZooKeeper zkc; + private final ExecutorService executor; public ZKMetadataStore(String metadataURL, MetadataStoreConfig metadataStoreConfig) throws IOException { try { - zkc = ZooKeeperClient.newBuilder() - .connectString(metadataURL) + zkc = ZooKeeperClient.newBuilder().connectString(metadataURL) .connectRetryPolicy(new BoundExponentialBackoffRetryPolicy(100, 60_000, Integer.MAX_VALUE)) .allowReadOnlyMode(metadataStoreConfig.isAllowReadOnlyOperations()) - .sessionTimeoutMs(metadataStoreConfig.getSessionTimeoutMillis()) - .build(); + .sessionTimeoutMs(metadataStoreConfig.getSessionTimeoutMillis()).build(); } catch (KeeperException | InterruptedException e) { throw new IOException(e); } + + this.executor = Executors.newSingleThreadExecutor(new DefaultThreadFactory("zk-metadata-store-callback")); } @VisibleForTesting public ZKMetadataStore(ZooKeeper zkc) { this.zkc = zkc; + this.executor = Executors.newSingleThreadExecutor(new DefaultThreadFactory("zk-metadata-store-callback")); } @Override @@ -71,14 +77,16 @@ public CompletableFuture> get(String path) { try { zkc.getData(path, null, (rc, path1, ctx, data, stat) -> { - Code code = Code.get(rc); - if (code == Code.OK) { - future.complete(Optional.of(new GetResult(data, getStat(stat)))); - } else if (code == Code.NONODE) { - future.complete(Optional.empty()); - } else { - future.completeExceptionally(getException(code, path)); - } + executor.execute(() -> { + Code code = Code.get(rc); + if (code == Code.OK) { + future.complete(Optional.of(new GetResult(data, getStat(stat)))); + } else if (code == Code.NONODE) { + future.complete(Optional.empty()); + } else { + future.completeExceptionally(getException(code, path)); + } + }); }, null); } catch (Throwable t) { future.completeExceptionally(new MetadataStoreException(t)); @@ -93,36 +101,37 @@ public CompletableFuture> getChildren(String path) { try { zkc.getChildren(path, null, (rc, path1, ctx, children) -> { - Code code = Code.get(rc); - if (code == Code.OK) { - Collections.sort(children); - future.complete(children); - } else if (code == Code.NONODE) { - // The node we want may not exist yet, so put a watcher on its existence - // before throwing up the exception. Its possible that the node could have - // been created after the call to getChildren, but before the call to exists(). - // If this is the case, exists will return true, and we just call getChildren again. - exists(path).thenAccept(exists -> { - if (exists) { - getChildren(path) - .thenAccept(c -> future.complete(c)) - .exceptionally(ex -> { - future.completeExceptionally(ex); - return null; - }); - } else { - // Z-node does not exist - future.complete(Collections.emptyList()); - } - }).exceptionally(ex -> { - future.completeExceptionally(ex); - return null; - }); + executor.execute(() -> { + Code code = Code.get(rc); + if (code == Code.OK) { + Collections.sort(children); + future.complete(children); + } else if (code == Code.NONODE) { + // The node we want may not exist yet, so put a watcher on its existence + // before throwing up the exception. Its possible that the node could have + // been created after the call to getChildren, but before the call to exists(). + // If this is the case, exists will return true, and we just call getChildren + // again. + exists(path).thenAccept(exists -> { + if (exists) { + getChildren(path).thenAccept(c -> future.complete(c)).exceptionally(ex -> { + future.completeExceptionally(ex); + return null; + }); + } else { + // Z-node does not exist + future.complete(Collections.emptyList()); + } + }).exceptionally(ex -> { + future.completeExceptionally(ex); + return null; + }); - future.complete(Collections.emptyList()); - } else { - future.completeExceptionally(getException(code, path)); - } + future.complete(Collections.emptyList()); + } else { + future.completeExceptionally(getException(code, path)); + } + }); }, null); } catch (Throwable t) { future.completeExceptionally(new MetadataStoreException(t)); @@ -137,14 +146,16 @@ public CompletableFuture exists(String path) { try { zkc.exists(path, null, (StatCallback) (rc, path1, ctx, stat) -> { - Code code = Code.get(rc); - if (code == Code.OK) { - future.complete(true); - } else if (code == Code.NONODE) { - future.complete(false); - } else { - future.completeExceptionally(getException(code, path)); - } + executor.execute(() -> { + Code code = Code.get(rc); + if (code == Code.OK) { + future.complete(true); + } else if (code == Code.NONODE) { + future.complete(false); + } else { + future.completeExceptionally(getException(code, path)); + } + }); }, future); } catch (Throwable t) { future.completeExceptionally(new MetadataStoreException(t)); @@ -163,39 +174,42 @@ public CompletableFuture put(String path, byte[] value, Optional opt try { if (hasVersion && expectedVersion == -1) { ZkUtils.asyncCreateFullPathOptimistic(zkc, path, value, ZooDefs.Ids.OPEN_ACL_UNSAFE, - CreateMode.PERSISTENT, - (rc, path1, ctx, name) -> { - Code code = Code.get(rc); - if (code == Code.OK) { - future.complete(new Stat(0, 0, 0)); - } else if (code == Code.NODEEXISTS) { - // We're emulating a request to create node, so the version is invalid - future.completeExceptionally(getException(Code.BADVERSION, path)); - } else { - future.completeExceptionally(getException(code, path)); - } + CreateMode.PERSISTENT, (rc, path1, ctx, name) -> { + executor.execute(() -> { + Code code = Code.get(rc); + if (code == Code.OK) { + future.complete(new Stat(0, 0, 0)); + } else if (code == Code.NODEEXISTS) { + // We're emulating a request to create node, so the version is invalid + future.completeExceptionally(getException(Code.BADVERSION, path)); + } else { + future.completeExceptionally(getException(code, path)); + } + }); }, null); } else { zkc.setData(path, value, expectedVersion, (rc, path1, ctx, stat) -> { - Code code = Code.get(rc); - if (code == Code.OK) { - future.complete(getStat(stat)); - } else if (code == Code.NONODE) { - if (hasVersion) { - // We're emulating here a request to update or create the znode, depending on the version - future.completeExceptionally(getException(Code.BADVERSION, path)); + executor.execute(() -> { + Code code = Code.get(rc); + if (code == Code.OK) { + future.complete(getStat(stat)); + } else if (code == Code.NONODE) { + if (hasVersion) { + // We're emulating here a request to update or create the znode, depending on + // the version + future.completeExceptionally(getException(Code.BADVERSION, path)); + } else { + // The z-node does not exist, let's create it first + put(path, value, Optional.of(-1L)).thenAccept(s -> future.complete(s)) + .exceptionally(ex -> { + future.completeExceptionally(ex.getCause()); + return null; + }); + } } else { - // The z-node does not exist, let's create it first - put(path, value, Optional.of(-1L)) - .thenAccept(s -> future.complete(s)) - .exceptionally(ex -> { - future.completeExceptionally(ex.getCause()); - return null; - }); + future.completeExceptionally(getException(code, path)); } - } else { - future.completeExceptionally(getException(code, path)); - } + }); }, null); } } catch (Throwable t) { @@ -213,12 +227,14 @@ public CompletableFuture delete(String path, Optional optExpectedVer try { zkc.delete(path, expectedVersion, (rc, path1, ctx) -> { - Code code = Code.get(rc); - if (code == Code.OK) { - future.complete(null); - } else { - future.completeExceptionally(getException(code, path)); - } + executor.execute(() -> { + Code code = Code.get(rc); + if (code == Code.OK) { + future.complete(null); + } else { + future.completeExceptionally(getException(code, path)); + } + }); }, null); } catch (Throwable t) { future.completeExceptionally(new MetadataStoreException(t)); @@ -230,6 +246,7 @@ public CompletableFuture delete(String path, Optional optExpectedVer @Override public void close() throws Exception { zkc.close(); + executor.shutdownNow(); } private static Stat getStat(org.apache.zookeeper.data.Stat zkStat) { From 913cd991e5b408b6c7bcbeead0c9d5c08b4f1d22 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Sat, 4 Jan 2020 18:38:50 -0800 Subject: [PATCH 3/5] Properly close metadata store in ml factory --- .../mledger/impl/ManagedLedgerFactoryImpl.java | 10 +++++++++- .../metadata/impl/zookeeper/ZKMetadataStore.java | 7 ++++++- 2 files changed, 15 insertions(+), 2 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index 85164642f1b1e..b88bc7f3444f1 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -75,6 +75,7 @@ import org.apache.bookkeeper.mledger.util.Futures; import org.apache.bookkeeper.zookeeper.ZooKeeperClient; import org.apache.pulsar.common.util.DateFormatter; +import org.apache.pulsar.metadata.api.MetadataStore; import org.apache.pulsar.metadata.api.Stat; import org.apache.pulsar.metadata.impl.zookeeper.ZKMetadataStore; import org.apache.zookeeper.ZooKeeper; @@ -101,6 +102,7 @@ public class ManagedLedgerFactoryImpl implements ManagedLedgerFactory { private final ScheduledFuture statsTask; private final long cacheEvictionTimeThresholdNanos; + private final MetadataStore metadataStore; private static final int StatsPeriodSeconds = 60; @@ -160,7 +162,8 @@ private ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPoli this.bookkeeperFactory = bookKeeperGroupFactory; this.isBookkeeperManaged = isBookkeeperManaged; this.zookeeper = isBookkeeperManaged ? zooKeeper : null; - this.store = new MetaStoreImpl(new ZKMetadataStore(zooKeeper), orderedExecutor); + this.metadataStore = new ZKMetadataStore(zooKeeper); + this.store = new MetaStoreImpl(metadataStore, orderedExecutor); this.config = config; this.mbean = new ManagedLedgerFactoryMBeanImpl(this); this.entryCacheManager = new EntryCacheManager(this); @@ -464,6 +467,11 @@ public void closeFailed(ManagedLedgerException exception, Object ctx) { cacheEvictionExecutor.shutdownNow(); entryCacheManager.clear(); + try { + metadataStore.close(); + } catch (Exception e) { + throw new ManagedLedgerException(e); + } } @Override diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/zookeeper/ZKMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/zookeeper/ZKMetadataStore.java index 348b6f5b71e56..77b3ba2494643 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/zookeeper/ZKMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/zookeeper/ZKMetadataStore.java @@ -49,11 +49,13 @@ public class ZKMetadataStore implements MetadataStore { + private final boolean isZkManaged; private final ZooKeeper zkc; private final ExecutorService executor; public ZKMetadataStore(String metadataURL, MetadataStoreConfig metadataStoreConfig) throws IOException { try { + isZkManaged = true; zkc = ZooKeeperClient.newBuilder().connectString(metadataURL) .connectRetryPolicy(new BoundExponentialBackoffRetryPolicy(100, 60_000, Integer.MAX_VALUE)) .allowReadOnlyMode(metadataStoreConfig.isAllowReadOnlyOperations()) @@ -67,6 +69,7 @@ public ZKMetadataStore(String metadataURL, MetadataStoreConfig metadataStoreConf @VisibleForTesting public ZKMetadataStore(ZooKeeper zkc) { + this.isZkManaged = false; this.zkc = zkc; this.executor = Executors.newSingleThreadExecutor(new DefaultThreadFactory("zk-metadata-store-callback")); } @@ -245,7 +248,9 @@ public CompletableFuture delete(String path, Optional optExpectedVer @Override public void close() throws Exception { - zkc.close(); + if (isZkManaged) { + zkc.close(); + } executor.shutdownNow(); } From 93716cd3c0b5565a07a875a0568e49e4fb7bc88b Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Mon, 10 Feb 2020 15:05:25 -0800 Subject: [PATCH 4/5] Fixed licenses --- pulsar-sql/presto-distribution/LICENSE | 2 ++ pulsar-sql/presto-pulsar/pom.xml | 1 + 2 files changed, 3 insertions(+) diff --git a/pulsar-sql/presto-distribution/LICENSE b/pulsar-sql/presto-distribution/LICENSE index 1e8971a167865..596cca3a29a27 100644 --- a/pulsar-sql/presto-distribution/LICENSE +++ b/pulsar-sql/presto-distribution/LICENSE @@ -356,6 +356,8 @@ The Apache Software License, Version 2.0 * Avro - avro-1.9.1.jar - avro-protobuf-1.9.1.jar + * Caffeine + - caffeine-2.6.2.jar * Javax - javax.inject-1.jar - javax.inject-1.jar diff --git a/pulsar-sql/presto-pulsar/pom.xml b/pulsar-sql/presto-pulsar/pom.xml index 87dfa2838534c..4d184d58c37bd 100644 --- a/pulsar-sql/presto-pulsar/pom.xml +++ b/pulsar-sql/presto-pulsar/pom.xml @@ -124,6 +124,7 @@ org.apache.pulsar:pulsar-client-original org.apache.pulsar:pulsar-client-admin-original org.apache.pulsar:managed-ledger + org.apache.pulsar:pulsar-metadata org.glassfish.jersey*:* javax.ws.rs:* From ee58ff6a7e6954940b833621c592e366dba62b6d Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Wed, 12 Feb 2020 10:46:16 -0800 Subject: [PATCH 5/5] Fixed problem in MockZookeeper --- .../src/test/java/org/apache/zookeeper/MockZooKeeper.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/managed-ledger/src/test/java/org/apache/zookeeper/MockZooKeeper.java b/managed-ledger/src/test/java/org/apache/zookeeper/MockZooKeeper.java index f4160edf3737b..11c84f85a93f3 100644 --- a/managed-ledger/src/test/java/org/apache/zookeeper/MockZooKeeper.java +++ b/managed-ledger/src/test/java/org/apache/zookeeper/MockZooKeeper.java @@ -354,7 +354,7 @@ public void getChildren(final String path, final Watcher watcher, final Children } String child = item.substring(path.length() + 1); - if (!child.contains("/")) { + if (item.charAt(path.length()) == '/' && !child.contains("/")) { children.add(child); } }