diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java index d39af92008b3b..27faf85195601 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java @@ -20,11 +20,9 @@ import com.google.common.base.Predicate; import com.google.common.collect.Range; - import java.util.List; import java.util.Map; import java.util.Set; - import org.apache.bookkeeper.common.annotation.InterfaceAudience; import org.apache.bookkeeper.common.annotation.InterfaceStability; import org.apache.bookkeeper.mledger.AsyncCallbacks.ClearBacklogCallback; @@ -234,6 +232,8 @@ void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallba */ boolean hasMoreEntries(); + boolean hasMoreEntries(PositionImpl position); + /** * Return the number of messages that this cursor still has to read. * diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java index a2e0ba00b5924..65a47cd0c8696 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java @@ -20,6 +20,7 @@ import io.netty.buffer.ByteBuf; import java.util.Map; +import java.util.NavigableMap; import java.util.concurrent.CompletableFuture; import org.apache.bookkeeper.common.annotation.InterfaceAudience; import org.apache.bookkeeper.common.annotation.InterfaceStability; @@ -30,6 +31,8 @@ import org.apache.bookkeeper.mledger.AsyncCallbacks.OffloadCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.OpenCursorCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.TerminateCallback; +import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; +import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.bookkeeper.mledger.intercept.ManagedLedgerInterceptor; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo.LedgerInfo; import org.apache.pulsar.common.api.proto.CommandSubscribe.InitialPosition; @@ -59,6 +62,19 @@ @InterfaceStability.Stable public interface ManagedLedger { + /** + * Make ManagedLedger ready to work + * @param callback + * @param ctx + */ + void initialize(final ManagedLedgerInitializeLedgerCallback callback, final Object ctx); + + boolean isValidPosition(PositionImpl nextReadPosition); + + boolean hasMoreEntries(PositionImpl nextReadPosition); + + void addWaitingEntryCallBack(WaitingEntryCallBack streamingEntryReader); + /** * @return the unique name of this ManagedLedger */ @@ -100,6 +116,8 @@ public interface ManagedLedger { */ void asyncAddEntry(byte[] data, AddEntryCallback callback, Object ctx); + void asyncReadEntry(PositionImpl position, AsyncCallbacks.ReadEntryCallback callback, Object ctx); + /** * Append a new entry to the end of a managed ledger. * @@ -357,6 +375,26 @@ public interface ManagedLedger { */ long getNumberOfEntries(); + long getEntriesAddedCounter(); + + long getLastLedgerCreatedTimestamp(); + + long getLastLedgerCreationFailureTimestamp(); + + int getWaitingCursorsCount(); + + long getCurrentLedgerEntries(); + + long getCurrentLedgerSize(); + + NavigableMap getLedgersInfo(); + + CompletableFuture getLedgerMetadata(long ledgerId); + + boolean ledgerExists(long ledgerId); + + void asyncDeleteLedgerFromBookKeeper(long ledgerId); + /** * Get the total number of active entries for this managed ledger. * @@ -387,6 +425,23 @@ public interface ManagedLedger { */ long getEstimatedBacklogSize(); + /** + * Get estimated backlog size from a specific position. + * @param pos + * @return + */ + long getEstimatedBacklogSize(PositionImpl pos); + + /** + * number of entries are in add progress + */ + int getPendingAddEntriesCount(); + + /** + * Get the total size in bytes of all the entries stored in this cache. + */ + long getCacheSize(); + /** * Return the size of all ledgers offloaded to 2nd tier storage */ @@ -430,6 +485,13 @@ public interface ManagedLedger { */ ManagedLedgerMXBean getStats(); + /** + * Remove all entries already read by active cursors and + * remove entries older than the cutoff threshold + * @param maxTimestamp + */ + void doCacheEviction(long maxTimestamp); + /** * Delete the ManagedLedger. * @@ -504,6 +566,12 @@ public interface ManagedLedger { */ Position getLastConfirmedEntry(); + /** + * Get state of the managed ledger. + * @return + */ + String getState(); + /** * Signaling managed ledger that we can resume after BK write failure */ @@ -590,6 +658,27 @@ void asyncSetProperties(Map properties, final AsyncCallbacks.Upd * */ CompletableFuture asyncFindPosition(com.google.common.base.Predicate predicate); + /** + * Get the entry position at a given distance from a given position. + * + * @param startPosition + * starting position + * @param n + * number of entries to skip ahead + * @param startRange + * specifies whether or not to include the start position in calculating the distance + * @return the new position that is n entries ahead + */ + PositionImpl getPositionAfterN(final PositionImpl startPosition, long n, + ManagedLedgerImpl.PositionBound startRange); + + /** + * first position of current topic + */ + PositionImpl getFirstPosition(); + + PositionImpl getLastPosition(); + /** * Get the ManagedLedgerInterceptor for ManagedLedger. * */ @@ -600,4 +689,21 @@ void asyncSetProperties(Map properties, final AsyncCallbacks.Upd * will got null if corresponding ledger not exists. */ CompletableFuture getLedgerInfo(long ledgerId); + + /** + * get the valid position next to the given one + * @param position current postion + */ + PositionImpl getNextValidPosition(final PositionImpl position); + + interface ManagedLedgerInitializeLedgerCallback { + void initializeComplete(); + + void initializeFailed(ManagedLedgerException e); + } + + // define boundaries for position based seeks and searches + enum PositionBound { + startIncluded, startExcluded + } } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactory.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactory.java index 4e103d09a2fad..c02422870f2e9 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactory.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactory.java @@ -19,7 +19,6 @@ package org.apache.bookkeeper.mledger; import java.util.function.Supplier; - import org.apache.bookkeeper.common.annotation.InterfaceAudience; import org.apache.bookkeeper.common.annotation.InterfaceStability; import org.apache.bookkeeper.mledger.AsyncCallbacks.DeleteLedgerCallback; @@ -113,7 +112,9 @@ ReadOnlyCursor openReadOnlyCursor(String managedLedgerName, Position startPositi * @param ctx */ void asyncOpenReadOnlyCursor(String managedLedgerName, Position startPosition, ManagedLedgerConfig config, - OpenReadOnlyCursorCallback callback, Object ctx); + OpenReadOnlyCursorCallback callback, Object ctx); + + void close(ManagedLedger ledger); /** * Get the current metadata info for a managed ledger. 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 bdc0bb87e7fd3..35173620b5937 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 @@ -26,7 +26,6 @@ import static org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.createManagedLedgerException; import static org.apache.bookkeeper.mledger.util.Errors.isNoSuchLedgerExistsException; import static org.apache.bookkeeper.mledger.util.SafeRun.safeRun; - import com.google.common.annotations.VisibleForTesting; import com.google.common.base.MoreObjects; import com.google.common.base.Predicate; @@ -38,9 +37,7 @@ import com.google.common.collect.Sets; import com.google.common.util.concurrent.RateLimiter; import com.google.protobuf.InvalidProtocolBufferException; - import io.netty.util.concurrent.FastThreadLocal; - import java.time.Clock; import java.util.ArrayDeque; import java.util.ArrayList; @@ -62,7 +59,6 @@ import java.util.concurrent.atomic.AtomicReferenceFieldUpdater; import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantReadWriteLock; - import org.apache.bookkeeper.client.AsyncCallback.CloseCallback; import org.apache.bookkeeper.client.AsyncCallback.DeleteCallback; import org.apache.bookkeeper.client.AsyncCallback.OpenCallback; @@ -80,15 +76,15 @@ import org.apache.bookkeeper.mledger.AsyncCallbacks.SkipEntriesCallback; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; +import org.apache.bookkeeper.mledger.ManagedCursorMXBean; import org.apache.bookkeeper.mledger.ManagedLedger; +import org.apache.bookkeeper.mledger.ManagedLedger.PositionBound; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.ManagedLedgerException.CursorAlreadyClosedException; import org.apache.bookkeeper.mledger.ManagedLedgerException.MetaStoreException; import org.apache.bookkeeper.mledger.ManagedLedgerException.NoMoreEntriesToReadException; -import org.apache.bookkeeper.mledger.ManagedCursorMXBean; 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.proto.MLDataFormats; import org.apache.bookkeeper.mledger.proto.MLDataFormats.LongProperty; @@ -584,6 +580,26 @@ public void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadE ledger.asyncReadEntries(op); } + public void asyncReadEntries(int restNumOfEntries, long restMaxSizeBytes, ReadEntriesCallback callback, + Object ctx, PositionImpl maxPosition, List alreadyRead) { + checkArgument(restNumOfEntries >= 0); + if (isClosed()) { + callback.readEntriesFailed(new ManagedLedgerException("Cursor was already closed"), ctx); + return; + } + + int restNumOfEntriesToRead = applyMaxSizeCap(restNumOfEntries, restMaxSizeBytes); + + PENDING_READ_OPS_UPDATER.incrementAndGet(this); + OpReadEntry op = OpReadEntry + .create(this, readPosition, restNumOfEntriesToRead + alreadyRead.size(), callback, ctx, maxPosition); + if (!alreadyRead.isEmpty()) { + op.readEntriesComplete(alreadyRead, ctx); + } else { + ledger.asyncReadEntries(op); + } + } + @Override public Entry getNthEntry(int n, IndividualDeletedEntries deletedEntries) throws InterruptedException, ManagedLedgerException { @@ -1191,17 +1207,7 @@ public Set asyncReplayEntries(Set positi } // filters out messages which are already acknowledged - Set alreadyAcknowledgedPositions = Sets.newHashSet(); - lock.readLock().lock(); - try { - positions.stream() - .filter(position -> individualDeletedMessages.contains(((PositionImpl) position).getLedgerId(), - ((PositionImpl) position).getEntryId()) - || ((PositionImpl) position).compareTo(markDeletePosition) <= 0) - .forEach(alreadyAcknowledgedPositions::add); - } finally { - lock.readLock().unlock(); - } + Set alreadyAcknowledgedPositions = filterAlreadyAcknowledgedCallback(positions); final int totalValidPositions = positions.size() - alreadyAcknowledgedPositions.size(); final AtomicReference exception = new AtomicReference<>(); @@ -1256,6 +1262,21 @@ public synchronized void readEntryFailed(ManagedLedgerException mle, Object ctx) return alreadyAcknowledgedPositions; } + protected Set filterAlreadyAcknowledgedCallback(Set positions) { + Set alreadyAcknowledgedPositions = Sets.newHashSet(); + lock.readLock().lock(); + try { + positions.stream() + .filter(position -> individualDeletedMessages.contains(position.getLedgerId(), + position.getEntryId()) + || ((PositionImpl) position).compareTo(markDeletePosition) <= 0) + .forEach(alreadyAcknowledgedPositions::add); + } finally { + lock.readLock().unlock(); + } + return alreadyAcknowledgedPositions; + } + protected long getNumberOfEntries(Range range) { long allEntries = ledger.getNumberOfEntries(range); @@ -1494,7 +1515,7 @@ long getNumIndividualDeletedEntriesToSkip(long numEntries) { return tempDeletedMessages.get(); } - boolean hasMoreEntries(PositionImpl position) { + public boolean hasMoreEntries(PositionImpl position) { PositionImpl lastPositionInLedger = ledger.getLastPosition(); if (position.compareTo(lastPositionInLedger) <= 0) { return getNumberOfEntries(Range.closed(position, lastPositionInLedger)) > 0; 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 ab280464219a2..7d5591cb09459 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 @@ -20,12 +20,9 @@ import static com.google.common.base.Preconditions.checkArgument; import static org.apache.bookkeeper.mledger.ManagedLedgerException.getManagedLedgerException; - import com.google.common.base.Predicates; import com.google.common.collect.Maps; - import io.netty.util.concurrent.DefaultThreadFactory; - import java.io.IOException; import java.util.ArrayList; import java.util.List; @@ -41,7 +38,6 @@ import java.util.concurrent.TimeUnit; import java.util.function.Supplier; import java.util.stream.Collectors; - import org.apache.bookkeeper.client.BKException; import org.apache.bookkeeper.client.BookKeeper; import org.apache.bookkeeper.common.util.OrderedExecutor; @@ -54,6 +50,7 @@ import org.apache.bookkeeper.mledger.AsyncCallbacks.OpenLedgerCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.OpenReadOnlyCursorCallback; import org.apache.bookkeeper.mledger.ManagedLedger; +import org.apache.bookkeeper.mledger.ManagedLedger.ManagedLedgerInitializeLedgerCallback; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.ManagedLedgerException.MetaStoreException; @@ -67,7 +64,6 @@ import org.apache.bookkeeper.mledger.ManagedLedgerInfo.PositionInfo; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.ReadOnlyCursor; -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.proto.MLDataFormats; @@ -78,8 +74,8 @@ import org.apache.bookkeeper.stats.NullStatsLogger; import org.apache.bookkeeper.stats.StatsLogger; import org.apache.bookkeeper.zookeeper.ZooKeeperClient; -import org.apache.pulsar.common.util.DateFormatter; import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; +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.ZKMetadataStore; @@ -100,9 +96,9 @@ public class ManagedLedgerFactoryImpl implements ManagedLedgerFactory { protected final ManagedLedgerFactoryMBeanImpl mbean; - protected final ConcurrentHashMap> ledgers = new ConcurrentHashMap<>(); + protected final ConcurrentHashMap> ledgers = new ConcurrentHashMap<>(); protected final ConcurrentHashMap pendingInitializeLedgers = - new ConcurrentHashMap<>(); + new ConcurrentHashMap<>(); private final EntryCacheManager entryCacheManager; private long lastStatTimestamp = System.nanoTime(); @@ -116,10 +112,10 @@ public class ManagedLedgerFactoryImpl implements ManagedLedgerFactory { private static class PendingInitializeManagedLedger { - private final ManagedLedgerImpl ledger; + private final ManagedLedger ledger; private final long createTimeMs; - PendingInitializeManagedLedger(ManagedLedgerImpl ledger) { + PendingInitializeManagedLedger(ManagedLedger ledger) { this.ledger = ledger; this.createTimeMs = System.currentTimeMillis(); } @@ -235,7 +231,7 @@ public BookKeeper get(EnsemblePlacementPolicyConfig policy) { private synchronized void flushCursors() { ledgers.values().forEach(mlfuture -> { if (mlfuture.isDone() && !mlfuture.isCompletedExceptionally()) { - ManagedLedgerImpl ml = mlfuture.getNow(null); + ManagedLedger ml = mlfuture.getNow(null); if (ml != null) { ml.getCursors().forEach(c -> ((ManagedCursorImpl) c).flush()); } @@ -250,9 +246,9 @@ private synchronized void refreshStats() { mbean.refreshStats(period, TimeUnit.NANOSECONDS); ledgers.values().forEach(mlfuture -> { if (mlfuture.isDone() && !mlfuture.isCompletedExceptionally()) { - ManagedLedgerImpl ml = mlfuture.getNow(null); + ManagedLedger ml = mlfuture.getNow(null); if (ml != null) { - ml.mbean.refreshStats(period, TimeUnit.NANOSECONDS); + ((ManagedLedgerMBeanImpl) ml.getStats()).refreshStats(period, TimeUnit.NANOSECONDS); } } }); @@ -283,7 +279,7 @@ private synchronized void doCacheEviction() { ledgers.values().forEach(mlfuture -> { if (mlfuture.isDone() && !mlfuture.isCompletedExceptionally()) { - ManagedLedgerImpl ml = mlfuture.getNow(null); + ManagedLedger ml = mlfuture.getNow(null); if (ml != null) { ml.doCacheEviction(maxTimestamp); } @@ -296,7 +292,7 @@ private synchronized void doCacheEviction() { * * @return */ - public Map getManagedLedgers() { + public Map getManagedLedgers() { // Return a view of already created ledger by filtering futures not yet completed return Maps.filterValues(Maps.transformValues(ledgers, future -> future.getNow(null)), Predicates.notNull()); } @@ -347,11 +343,11 @@ public void asyncOpen(final String name, final ManagedLedgerConfig config, final Supplier mlOwnershipChecker, final Object ctx) { // If the ledger state is bad, remove it from the map. - CompletableFuture existingFuture = ledgers.get(name); + CompletableFuture existingFuture = ledgers.get(name); if (existingFuture != null) { if (existingFuture.isDone()) { try { - ManagedLedgerImpl l = existingFuture.get(); + ManagedLedger l = existingFuture.get(); if (l.getState().equals(State.Fenced.toString()) || l.getState().equals(State.Closed.toString())) { // Managed ledger is in unusable state. Recreate it. log.warn("[{}] Attempted to open ledger in {} state. Removing from the map to recreate it", name, @@ -380,8 +376,8 @@ public void asyncOpen(final String name, final ManagedLedgerConfig config, final // Ensure only one managed ledger is created and initialized ledgers.computeIfAbsent(name, (mlName) -> { // Create the managed ledger - CompletableFuture future = new CompletableFuture<>(); - final ManagedLedgerImpl newledger = new ManagedLedgerImpl(this, + CompletableFuture future = new CompletableFuture<>(); + final ManagedLedger newledger = new ManagedLedgerImpl(this, bookkeeperFactory.get( new EnsemblePlacementPolicyConfig(config.getBookKeeperEnsemblePlacementPolicyClassName(), config.getBookKeeperEnsemblePlacementPolicyProperties())), @@ -428,10 +424,9 @@ public void closeFailed(ManagedLedgerException exception, Object ctx) { }); } - - @Override - public ReadOnlyCursor openReadOnlyCursor(String managedLedgerName, Position startPosition, ManagedLedgerConfig config) + public ReadOnlyCursor openReadOnlyCursor(String managedLedgerName, Position startPosition, + ManagedLedgerConfig config) throws InterruptedException, ManagedLedgerException { class Result { ReadOnlyCursor c = null; @@ -487,7 +482,7 @@ public void asyncOpenReadOnlyCursor(String managedLedgerName, Position startPosi }); } - void close(ManagedLedger ledger) { + public void close(ManagedLedger ledger) { // Remove the ledger from the internal factory cache ledgers.remove(ledger.getName()); entryCacheManager.removeEntryCache(ledger.getName()); @@ -502,8 +497,8 @@ public void shutdown() throws InterruptedException, ManagedLedgerException { final CountDownLatch latch = new CountDownLatch(numLedgers); log.info("Closing {} ledgers", numLedgers); - for (CompletableFuture ledgerFuture : ledgers.values()) { - ManagedLedgerImpl ledger = ledgerFuture.getNow(null); + for (CompletableFuture ledgerFuture : ledgers.values()) { + ManagedLedger ledger = ledgerFuture.getNow(null); if (ledger == null) { latch.countDown(); continue; @@ -732,7 +727,7 @@ public void deleteLedgerFailed(ManagedLedgerException exception, Object ctx) { @Override public void asyncDelete(String name, DeleteLedgerCallback callback, Object ctx) { - CompletableFuture future = ledgers.get(name); + CompletableFuture future = ledgers.get(name); if (future == null) { // Managed ledger does not exist and we're not currently trying to open it deleteManagedLedger(name, callback, ctx); 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 0e5ba685e9741..78ce6f18ee71f 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 @@ -23,8 +23,6 @@ import static java.lang.Math.min; import static org.apache.bookkeeper.mledger.util.Errors.isNoSuchLedgerExistsException; import static org.apache.bookkeeper.mledger.util.SafeRun.safeRun; - -import com.fasterxml.jackson.core.JsonProcessingException; import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.BoundType; import com.google.common.collect.ImmutableMap; @@ -77,7 +75,6 @@ import org.apache.bookkeeper.client.LedgerHandle; import org.apache.bookkeeper.client.api.ReadHandle; import org.apache.bookkeeper.common.util.Backoff; -import org.apache.bookkeeper.common.util.JsonUtil; import org.apache.bookkeeper.common.util.OrderedExecutor; import org.apache.bookkeeper.common.util.OrderedScheduler; import org.apache.bookkeeper.common.util.Retries; @@ -232,11 +229,6 @@ enum State { // After handling the BK write failure, managed ledger will get signalled to create a new ledger } - // define boundaries for position based seeks and searches - public enum PositionBound { - startIncluded, startExcluded - } - private static final AtomicReferenceFieldUpdater STATE_UPDATER = AtomicReferenceFieldUpdater .newUpdater(ManagedLedgerImpl.class, State.class, "state"); protected volatile State state = null; @@ -310,7 +302,7 @@ public ManagedLedgerImpl(ManagedLedgerFactoryImpl factory, BookKeeper bookKeeper } } - synchronized void initialize(final ManagedLedgerInitializeLedgerCallback callback, final Object ctx) { + public synchronized void initialize(final ManagedLedgerInitializeLedgerCallback callback, final Object ctx) { log.info("Opening managed ledger {}", name); // Fetch the list of existing ledgers in the managed ledger @@ -2053,7 +2045,7 @@ void discardEntriesFromCache(ManagedCursorImpl cursor, PositionImpl newPosition) } } - void doCacheEviction(long maxTimestamp) { + public void doCacheEviction(long maxTimestamp) { // Always remove all entries already read by active cursors PositionImpl slowestReaderPos = getEarlierReadPositionForActiveCursors(); if (slowestReaderPos != null) { @@ -2529,7 +2521,7 @@ public void deleteCursorFailed(ManagedLedgerException exception, Object ctx) { } } - private void asyncDeleteLedgerFromBookKeeper(long ledgerId) { + public void asyncDeleteLedgerFromBookKeeper(long ledgerId) { asyncDeleteLedger(ledgerId, DEFAULT_LEDGER_DELETE_RETRIES); } @@ -2954,15 +2946,16 @@ private void cleanupOffloaded(long ledgerId, UUID uuid, String offloadDriverName * identify offloader */ Map offloadDriverMetadata, String cleanupReason) { + offloadDriverMetadata.put("ManagedLedgerName", name); Retries.run(Backoff.exponentialJittered(TimeUnit.SECONDS.toMillis(1), TimeUnit.SECONDS.toHours(1)).limit(10), Retries.NonFatalPredicate, () -> config.getLedgerOffloader().deleteOffloaded(ledgerId, uuid, offloadDriverMetadata), scheduledExecutor, name).whenComplete((ignored, exception) -> { - if (exception != null) { - log.warn("Error cleaning up offload for {}, (cleanup reason: {})", ledgerId, cleanupReason, - exception); - } - }); + if (exception != null) { + log.warn("Error cleaning up offload for {}, (cleanup reason: {})", ledgerId, cleanupReason, + exception); + } + }); } /** @@ -3194,7 +3187,7 @@ public PositionImpl getFirstPosition() { return new PositionImpl(ledgerId, -1); } - PositionImpl getLastPosition() { + public PositionImpl getLastPosition() { return lastConfirmedEntry; } @@ -3381,12 +3374,6 @@ public void setConfig(ManagedLedgerConfig config) { this.cursors.forEach(c -> c.setThrottleMarkDelete(config.getThrottleMarkDelete())); } - interface ManagedLedgerInitializeLedgerCallback { - void initializeComplete(); - - void initializeFailed(ManagedLedgerException e); - } - // Expose internal values for debugging purposes public long getEntriesAddedCounter() { return ENTRIES_ADDED_COUNTER_UPDATER.get(this); 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 index caa21a149c764..96a081b252377 100644 --- 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 @@ -19,14 +19,11 @@ 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; @@ -177,9 +174,20 @@ public void asyncUpdateCursorInfo(String ledgerName, String cursorName, ManagedC } store.put(path, content, Optional.of(expectedVersion)) - .thenAcceptAsync(optStat -> callback.operationComplete(null, optStat), executor.chooseThread(ledgerName)) + .thenAcceptAsync(optStat -> { + if (log.isDebugEnabled()) { + log.debug("[{}] Updating consumer {} on meta-data store with {} success", ledgerName, + cursorName, info); + } + callback.operationComplete(null, optStat); + }, executor.chooseThread(ledgerName)) .exceptionally(ex -> { - executor.executeOrdered(ledgerName, SafeRunnable.safeRun(() -> callback.operationFailed(getException(ex)))); + if (log.isDebugEnabled()) { + log.debug("[{}] Updating consumer {} on meta-data store with {} failed", ledgerName, cursorName, + info, ex); + } + executor.executeOrdered(ledgerName, + SafeRunnable.safeRun(() -> callback.operationFailed(getException(ex)))); return null; }); } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpFindNewest.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpFindNewest.java index 91567fc9bfb57..049fbfffeec72 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpFindNewest.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpFindNewest.java @@ -19,21 +19,21 @@ package org.apache.bookkeeper.mledger.impl; import com.google.common.base.Predicate; -import org.apache.bookkeeper.mledger.AsyncCallbacks.FindEntryCallback; -import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntryCallback; - import java.util.Optional; import lombok.extern.slf4j.Slf4j; - +import org.apache.bookkeeper.mledger.AsyncCallbacks.FindEntryCallback; +import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntryCallback; import org.apache.bookkeeper.mledger.Entry; +import org.apache.bookkeeper.mledger.ManagedCursor; +import org.apache.bookkeeper.mledger.ManagedLedger; +import org.apache.bookkeeper.mledger.ManagedLedger.PositionBound; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.Position; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.PositionBound; @Slf4j class OpFindNewest implements ReadEntryCallback { - private final ManagedCursorImpl cursor; - private final ManagedLedgerImpl ledger; + private final ManagedCursor cursor; + private final ManagedLedger ledger; private final PositionImpl startPosition; private final FindEntryCallback callback; private final Predicate condition; @@ -49,10 +49,10 @@ enum State { Position lastMatchedPosition = null; State state; - public OpFindNewest(ManagedCursorImpl cursor, PositionImpl startPosition, Predicate condition, - long numberOfEntries, FindEntryCallback callback, Object ctx) { + public OpFindNewest(ManagedCursor cursor, PositionImpl startPosition, Predicate condition, + long numberOfEntries, FindEntryCallback callback, Object ctx) { this.cursor = cursor; - this.ledger = cursor.ledger; + this.ledger = cursor.getManagedLedger(); this.startPosition = startPosition; this.callback = callback; this.condition = condition; diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ReadOnlyCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ReadOnlyCursorImpl.java index 7a0445ef5ebff..e260c20a9b4b3 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ReadOnlyCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ReadOnlyCursorImpl.java @@ -19,14 +19,12 @@ package org.apache.bookkeeper.mledger.impl; import com.google.common.collect.Range; - import lombok.extern.slf4j.Slf4j; - import org.apache.bookkeeper.client.BookKeeper; import org.apache.bookkeeper.mledger.AsyncCallbacks; +import org.apache.bookkeeper.mledger.ManagedLedger.PositionBound; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ReadOnlyCursor; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.PositionBound; import org.apache.bookkeeper.mledger.proto.MLDataFormats; @Slf4j diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java index 60043502023cb..ca9c62d998f5b 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java @@ -23,7 +23,6 @@ import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; import static org.testng.Assert.fail; - import com.google.common.base.Predicate; import com.google.common.collect.Lists; import com.google.common.collect.Range; @@ -93,6 +92,11 @@ public boolean hasMoreEntries() { return true; } + @Override + public boolean hasMoreEntries(PositionImpl position) { + return false; + } + @Override public long getNumberOfEntries() { return 0; diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryTest.java index 1381542e66ad9..acad30b9d8db7 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryTest.java @@ -19,23 +19,14 @@ package org.apache.bookkeeper.mledger.impl; import static org.testng.Assert.assertEquals; - -import org.apache.bookkeeper.conf.ClientConfiguration; -import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; -import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; -import org.apache.bookkeeper.mledger.ManagedLedgerFactory; import org.apache.bookkeeper.mledger.ManagedLedgerInfo; import org.apache.bookkeeper.mledger.ManagedLedgerInfo.CursorInfo; import org.apache.bookkeeper.mledger.ManagedLedgerInfo.MessageRangeInfo; import org.apache.bookkeeper.test.MockedBookKeeperTestCase; -import org.apache.bookkeeper.test.ZooKeeperUtil; -import org.testng.Assert; import org.testng.annotations.Test; -import java.util.List; - public class ManagedLedgerFactoryTest extends MockedBookKeeperTestCase { @Test(timeOut = 20000) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java index f174433b43f71..8283f4d08082c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java @@ -102,7 +102,15 @@ public void initialize(ServiceConfiguration conf, ZooKeeper zkClient, }; this.managedLedgerFactory = - new ManagedLedgerFactoryImpl(bkFactory, zkClient, managedLedgerFactoryConfig, statsLogger); + createManagedLedgerFactory(zkClient, managedLedgerFactoryConfig, statsLogger, bkFactory); + } + + protected ManagedLedgerFactoryImpl + createManagedLedgerFactory(ZooKeeper zkClient, + ManagedLedgerFactoryConfig managedLedgerFactoryConfig, + StatsLogger statsLogger, + BookkeeperFactoryForCustomEnsemblePlacementPolicy bkFactory) throws Exception { + return new ManagedLedgerFactoryImpl(bkFactory, zkClient, managedLedgerFactoryConfig, statsLogger); } public ManagedLedgerFactory getManagedLedgerFactory() { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index 159a6be03c989..65eb6e6781fd9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -53,12 +53,12 @@ import org.apache.bookkeeper.mledger.AsyncCallbacks.ManagedLedgerInfoCallback; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.LedgerOffloader; +import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.ManagedLedgerException.MetadataNotFoundException; import org.apache.bookkeeper.mledger.ManagedLedgerInfo; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerOfflineBacklog; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.commons.lang3.StringUtils; @@ -2193,7 +2193,7 @@ private void getEntryBatchSize(CompletableFuture batchSizeFuture, Persi MessageIdImpl messageId, int batchIndex) { if (batchIndex >= 0) { try { - ManagedLedgerImpl ledger = (ManagedLedgerImpl) topic.getManagedLedger(); + ManagedLedger ledger = topic.getManagedLedger(); ledger.asyncReadEntry(new PositionImpl(messageId.getLedgerId(), messageId.getEntryId()), new AsyncCallbacks.ReadEntryCallback() { @Override @@ -2278,7 +2278,7 @@ protected void internalGetMessageById(AsyncResponse asyncResponse, long ledgerId // will redirect if the topic not owned by current broker validateReadOperationOnTopic(authoritative); PersistentTopic topic = (PersistentTopic) getTopicReference(topicName); - ManagedLedgerImpl ledger = (ManagedLedgerImpl) topic.getManagedLedger(); + ManagedLedger ledger = topic.getManagedLedger(); ledger.asyncReadEntry(new PositionImpl(ledgerId, entryId), new AsyncCallbacks.ReadEntryCallback() { @Override public void readEntryFailed(ManagedLedgerException exception, Object ctx) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BacklogQuotaManager.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BacklogQuotaManager.java index c3f22815b060a..da89282b9b2b7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BacklogQuotaManager.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BacklogQuotaManager.java @@ -26,7 +26,7 @@ import java.util.concurrent.CompletableFuture; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.ManagedCursor.IndividualDeletedEntries; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; +import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.admin.AdminResource; import org.apache.pulsar.broker.service.persistent.PersistentTopic; @@ -132,7 +132,7 @@ private void dropBacklog(PersistentTopic persistentTopic, BacklogQuota quota) { // Get estimated unconsumed size for the managed ledger associated with this topic. Estimated size is more // useful than the actual storage size. Actual storage size gets updated only when managed ledger is trimmed. - ManagedLedgerImpl mLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); + ManagedLedger mLedger = persistentTopic.getManagedLedger(); long backlogSize = mLedger.getEstimatedBacklogSize(); if (log.isDebugEnabled()) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index d74af68b5c238..1aa96d3a534e1 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -48,9 +48,9 @@ import javax.net.ssl.SSLSession; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.Entry; +import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.Position; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.bookkeeper.mledger.util.SafeRun; import org.apache.commons.lang3.StringUtils; @@ -949,6 +949,7 @@ protected void handleSubscribe(final CommandSubscribe subscribe) { } if (schema != null) { + log.info("subscribe with schema {}", schema); return topic.addSchemaIfIdleOrCheckCompatible(schema) .thenCompose(v -> topic.subscribe( ServerCnx.this, subscriptionName, consumerId, @@ -957,6 +958,7 @@ protected void handleSubscribe(final CommandSubscribe subscribe) { readCompacted, initialPosition, startMessageRollbackDurationSec, isReplicated, keySharedMeta)); } else { + log.info("subscribe with out schema"); return topic.subscribe(ServerCnx.this, subscriptionName, consumerId, subType, priorityLevel, consumerName, isDurable, startMessageId, metadata, readCompacted, initialPosition, @@ -1604,7 +1606,7 @@ private void getLargestBatchIndexWhenPossible( String subscriptionName) { PersistentTopic persistentTopic = (PersistentTopic) topic; - ManagedLedgerImpl ml = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); + ManagedLedger ml = persistentTopic.getManagedLedger(); // If it's not pointing to a valid entry, respond messageId of the current position. if (lastPosition.getEntryId() == -1) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java index 3107f13c2865c..57f498afbce9d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java @@ -32,7 +32,6 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.Position; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.service.BrokerServiceException; @@ -338,7 +337,7 @@ private boolean removeConsumersFromRecentJoinedConsumers() { PositionImpl mdp = (PositionImpl) cursor.getMarkDeletedPosition(); if (mdp != null) { PositionImpl nextPositionOfTheMarkDeletePosition = - ((ManagedLedgerImpl) cursor.getManagedLedger()).getNextValidPosition(mdp); + cursor.getManagedLedger().getNextValidPosition(mdp); while (itr.hasNext()) { Map.Entry entry = itr.next(); if (entry.getValue().compareTo(nextPositionOfTheMarkDeletePosition) <= 0) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStreamingDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStreamingDispatcherMultipleConsumers.java index 9340e17aab22c..1939036184f7e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStreamingDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStreamingDispatcherMultipleConsumers.java @@ -26,7 +26,6 @@ import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.bookkeeper.mledger.util.SafeRun; import org.apache.pulsar.broker.service.Consumer; @@ -88,7 +87,7 @@ public synchronized void readEntryComplete(Entry entry, PendingReadEntryRequest log.debug("[{}] Distributing a messages to {} consumers", name, consumerList.size()); } - cursor.seek(((ManagedLedgerImpl) cursor.getManagedLedger()) + cursor.seek(cursor.getManagedLedger() .getNextValidPosition((PositionImpl) entry.getPosition())); sendMessagesToConsumers(readType, Lists.newArrayList(entry)); ctx.recycle(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStreamingDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStreamingDispatcherSingleActiveConsumer.java index b4e4ed37e46fc..a64824d21148f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStreamingDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStreamingDispatcherSingleActiveConsumer.java @@ -25,7 +25,6 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.bookkeeper.mledger.util.SafeRun; import org.apache.pulsar.broker.service.Consumer; @@ -163,7 +162,7 @@ public synchronized void internalReadEntryComplete(Entry entry, PendingReadEntry filterEntriesForConsumer(Lists.newArrayList(entry), batchSizes, sendMessageInfo, batchIndexesAcks, cursor, false); // Update cursor's read position. - cursor.seek(((ManagedLedgerImpl) cursor.getManagedLedger()) + cursor.seek(cursor.getManagedLedger() .getNextValidPosition((PositionImpl) entry.getPosition())); dispatchEntriesToConsumer(currentConsumer, Lists.newArrayList(entry), batchSizes, batchIndexesAcks, sendMessageInfo); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java index 4e05c58dceed8..7185fb459e4cc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java @@ -36,12 +36,12 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.ManagedCursor.IndividualDeletedEntries; +import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.ManagedLedgerException.ConcurrentFindCursorPositionException; import org.apache.bookkeeper.mledger.ManagedLedgerException.InvalidCursorPositionException; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.commons.lang3.tuple.MutablePair; import org.apache.pulsar.broker.intercept.BrokerInterceptor; @@ -380,7 +380,7 @@ public void acknowledgeMessage(List positions, AckType ackType, Map properties) { if (position != null) { - ManagedLedgerImpl managedLedger = ((ManagedLedgerImpl) cursor.getManagedLedger()); + ManagedLedger managedLedger = cursor.getManagedLedger(); PositionImpl nextPosition = managedLedger.getNextValidPosition(position); managedLedger.asyncReadEntry(nextPosition, new ReadEntryCallback() { @Override @@ -963,7 +963,7 @@ public SubscriptionStats getStats(Boolean getPreciseBacklog, boolean subscriptio } subStats.msgBacklog = getNumberOfEntriesInBacklog(getPreciseBacklog); if (subscriptionBacklogSize) { - subStats.backlogSize = ((ManagedLedgerImpl) topic.getManagedLedger()) + subStats.backlogSize = topic.getManagedLedger() .getEstimatedBacklogSize((PositionImpl) cursor.getMarkDeletedPosition()); } subStats.msgBacklogNoDelayed = subStats.msgBacklog - subStats.msgDelayed; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 56689a27b3f5b..6bf7c00120fce 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -62,7 +62,6 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.ManagedLedgerTerminatedException; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.ServiceConfiguration; @@ -403,32 +402,31 @@ private void asyncAddEntry(ByteBuf headersAndPayload, PublishContext publishCont } public void asyncReadEntry(PositionImpl position, AsyncCallbacks.ReadEntryCallback callback, Object ctx) { - if (ledger instanceof ManagedLedgerImpl) { - ((ManagedLedgerImpl) ledger).asyncReadEntry(position, callback, ctx); - } else { - callback.readEntryFailed(new ManagedLedgerException( - "Unexpected managedledger implementation, doesn't support " - + "direct read entry operation."), ctx); + try { + ledger.newNonDurableCursor(position).asyncReadEntries(1, new AsyncCallbacks.ReadEntriesCallback() { + @Override + public void readEntriesComplete(List entries, Object ctx) { + callback.readEntryComplete(entries.get(0), ctx); + } + + @Override + public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { + callback.readEntryFailed(exception, ctx); + + } + }, null, PositionImpl.latest); + } catch (ManagedLedgerException e) { + callback.readEntryFailed(e, null); } } - public PositionImpl getPositionAfterN(PositionImpl startPosition, long n) throws ManagedLedgerException { - if (ledger instanceof ManagedLedgerImpl) { - return ((ManagedLedgerImpl) ledger).getPositionAfterN(startPosition, n, - ManagedLedgerImpl.PositionBound.startExcluded); - } else { - throw new ManagedLedgerException("Unexpected managedledger implementation, doesn't support " - + "getPositionAfterN operation."); - } + public PositionImpl getPositionAfterN(PositionImpl startPosition, long n) { + return ledger.getPositionAfterN(startPosition, n, + ManagedLedger.PositionBound.startExcluded); } public PositionImpl getFirstPosition() throws ManagedLedgerException { - if (ledger instanceof ManagedLedgerImpl) { - return ((ManagedLedgerImpl) ledger).getFirstPosition(); - } else { - throw new ManagedLedgerException("Unexpected managedledger implementation, doesn't support " - + "getFirstPosition operation."); - } + return ledger.getFirstPosition(); } public long getNumberOfEntries() { @@ -1629,7 +1627,7 @@ public void updateRates(NamespaceStats nsStats, NamespaceBundleStats bundleStats topicStatsStream.writePair("msgThroughputOut", topicStatsHelper.aggMsgThroughputOut); topicStatsStream.writePair("storageSize", ledger.getTotalSize()); topicStatsStream.writePair("backlogSize", ledger.getEstimatedBacklogSize()); - topicStatsStream.writePair("pendingAddEntriesCount", ((ManagedLedgerImpl) ledger).getPendingAddEntriesCount()); + topicStatsStream.writePair("pendingAddEntriesCount", ledger.getPendingAddEntriesCount()); nsStats.msgRateIn += topicStatsHelper.aggMsgRateIn; nsStats.msgRateOut += topicStatsHelper.aggMsgRateOut; @@ -1641,7 +1639,7 @@ public void updateRates(NamespaceStats nsStats, NamespaceBundleStats bundleStats bundleStats.msgRateOut += topicStatsHelper.aggMsgRateOut; bundleStats.msgThroughputIn += topicStatsHelper.aggMsgThroughputIn; bundleStats.msgThroughputOut += topicStatsHelper.aggMsgThroughputOut; - bundleStats.cacheSize += ((ManagedLedgerImpl) ledger).getCacheSize(); + bundleStats.cacheSize += ledger.getCacheSize(); // Close topic object topicStatsStream.endObject(); @@ -1729,8 +1727,7 @@ public CompletableFuture getInternalStats(boolean CompletableFuture statFuture = new CompletableFuture<>(); PersistentTopicInternalStats stats = new PersistentTopicInternalStats(); - - ManagedLedgerImpl ml = (ManagedLedgerImpl) ledger; + ManagedLedger ml = ledger; stats.entriesAddedCounter = ml.getEntriesAddedCounter(); stats.numberOfEntries = ml.getNumberOfEntries(); stats.totalSize = ml.getTotalSize(); @@ -2365,7 +2362,7 @@ public CompletableFuture getLastMessageId() { .complete(new MessageIdImpl(position.getLedgerId(), position.getEntryId(), partitionIndex)); return completableFuture; } - ManagedLedgerImpl ledgerImpl = (ManagedLedgerImpl) ledger; + ManagedLedger ledgerImpl = ledger; if (!ledgerImpl.ledgerExists(position.getLedgerId())) { completableFuture .complete(MessageId.earliest); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/SchemaRegistryServiceImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/SchemaRegistryServiceImpl.java index d9dd91f34f168..cbb291ed0a983 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/SchemaRegistryServiceImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/SchemaRegistryServiceImpl.java @@ -138,15 +138,18 @@ public CompletableFuture>> getAllSchem @NotNull public CompletableFuture putSchemaIfAbsent(String schemaId, SchemaData schema, SchemaCompatibilityStrategy strategy) { + log.debug("call put schema if absent {} {} {}", schema, schema, strategy); + return trimDeletedSchemaAndGetList(schemaId).thenCompose(schemaAndMetadataList -> getSchemaVersionBySchemaData(schemaAndMetadataList, schema).thenCompose(schemaVersion -> { - if (schemaVersion != null) { - return CompletableFuture.completedFuture(schemaVersion); - } - CompletableFuture checkCompatibilityFuture = new CompletableFuture<>(); - if (schemaAndMetadataList.size() != 0) { - if (isTransitiveStrategy(strategy)) { - checkCompatibilityFuture = checkCompatibilityWithAll(schema, strategy, schemaAndMetadataList); + if (schemaVersion != null) { + return CompletableFuture.completedFuture(schemaVersion); + } + CompletableFuture checkCompatibilityFuture = new CompletableFuture<>(); + if (schemaAndMetadataList.size() != 0) { + if (isTransitiveStrategy(strategy)) { + checkCompatibilityFuture = checkCompatibilityWithAll(schema, strategy, + schemaAndMetadataList); } else { checkCompatibilityFuture = checkCompatibilityWithLatest(schemaId, schema, strategy); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/streamingdispatch/StreamingEntryReader.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/streamingdispatch/StreamingEntryReader.java index 24f9bcc236e72..1eacef174a393 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/streamingdispatch/StreamingEntryReader.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/streamingdispatch/StreamingEntryReader.java @@ -29,11 +29,11 @@ import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.Entry; +import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.WaitingEntryCallBack; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.bookkeeper.mledger.util.SafeRun; import org.apache.pulsar.broker.service.persistent.PersistentTopic; @@ -95,7 +95,7 @@ public synchronized void asyncReadEntries(int numEntriesToRead, int maxReadSizeB } PositionImpl nextReadPosition = (PositionImpl) cursor.getReadPosition(); - ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) cursor.getManagedLedger(); + ManagedLedger managedLedger = cursor.getManagedLedger(); // Edge case, when a old ledger is full and new ledger is not yet opened, position can point to next // position of the last confirmed position, but it'll be an invalid position. So try to update the position. if (!managedLedger.isValidPosition(nextReadPosition)) { @@ -272,9 +272,9 @@ private void retryReadRequest(PendingReadEntryRequest pendingReadEntryRequest, l // Jump again into dispatcher dedicated thread topic.getBrokerService().getTopicOrderedExecutor().executeOrdered(dispatcher.getName(), SafeRun.safeRun(() -> { - ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) cursor.getManagedLedger(); - managedLedger.asyncReadEntry(pendingReadEntryRequest.position, this, pendingReadEntryRequest); - })); + ManagedLedger managedLedger = cursor.getManagedLedger(); + managedLedger.asyncReadEntry(pendingReadEntryRequest.position, this, pendingReadEntryRequest); + })); }, delay, TimeUnit.MILLISECONDS); } @@ -290,7 +290,7 @@ private synchronized void internalEntriesAvailable() { log.debug("[{}} Streaming entry reader get notification of newly added entries from managed ledger," + " trying to issued pending read requests.", cursor.getName()); } - ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) cursor.getManagedLedger(); + ManagedLedger managedLedger = cursor.getManagedLedger(); List newlyIssuedRequests = new ArrayList<>(); if (!pendingReads.isEmpty()) { // Edge case, when a old ledger is full and new ledger is not yet opened, position can point to next diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/metrics/AbstractMetrics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/metrics/AbstractMetrics.java index 21b941ec5d728..ea4505e04379d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/metrics/AbstractMetrics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/metrics/AbstractMetrics.java @@ -25,9 +25,9 @@ import java.util.Map; import java.util.regex.Matcher; import java.util.regex.Pattern; +import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerFactoryMXBean; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerMBeanImpl; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.common.policies.data.TopicStats; @@ -97,7 +97,7 @@ protected ManagedLedgerFactoryMXBean getManagedLedgerCacheStats() { * * @return */ - protected Map getManagedLedgers() { + protected Map getManagedLedgers() { return ((ManagedLedgerFactoryImpl) pulsar.getManagedLedgerFactory()).getManagedLedgers(); } @@ -224,8 +224,8 @@ protected void populateMaxMap(Map map, String mkey, long value) { * @param metrics * @param ledger */ - protected void populateDimensionMap(Map> ledgersByDimensionMap, Metrics metrics, - ManagedLedgerImpl ledger) { + protected void populateDimensionMap(Map> ledgersByDimensionMap, Metrics metrics, + ManagedLedger ledger) { if (!ledgersByDimensionMap.containsKey(metrics)) { // create new list ledgersByDimensionMap.put(metrics, Lists.newArrayList(ledger)); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/metrics/ManagedCursorMetrics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/metrics/ManagedCursorMetrics.java index e888c933461c3..dfbe8c4b8087b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/metrics/ManagedCursorMetrics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/metrics/ManagedCursorMetrics.java @@ -25,9 +25,8 @@ import java.util.Map; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.ManagedCursorMXBean; -import org.apache.bookkeeper.mledger.impl.ManagedCursorContainer; +import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.common.stats.Metrics; @@ -55,13 +54,12 @@ public synchronized List generate() { */ private List aggregate() { metricsCollection.clear(); - for (Map.Entry e : getManagedLedgers().entrySet()) { + for (Map.Entry e : getManagedLedgers().entrySet()) { String ledgerName = e.getKey(); - ManagedLedgerImpl ledger = e.getValue(); + ManagedLedger ledger = e.getValue(); String namespace = parseNamespaceFromLedgerName(ledgerName); - ManagedCursorContainer cursorContainer = ledger.getCursors(); - Iterator cursorIterator = cursorContainer.iterator(); + Iterator cursorIterator = ledger.getCursors().iterator(); while (cursorIterator.hasNext()) { ManagedCursorImpl cursor = (ManagedCursorImpl) cursorIterator.next(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/metrics/ManagedLedgerMetrics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/metrics/ManagedLedgerMetrics.java index 273447cc0e4d8..79c940dc89de4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/metrics/ManagedLedgerMetrics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/metrics/ManagedLedgerMetrics.java @@ -23,16 +23,16 @@ import java.util.List; import java.util.Map; import java.util.Map.Entry; +import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerMXBean; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.common.stats.Metrics; public class ManagedLedgerMetrics extends AbstractMetrics { private List metricsCollection; - private Map> ledgersByDimensionMap; + private Map> ledgersByDimensionMap; // temp map to prepare aggregation metrics private Map tempAggregatedMetricsMap; @@ -57,20 +57,19 @@ public synchronized List generate() { * @param ledgersByDimension * @return */ - private List aggregate(Map> ledgersByDimension) { - + private List aggregate(Map> ledgersByDimension) { metricsCollection.clear(); - for (Entry> e : ledgersByDimension.entrySet()) { + for (Entry> e : ledgersByDimension.entrySet()) { Metrics metrics = e.getKey(); - List ledgers = e.getValue(); + List ledgers = e.getValue(); // prepare aggregation map tempAggregatedMetricsMap.clear(); // generate the collections by each metrics and then apply the aggregation - for (ManagedLedgerImpl ledger : ledgers) { + for (ManagedLedger ledger : ledgers) { ManagedLedgerMXBean lStats = ledger.getStats(); populateAggregationMapWithSum(tempAggregatedMetricsMap, "brk_ml_AddEntryBytesRate", @@ -134,17 +133,17 @@ private List aggregate(Map> ledgersByD * * @return */ - private Map> groupLedgersByDimension() { + private Map> groupLedgersByDimension() { ledgersByDimensionMap.clear(); // get the current topics statistics from StatsBrokerFilter // Map : topic-name->dest-stat - for (Entry e : getManagedLedgers().entrySet()) { + for (Entry e : getManagedLedgers().entrySet()) { String ledgerName = e.getKey(); - ManagedLedgerImpl ledger = e.getValue(); + ManagedLedger ledger = e.getValue(); // we want to aggregate by NS dimension diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ManagedLedgerMetricsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ManagedLedgerMetricsTest.java index cb8ae67811e5a..bf1778dc57a68 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ManagedLedgerMetricsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ManagedLedgerMetricsTest.java @@ -21,9 +21,8 @@ import java.util.List; import java.util.Map.Entry; import java.util.concurrent.TimeUnit; - +import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; -import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerMBeanImpl; import org.apache.pulsar.broker.service.BrokerTestBase; import org.apache.pulsar.broker.stats.metrics.ManagedLedgerMetrics; @@ -65,7 +64,7 @@ public void testManagedLedgerMetrics() throws Exception { producer.send(message.getBytes()); } - for (Entry ledger : ((ManagedLedgerFactoryImpl) pulsar.getManagedLedgerFactory()) + for (Entry ledger : ((ManagedLedgerFactoryImpl) pulsar.getManagedLedgerFactory()) .getManagedLedgers().entrySet()) { ManagedLedgerMBeanImpl stats = (ManagedLedgerMBeanImpl) ledger.getValue().getStats(); stats.refreshStats(1, TimeUnit.SECONDS); @@ -78,7 +77,7 @@ public void testManagedLedgerMetrics() throws Exception { String message = "my-message-" + i; producer.send(message.getBytes()); } - for (Entry ledger : ((ManagedLedgerFactoryImpl) pulsar.getManagedLedgerFactory()) + for (Entry ledger : ((ManagedLedgerFactoryImpl) pulsar.getManagedLedgerFactory()) .getManagedLedgers().entrySet()) { ManagedLedgerMBeanImpl stats = (ManagedLedgerMBeanImpl) ledger.getValue().getStats(); stats.refreshStats(1, TimeUnit.SECONDS); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/naming/TopicName.java b/pulsar-common/src/main/java/org/apache/pulsar/common/naming/TopicName.java index 8e30854d7c8a5..c96ad2178312f 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/naming/TopicName.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/naming/TopicName.java @@ -27,6 +27,8 @@ import java.util.List; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; +import java.util.stream.Stream; import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.common.util.Codec; import org.slf4j.Logger; @@ -78,7 +80,7 @@ public static TopicName get(String domain, String tenant, String namespace, Stri } public static TopicName get(String domain, String tenant, String cluster, String namespace, - String topic) { + String topic) { String name = domain + "://" + tenant + '/' + cluster + '/' + namespace + '/' + topic; return TopicName.get(name); } @@ -113,11 +115,11 @@ private TopicName(String completeTopicName) { completeTopicName = TopicDomain.persistent.name() + "://" + completeTopicName; } else if (parts.length == 1) { completeTopicName = TopicDomain.persistent.name() + "://" - + PUBLIC_TENANT + "/" + DEFAULT_NAMESPACE + "/" + parts[0]; + + PUBLIC_TENANT + "/" + DEFAULT_NAMESPACE + "/" + parts[0]; } else { throw new IllegalArgumentException( - "Invalid short topic name '" + completeTopicName + "', it should be in the format of " - + "// or "); + "Invalid short topic name '" + completeTopicName + "', it should be in the format of " + + "// or "); } } @@ -169,11 +171,11 @@ private TopicName(String completeTopicName) { } if (isV2()) { this.completeTopicName = String.format("%s://%s/%s/%s", - domain, tenant, namespacePortion, localName); + domain, tenant, namespacePortion, localName); } else { this.completeTopicName = String.format("%s://%s/%s/%s/%s", - domain, tenant, cluster, - namespacePortion, localName); + domain, tenant, cluster, + namespacePortion, localName); } } @@ -320,6 +322,33 @@ public String getPersistenceNamingEncoding() { } } + public static TopicName fromPersistenceNamingEncoding(String name) throws Exception { + if (name == null) { + return null; + } + final String[] arr = name.split("/"); + + if (arr.length == 4) { + String tenant = arr[0]; + String namespacePortion = arr[1]; + String domain = arr[2]; + String encodedLocalName = arr[3]; + final String decodedName = Codec.decode(encodedLocalName); + return TopicName.get(domain, tenant, namespacePortion, decodedName); + } else if (arr.length == 5) { + String tenant = arr[0]; + String cluster = arr[1]; + String namespacePortion = arr[2]; + String domain = arr[3]; + String encodedLocalName = arr[4]; + final String decodedName = Codec.decode(encodedLocalName); + return TopicName.get(domain, tenant, cluster, namespacePortion, decodedName); + } else { + log.error("arr.length = {}, arr = {}", arr.length, Stream.of(arr).collect(Collectors.toList())); + throw new Exception("not valid name: " + name); + } + } + /** * Get a string suitable for completeTopicName lookup. * @@ -344,8 +373,8 @@ public boolean isGlobal() { public String getSchemaName() { return getTenant() - + "/" + getNamespacePortion() - + "/" + TopicName.get(getPartitionedTopicName()).getEncodedLocalName(); + + "/" + getNamespacePortion() + + "/" + TopicName.get(getPartitionedTopicName()).getEncodedLocalName(); } @Override diff --git a/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/MockManagedLedger.java b/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/MockManagedLedger.java index fc4337689202e..547577c388544 100644 --- a/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/MockManagedLedger.java +++ b/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/MockManagedLedger.java @@ -21,6 +21,7 @@ import com.google.common.base.Predicate; import io.netty.buffer.ByteBuf; import java.util.Map; +import java.util.NavigableMap; import java.util.concurrent.CompletableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.AsyncCallbacks; @@ -31,12 +32,34 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.ManagedLedgerMXBean; import org.apache.bookkeeper.mledger.Position; +import org.apache.bookkeeper.mledger.WaitingEntryCallBack; +import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.bookkeeper.mledger.intercept.ManagedLedgerInterceptor; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo.LedgerInfo; import org.apache.pulsar.common.api.proto.CommandSubscribe; @Slf4j public class MockManagedLedger implements ManagedLedger { + @Override + public void initialize(ManagedLedgerInitializeLedgerCallback callback, Object ctx) { + + } + + @Override + public boolean isValidPosition(PositionImpl nextReadPosition) { + return false; + } + + @Override + public boolean hasMoreEntries(PositionImpl nextReadPosition) { + return false; + } + + @Override + public void addWaitingEntryCallBack(WaitingEntryCallBack streamingEntryReader) { + + } + @Override public String getName() { return null; @@ -57,6 +80,11 @@ public void asyncAddEntry(byte[] data, AsyncCallbacks.AddEntryCallback callback, } + @Override + public void asyncReadEntry(PositionImpl position, AsyncCallbacks.ReadEntryCallback callback, Object ctx) { + + } + @Override public Position addEntry(byte[] data, int offset, int length) throws InterruptedException, ManagedLedgerException { return null; @@ -168,6 +196,57 @@ public long getNumberOfEntries() { return 0; } + @Override + public long getEntriesAddedCounter() { + return 0; + } + + @Override + public long getLastLedgerCreatedTimestamp() { + return 0; + } + + @Override + public long getLastLedgerCreationFailureTimestamp() { + return 0; + } + + @Override + public int getWaitingCursorsCount() { + return 0; + } + + @Override + public long getCurrentLedgerEntries() { + return 0; + } + + @Override + public long getCurrentLedgerSize() { + return 0; + } + + + @Override + public NavigableMap getLedgersInfo() { + return null; + } + + @Override + public CompletableFuture getLedgerMetadata(long ledgerId) { + return null; + } + + @Override + public boolean ledgerExists(long ledgerId) { + return false; + } + + @Override + public void asyncDeleteLedgerFromBookKeeper(long ledgerId) { + + } + @Override public long getNumberOfActiveEntries() { return 0; @@ -183,6 +262,21 @@ public long getEstimatedBacklogSize() { return 0; } + @Override + public long getEstimatedBacklogSize(PositionImpl pos) { + return 0; + } + + @Override + public int getPendingAddEntriesCount() { + return 0; + } + + @Override + public long getCacheSize() { + return 0; + } + @Override public long getOffloadedSize() { return 0; @@ -213,6 +307,11 @@ public ManagedLedgerMXBean getStats() { return null; } + @Override + public void doCacheEviction(long maxTimestamp) { + + } + @Override public void delete() throws InterruptedException, ManagedLedgerException { @@ -258,6 +357,11 @@ public Position getLastConfirmedEntry() { return null; } + @Override + public String getState() { + return null; + } + @Override public void readyToCreateNewLedger() { @@ -315,6 +419,21 @@ public CompletableFuture asyncFindPosition(Predicate predicate) return null; } + @Override + public PositionImpl getPositionAfterN(PositionImpl startPosition, long n, PositionBound startRange) { + return null; + } + + @Override + public PositionImpl getFirstPosition() { + return null; + } + + @Override + public PositionImpl getLastPosition() { + return null; + } + @Override public ManagedLedgerInterceptor getManagedLedgerInterceptor() { return null; @@ -325,4 +444,9 @@ public CompletableFuture getLedgerInfo(long ledgerId) { final LedgerInfo build = LedgerInfo.newBuilder().setLedgerId(ledgerId).setSize(100).setEntries(20).build(); return CompletableFuture.completedFuture(build); } + + @Override + public PositionImpl getNextValidPosition(PositionImpl position) { + return null; + } }