From 8f2a23b31008fb0219e2f9da44a1ad0596b00b14 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Tue, 2 May 2017 16:50:51 -0700 Subject: [PATCH 1/2] Remove Managed Ledger metadata text format --- .../mledger/ManagedLedgerFactoryConfig.java | 10 -- .../mledger/impl/ManagedCursorImpl.java | 50 ++------ .../mledger/impl/MetaStoreImplZookeeper.java | 83 +++---------- .../mledger/impl/ManagedCursorTest.java | 41 ++++--- .../ManagedLedgerBinaryFormatConversion.java | 109 ------------------ .../impl/ManagedLedgerFactoryTest.java | 32 +---- .../mledger/impl/ManagedLedgerTest.java | 10 +- .../test/MockedBookKeeperTestCase.java | 10 -- 8 files changed, 59 insertions(+), 286 deletions(-) delete mode 100644 managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerBinaryFormatConversion.java diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java index 29f94eca5575a..f534657c6bece 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java @@ -24,8 +24,6 @@ public class ManagedLedgerFactoryConfig { private long maxCacheSize = 128 * MB; private double cacheEvictionWatermark = 0.90; - private boolean useProtobufBinaryFormatInZK = false; - public long getMaxCacheSize() { return maxCacheSize; } @@ -55,12 +53,4 @@ public ManagedLedgerFactoryConfig setCacheEvictionWatermark(double cacheEviction return this; } - public boolean useProtobufBinaryFormatInZK() { - return useProtobufBinaryFormatInZK; - } - - public void setUseProtobufBinaryFormatInZK(boolean useProtobufBinaryFormatInZK) { - this.useProtobufBinaryFormatInZK = useProtobufBinaryFormatInZK; - } - } 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 8e02c9d2f6536..52bad592b8bd9 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 @@ -60,7 +60,6 @@ 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.impl.MetaStoreImplZookeeper.ZNodeProtobufFormat; import org.apache.bookkeeper.mledger.proto.MLDataFormats; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedCursorInfo; import org.apache.bookkeeper.mledger.proto.MLDataFormats.PositionInfo; @@ -91,15 +90,18 @@ public class ManagedCursorImpl implements ManagedCursor { protected static final AtomicReferenceFieldUpdater WAITING_READ_OP_UPDATER = AtomicReferenceFieldUpdater.newUpdater(ManagedCursorImpl.class, OpReadEntry.class, "waitingReadOp"); + @SuppressWarnings("unused") private volatile OpReadEntry waitingReadOp = null; private static final int FALSE = 0; private static final int TRUE = 1; private static final AtomicIntegerFieldUpdater RESET_CURSOR_IN_PROGRESS_UPDATER = AtomicIntegerFieldUpdater .newUpdater(ManagedCursorImpl.class, "resetCursorInProgress"); + @SuppressWarnings("unused") private volatile int resetCursorInProgress = FALSE; private static final AtomicIntegerFieldUpdater PENDING_READ_OPS_UPDATER = AtomicIntegerFieldUpdater.newUpdater(ManagedCursorImpl.class, "pendingReadOps"); + @SuppressWarnings("unused") private volatile int pendingReadOps = 0; // This counters are used to compute the numberOfEntries and numberOfEntriesInBacklog values, without having to look @@ -117,8 +119,6 @@ public class ManagedCursorImpl implements ManagedCursor { private final RateLimiter markDeleteLimiter; - private final ZNodeProtobufFormat protobufFormat; - class PendingMarkDeleteEntry { final PositionImpl newPosition; final MarkDeleteCallback callback; @@ -139,6 +139,7 @@ public PendingMarkDeleteEntry(PositionImpl newPosition, MarkDeleteCallback callb private final ArrayDeque pendingMarkDeleteOps = new ArrayDeque<>(); private static final AtomicIntegerFieldUpdater PENDING_MARK_DELETED_SUBMITTED_COUNT_UPDATER = AtomicIntegerFieldUpdater.newUpdater(ManagedCursorImpl.class, "pendingMarkDeletedSubmittedCount"); + @SuppressWarnings("unused") private volatile int pendingMarkDeletedSubmittedCount = 0; private long lastLedgerSwitchTimestamp; @@ -171,9 +172,6 @@ public interface VoidCallback { RESET_CURSOR_IN_PROGRESS_UPDATER.set(this, FALSE); WAITING_READ_OP_UPDATER.set(this, null); this.lastLedgerSwitchTimestamp = System.currentTimeMillis(); - this.protobufFormat = ledger.factory.getConfig().useProtobufBinaryFormatInZK() ? // - ZNodeProtobufFormat.Binary : // - ZNodeProtobufFormat.Text; if (config.getThrottleMarkDelete() > 0.0) { markDeleteLimiter = RateLimiter.create(config.getThrottleMarkDelete()); @@ -1234,7 +1232,7 @@ public void asyncMarkDelete(final Position position, final MarkDeleteCallback ca if (RESET_CURSOR_IN_PROGRESS_UPDATER.get(this) == TRUE) { if (log.isDebugEnabled()) { log.debug("[{}] cursor reset in progress - ignoring mark delete on position [{}] for cursor [{}]", - ledger.getName(), (PositionImpl) position, name); + ledger.getName(), position, name); } callback.markDeleteFailed( new ManagedLedgerException("Reset cursor in progress - unable to mark delete position " @@ -1674,9 +1672,7 @@ private void persistPositionMetaStore(long cursorsLedgerId, PositionImpl positio .setMarkDeleteLedgerId(position.getLedgerId()) // .setMarkDeleteEntryId(position.getEntryId()); // - if (protobufFormat == ZNodeProtobufFormat.Binary) { - info.addAllIndividualDeletedMessages(buildIndividualDeletedMessageRanges()); - } + info.addAllIndividualDeletedMessages(buildIndividualDeletedMessageRanges()); if (log.isDebugEnabled()) { log.debug("[{}][{}] Closing cursor at md-position: {}", ledger.getName(), name, markDeletePosition); @@ -1705,36 +1701,6 @@ public void asyncClose(final AsyncCallbacks.CloseCallback callback, final Object return; } - lock.readLock().lock(); - try { - if (cursorLedger != null && protobufFormat == ZNodeProtobufFormat.Text - && !individualDeletedMessages.isEmpty()) { - // To save individualDeletedMessages status, we don't want to dump the information in text format into - // the z-node. Until we switch to binary format, just flush the mark-delete + the - // individualDeletedMessages into the ledger. - persistPosition(cursorLedger, markDeletePosition, new VoidCallback() { - @Override - public void operationComplete() { - cursorLedger.asyncClose(new CloseCallback() { - @Override - public void closeComplete(int rc, LedgerHandle lh, Object ctx) { - callback.closeComplete(ctx); - } - }, ctx); - } - - @Override - public void operationFailed(ManagedLedgerException exception) { - callback.closeFailed(exception, ctx); - } - }); - - return; - } - } finally { - lock.readLock().unlock(); - } - persistPositionMetaStore(-1, markDeletePosition, new MetaStoreCallback() { @Override public void operationComplete(Void result, Stat stat) { @@ -2182,7 +2148,7 @@ public String getIndividuallyDeletedMessages() { /** * Checks given position is part of deleted-range and returns next position of upper-end as all the messages are * deleted up to that point - * + * * @param position * @return next available position */ @@ -2194,7 +2160,7 @@ public PositionImpl getNextAvailablePosition(PositionImpl position) { } return position.getNext(); } - + public boolean isIndividuallyDeletedEntriesEmpty() { lock.readLock().lock(); try { 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 index 4c2d8afd4fb5a..005b34a97a62d 100644 --- 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 @@ -47,10 +47,6 @@ public class MetaStoreImplZookeeper implements MetaStore { - public static enum ZNodeProtobufFormat { - Text, Binary - } - private static final Charset Encoding = Charsets.UTF_8; private static final List Acl = ZooDefs.Ids.OPEN_ACL_UNSAFE; @@ -58,7 +54,6 @@ public static enum ZNodeProtobufFormat { private static final String prefix = prefixName + "/"; private final ZooKeeper zk; - private final ZNodeProtobufFormat protobufFormat; private final OrderedSafeExecutor executor; private static class ZKStat implements Stat { @@ -94,14 +89,9 @@ public long getModificationTimestamp() { } } - public MetaStoreImplZookeeper(ZooKeeper zk, OrderedSafeExecutor executor) throws Exception { - this(zk, ZNodeProtobufFormat.Text, executor); - } - - public MetaStoreImplZookeeper(ZooKeeper zk, ZNodeProtobufFormat protobufFormat, OrderedSafeExecutor executor) + public MetaStoreImplZookeeper(ZooKeeper zk, OrderedSafeExecutor executor) throws Exception { this.zk = zk; - this.protobufFormat = protobufFormat; this.executor = executor; if (zk.exists(prefixName, false) == null) { @@ -177,9 +167,7 @@ public void asyncUpdateLedgerIds(String ledgerName, ManagedLedgerInfo mlInfo, St log.debug("[{}] Updating metadata version={} with content={}", ledgerName, zkStat.version, mlInfo); } - byte[] serializedMlInfo = protobufFormat == ZNodeProtobufFormat.Text ? // - mlInfo.toString().getBytes(Encoding) : // Text format - mlInfo.toByteArray(); // Binary format + byte[] serializedMlInfo = mlInfo.toByteArray(); // Binary format zk.setData(prefix + ledgerName, serializedMlInfo, zkStat.getVersion(), (rc, path, zkCtx, stat1) -> executor.submit(safeRun(() -> { @@ -255,9 +243,7 @@ public void asyncUpdateCursorInfo(final String ledgerName, final String cursorNa info.getCursorsLedgerId(), info.getMarkDeleteLedgerId(), info.getMarkDeleteEntryId()); String path = prefix + ledgerName + "/" + cursorName; - byte[] content = protobufFormat == ZNodeProtobufFormat.Text ? // - info.toString().getBytes(Encoding) : // Text format - info.toByteArray(); // Binary format + byte[] content = info.toByteArray(); // Binary format if (stat == null) { if (log.isDebugEnabled()) { @@ -336,60 +322,29 @@ public Iterable getManagedLedgers() throws MetaStoreException { private ManagedLedgerInfo parseManagedLedgerInfo(byte[] data) throws ParseException, InvalidProtocolBufferException { - if (protobufFormat == ZNodeProtobufFormat.Text) { - // First try text format, then fallback to binary - try { - return parseManagedLedgerInfoFromText(data); - } catch (ParseException e) { - return parseManagedLedgerInfoFromBinary(data); - } - } else { - // First try binary format, then fallback to text - try { - return parseManagedLedgerInfoFromBinary(data); - } catch (InvalidProtocolBufferException e) { - return parseManagedLedgerInfoFromText(data); - } + // 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 ManagedLedgerInfo parseManagedLedgerInfoFromText(byte[] data) throws ParseException { - ManagedLedgerInfo.Builder builder = ManagedLedgerInfo.newBuilder(); - TextFormat.merge(new String(data, Encoding), builder); - return builder.build(); - } - - private ManagedLedgerInfo parseManagedLedgerInfoFromBinary(byte[] data) throws InvalidProtocolBufferException { - return ManagedLedgerInfo.newBuilder().mergeFrom(data).build(); - } - private ManagedCursorInfo parseManagedCursorInfo(byte[] data) throws ParseException, InvalidProtocolBufferException { - if (protobufFormat == ZNodeProtobufFormat.Text) { - // First try text format, then fallback to binary - try { - return parseManagedCursorInfoFromText(data); - } catch (ParseException e) { - return parseManagedCursorInfoFromBinary(data); - } - } else { - // First try binary format, then fallback to text - try { - return parseManagedCursorInfoFromBinary(data); - } catch (InvalidProtocolBufferException e) { - return parseManagedCursorInfoFromText(data); - } + // 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(); } - } - - private ManagedCursorInfo parseManagedCursorInfoFromText(byte[] data) throws ParseException { - ManagedCursorInfo.Builder builder = ManagedCursorInfo.newBuilder(); - TextFormat.merge(new String(data, Encoding), builder); - return builder.build(); - } - private ManagedCursorInfo parseManagedCursorInfoFromBinary(byte[] data) throws InvalidProtocolBufferException { - return ManagedCursorInfo.newBuilder().mergeFrom(data).build(); } private static final Logger log = LoggerFactory.getLogger(MetaStoreImplZookeeper.class); 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 d5ee1f765f3cd..dc4152fea3d40 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 @@ -43,8 +43,6 @@ import java.util.stream.Collectors; import org.apache.bookkeeper.client.BKException; -import org.apache.bookkeeper.client.BookKeeper.DigestType; -import org.apache.bookkeeper.client.LedgerHandle; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.AsyncCallbacks.AddEntryCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.MarkDeleteCallback; @@ -52,21 +50,16 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.ManagedCursor.IndividualDeletedEntries; -import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.VoidCallback; -import org.apache.bookkeeper.mledger.impl.MetaStoreImplZookeeper.ZNodeProtobufFormat; import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.ManagedLedgerFactory; import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig; import org.apache.bookkeeper.mledger.Position; +import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.VoidCallback; import org.apache.bookkeeper.test.MockedBookKeeperTestCase; -import org.apache.zookeeper.CreateMode; -import org.apache.zookeeper.ZooDefs; -import org.apache.zookeeper.data.Stat; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.testng.annotations.Factory; import org.testng.annotations.Test; import com.google.common.base.Charsets; @@ -76,12 +69,6 @@ public class ManagedCursorTest extends MockedBookKeeperTestCase { private static final Charset Encoding = Charsets.UTF_8; - - @Factory(dataProvider = "protobufFormat") - public ManagedCursorTest(ZNodeProtobufFormat protobufFormat) { - super(); - this.protobufFormat = protobufFormat; - } @Test(timeOut = 20000) void readFromEmptyLedger() throws Exception { @@ -315,6 +302,7 @@ void asyncReadWithoutErrors() throws Exception { final CountDownLatch counter = new CountDownLatch(1); cursor.asyncReadEntries(100, new ReadEntriesCallback() { + @Override public void readEntriesComplete(List entries, Object ctx) { assertNull(ctx); assertEquals(entries.size(), 1); @@ -322,6 +310,7 @@ public void readEntriesComplete(List entries, Object ctx) { counter.countDown(); } + @Override public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { fail(exception.getMessage()); } @@ -343,11 +332,13 @@ void asyncReadWithErrors() throws Exception { stopBookKeeper(); cursor.asyncReadEntries(100, new ReadEntriesCallback() { + @Override public void readEntriesComplete(List entries, Object ctx) { entries.forEach(e -> e.release()); counter.countDown(); } + @Override public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { fail("async-call should not have failed"); } @@ -364,10 +355,12 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { final CountDownLatch counter2 = new CountDownLatch(1); cursor.asyncReadEntries(100, new ReadEntriesCallback() { + @Override public void readEntriesComplete(List entries, Object ctx) { fail("async-call should have failed"); } + @Override public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter2.countDown(); } @@ -389,10 +382,12 @@ void asyncReadWithInvalidParameter() throws Exception { stopBookKeeper(); cursor.asyncReadEntries(0, new ReadEntriesCallback() { + @Override public void readEntriesComplete(List entries, Object ctx) { fail("async-call should have failed"); } + @Override public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter.countDown(); } @@ -534,6 +529,7 @@ void testConcurrentResetCursor() throws Exception { final int idx = i; futures.add(executor.submit(new Callable() { + @Override public AtomicBoolean call() throws Exception { barrier.await(); @@ -845,15 +841,19 @@ public void asyncMarkDeleteBlocking() throws Exception { final CountDownLatch latch = new CountDownLatch(N); for (int i = 0; i < N; i++) { ledger.asyncAddEntry("entry".getBytes(Encoding), new AddEntryCallback() { + @Override public void addFailed(ManagedLedgerException exception, Object ctx) { } + @Override public void addComplete(Position position, Object ctx) { lastPosition.set(position); c1.asyncMarkDelete(position, new MarkDeleteCallback() { + @Override public void markDeleteFailed(ManagedLedgerException exception, Object ctx) { } + @Override public void markDeleteComplete(Object ctx) { latch.countDown(); } @@ -893,10 +893,12 @@ void cursorPersistenceAsyncMarkDeleteSameThread() throws Exception { final CountDownLatch latch = new CountDownLatch(N); for (final Position p : positions) { c1.asyncMarkDelete(p, new MarkDeleteCallback() { + @Override public void markDeleteComplete(Object ctx) { latch.countDown(); } + @Override public void markDeleteFailed(ManagedLedgerException exception, Object ctx) { log.error("Failed to markdelete", exception); latch.countDown(); @@ -944,20 +946,24 @@ void unorderedAsyncMarkDelete() throws Exception { final CountDownLatch latch = new CountDownLatch(2); c1.asyncMarkDelete(p2, new MarkDeleteCallback() { + @Override public void markDeleteFailed(ManagedLedgerException exception, Object ctx) { fail(); } + @Override public void markDeleteComplete(Object ctx) { latch.countDown(); } }, null); c1.asyncMarkDelete(p1, new MarkDeleteCallback() { + @Override public void markDeleteFailed(ManagedLedgerException exception, Object ctx) { latch.countDown(); } + @Override public void markDeleteComplete(Object ctx) { fail(); } @@ -1462,12 +1468,14 @@ void testReadEntriesOrWait() throws Exception { ManagedCursor c = ledger.openCursor("c" + i); c.asyncReadEntriesOrWait(1, new ReadEntriesCallback() { + @Override public void readEntriesComplete(List entries, Object ctx) { assertEquals(entries.size(), 1); entries.forEach(e -> e.release()); counter.countDown(); } + @Override public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { log.error("Error reading", exception); } @@ -1493,6 +1501,7 @@ void testReadEntriesOrWaitBlocking() throws Exception { final ManagedCursor cursor = ledger.openCursor("c" + i); futures.add(executor.submit(new Callable() { + @Override public Void call() throws Exception { barrier.await(); @@ -2188,8 +2197,8 @@ void testGetEntryAfterN() throws Exception { assertNull(e); // check that the mark delete and read positions have not been updated after all the previous operations - assertEquals((PositionImpl) c1.getMarkDeletedPosition(), new PositionImpl(currentLedger, -1)); - assertEquals((PositionImpl) c1.getReadPosition(), new PositionImpl(currentLedger, 4)); + assertEquals(c1.getMarkDeletedPosition(), new PositionImpl(currentLedger, -1)); + assertEquals(c1.getReadPosition(), new PositionImpl(currentLedger, 4)); c1.markDelete(pos4); assertEquals(c1.getMarkDeletedPosition(), pos4); diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerBinaryFormatConversion.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerBinaryFormatConversion.java deleted file mode 100644 index 3260e487c08a0..0000000000000 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerBinaryFormatConversion.java +++ /dev/null @@ -1,109 +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.testng.Assert.assertEquals; - -import java.nio.charset.Charset; - -import org.apache.bookkeeper.mledger.ManagedCursor; -import org.apache.bookkeeper.mledger.ManagedLedger; -import org.apache.bookkeeper.mledger.ManagedLedgerFactory; -import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig; -import org.apache.bookkeeper.mledger.Position; -import org.apache.bookkeeper.test.MockedBookKeeperTestCase; -import org.testng.annotations.DataProvider; -import org.testng.annotations.Test; - -import com.google.common.base.Charsets; - -public class ManagedLedgerBinaryFormatConversion extends MockedBookKeeperTestCase { - - private static final Charset Encoding = Charsets.UTF_8; - - @DataProvider(name = "gracefulClose") - public static Object[][] protobufFormat() { - return new Object[][] { { false }, { true } }; - } - - @Test(timeOut = 20000, dataProvider = "gracefulClose") - void textToBinary(boolean gracefulClose) throws Exception { - ManagedLedgerFactoryConfig textConf = new ManagedLedgerFactoryConfig(); - textConf.setUseProtobufBinaryFormatInZK(false); - ManagedLedgerFactory textFactory = new ManagedLedgerFactoryImpl(bkc, zkc, textConf); - ManagedLedger ledger = textFactory.open("my_test_ledger"); - - ManagedCursor c1 = ledger.openCursor("c1"); - ledger.addEntry("test-0".getBytes(Encoding)); - Position p1 = ledger.addEntry("test-1".getBytes(Encoding)); - ledger.addEntry("test-2".getBytes(Encoding)); - - c1.delete(p1); - - if (gracefulClose) { - ledger.close(); - } - - // Reopen with binary format - ManagedLedgerFactoryConfig binaryConf = new ManagedLedgerFactoryConfig(); - binaryConf.setUseProtobufBinaryFormatInZK(true); - ManagedLedgerFactory binaryFactory = new ManagedLedgerFactoryImpl(bkc, zkc, binaryConf); - ledger = binaryFactory.open("my_test_ledger"); - c1 = ledger.openCursor("c1"); - - // The 'p1' entry was already deleted - assertEquals(c1.getNumberOfEntriesInBacklog(), 2); - - textFactory.shutdown(); - binaryFactory.shutdown(); - } - - @Test(timeOut = 20000, dataProvider = "gracefulClose") - void binaryToText(boolean gracefulClose) throws Exception { - ManagedLedgerFactoryConfig binaryConf = new ManagedLedgerFactoryConfig(); - binaryConf.setUseProtobufBinaryFormatInZK(true); - ManagedLedgerFactory binaryFactory = new ManagedLedgerFactoryImpl(bkc, zkc, binaryConf); - ManagedLedger ledger = binaryFactory.open("my_test_ledger"); - ManagedCursor c1 = ledger.openCursor("c1"); - - ledger.addEntry("test-0".getBytes(Encoding)); - Position p1 = ledger.addEntry("test-1".getBytes(Encoding)); - ledger.addEntry("test-2".getBytes(Encoding)); - - c1.delete(p1); - - if (gracefulClose) { - ledger.close(); - } - - // Reopen with binary format - - ManagedLedgerFactoryConfig textConf = new ManagedLedgerFactoryConfig(); - textConf.setUseProtobufBinaryFormatInZK(false); - ManagedLedgerFactory textFactory = new ManagedLedgerFactoryImpl(bkc, zkc, textConf); - ledger = textFactory.open("my_test_ledger"); - c1 = ledger.openCursor("c1"); - - // The 'p1' entry was already deleted - assertEquals(c1.getNumberOfEntriesInBacklog(), 2); - - textFactory.shutdown(); - binaryFactory.shutdown(); - } -} 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 eea0c399b0c17..b704db7168f3b 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,33 +19,17 @@ package org.apache.bookkeeper.mledger.impl; import static org.testng.Assert.assertEquals; -import static org.testng.Assert.assertNull; - -import java.nio.charset.Charset; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ManagedLedgerInfo; -import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.ManagedLedgerInfo.CursorInfo; import org.apache.bookkeeper.mledger.ManagedLedgerInfo.MessageRangeInfo; -import org.apache.bookkeeper.mledger.impl.MetaStoreImplZookeeper.ZNodeProtobufFormat; import org.apache.bookkeeper.test.MockedBookKeeperTestCase; -import org.testng.annotations.Factory; import org.testng.annotations.Test; -import com.google.common.base.Charsets; - public class ManagedLedgerFactoryTest extends MockedBookKeeperTestCase { - private static final Charset Encoding = Charsets.UTF_8; - - @Factory(dataProvider = "protobufFormat") - public ManagedLedgerFactoryTest(ZNodeProtobufFormat protobufFormat) { - super(); - this.protobufFormat = protobufFormat; - } - @Test(timeOut = 20000) public void testGetManagedLedgerInfoWithClose() throws Exception { ManagedLedgerConfig conf = new ManagedLedgerConfig(); @@ -78,17 +62,13 @@ public void testGetManagedLedgerInfoWithClose() throws Exception { assertEquals(cursorInfo.markDelete.ledgerId, 3); assertEquals(cursorInfo.markDelete.entryId, -1); - if (protobufFormat == ZNodeProtobufFormat.Binary) { - assertEquals(cursorInfo.individualDeletedMessages.size(), 1); + assertEquals(cursorInfo.individualDeletedMessages.size(), 1); - MessageRangeInfo mri = cursorInfo.individualDeletedMessages.get(0); - assertEquals(mri.from.ledgerId, p1.getLedgerId()); - assertEquals(mri.from.entryId, p1.getEntryId()); - assertEquals(mri.to.ledgerId, p3.getLedgerId()); - assertEquals(mri.to.entryId, p3.getEntryId()); - } else { - assertNull(cursorInfo.individualDeletedMessages); - } + MessageRangeInfo mri = cursorInfo.individualDeletedMessages.get(0); + assertEquals(mri.from.ledgerId, p1.getLedgerId()); + assertEquals(mri.from.entryId, p1.getEntryId()); + assertEquals(mri.to.ledgerId, p3.getLedgerId()); + assertEquals(mri.to.entryId, p3.getEntryId()); } } 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 28612c818843e..20e6efb4477d3 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 @@ -62,7 +62,6 @@ 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.impl.MetaStoreImplZookeeper.ZNodeProtobufFormat; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo.LedgerInfo; import org.apache.bookkeeper.mledger.util.Pair; @@ -76,7 +75,6 @@ import org.apache.zookeeper.ZooKeeper; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.testng.annotations.Factory; import org.testng.annotations.Test; import com.google.common.base.Charsets; @@ -92,12 +90,6 @@ public class ManagedLedgerTest extends MockedBookKeeperTestCase { private static final Charset Encoding = Charsets.UTF_8; - @Factory(dataProvider = "protobufFormat") - public ManagedLedgerTest(ZNodeProtobufFormat protobufFormat) { - super(); - this.protobufFormat = protobufFormat; - } - @Test public void managedLedgerApi() throws Exception { ManagedLedger ledger = factory.open("my_test_ledger"); @@ -118,7 +110,7 @@ public void managedLedgerApi() throws Exception { // Acknowledge only on last entry Entry lastEntry = entries.get(entries.size() - 1); cursor.markDelete(lastEntry.getPosition()); - + for (Entry entry : entries) { log.info("Read entry. Position={} Content='{}'", entry.getPosition(), new String(entry.getData())); entry.release(); diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/test/MockedBookKeeperTestCase.java b/managed-ledger/src/test/java/org/apache/bookkeeper/test/MockedBookKeeperTestCase.java index 09c0a0daeee0e..040fb27ca2fb7 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/test/MockedBookKeeperTestCase.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/test/MockedBookKeeperTestCase.java @@ -26,7 +26,6 @@ import org.apache.bookkeeper.conf.ClientConfiguration; import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; -import org.apache.bookkeeper.mledger.impl.MetaStoreImplZookeeper.ZNodeProtobufFormat; import org.apache.bookkeeper.util.OrderedSafeExecutor; import org.apache.bookkeeper.util.ZkUtils; import org.apache.zookeeper.MockZooKeeper; @@ -34,7 +33,6 @@ import org.slf4j.LoggerFactory; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; -import org.testng.annotations.DataProvider; /** * A class runs several bookie servers for testing. @@ -56,8 +54,6 @@ public abstract class MockedBookKeeperTestCase { protected OrderedSafeExecutor executor; protected ExecutorService cachedExecutor; - - protected ZNodeProtobufFormat protobufFormat = ZNodeProtobufFormat.Text; public MockedBookKeeperTestCase() { // By default start a 3 bookies cluster @@ -67,11 +63,6 @@ public MockedBookKeeperTestCase() { public MockedBookKeeperTestCase(int numBookies) { this.numBookies = numBookies; } - - @DataProvider(name = "protobufFormat") - public static Object[][] protobufFormat() { - return new Object[][] { { ZNodeProtobufFormat.Text }, { ZNodeProtobufFormat.Binary } }; - } @BeforeMethod public void setUp(Method method) throws Exception { @@ -87,7 +78,6 @@ public void setUp(Method method) throws Exception { executor = new OrderedSafeExecutor(2, "test"); cachedExecutor = Executors.newCachedThreadPool(); ManagedLedgerFactoryConfig conf = new ManagedLedgerFactoryConfig(); - conf.setUseProtobufBinaryFormatInZK(protobufFormat == ZNodeProtobufFormat.Binary); factory = new ManagedLedgerFactoryImpl(bkc, zkc, conf); } From ec7a7e2416800c47f6c8e619331d706b28bb797a Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Wed, 3 May 2017 11:23:12 -0700 Subject: [PATCH 2/2] MockZookeeper should store values as byte[] and not convert to String --- .../org/apache/zookeeper/MockZooKeeper.java | 25 ++++++++++--------- 1 file changed, 13 insertions(+), 12 deletions(-) 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 31af9c5b0572b..51041b4752cdb 100644 --- a/managed-ledger/src/test/java/org/apache/zookeeper/MockZooKeeper.java +++ b/managed-ledger/src/test/java/org/apache/zookeeper/MockZooKeeper.java @@ -53,7 +53,7 @@ @SuppressWarnings({ "deprecation", "restriction", "rawtypes" }) public class MockZooKeeper extends ZooKeeper { - private TreeMap> tree; + private TreeMap> tree; private SetMultimap watchers; private volatile boolean stopped; private boolean alwaysFail = false; @@ -149,7 +149,7 @@ public String create(String path, byte[] data, List acl, CreateMode createM } if (createMode == CreateMode.EPHEMERAL_SEQUENTIAL || createMode == CreateMode.PERSISTENT_SEQUENTIAL) { - String parentData = tree.get(parent).first; + byte[] parentData = tree.get(parent).first; int parentVersion = tree.get(parent).second; path = path + parentVersion; @@ -157,7 +157,7 @@ public String create(String path, byte[] data, List acl, CreateMode createM tree.put(parent, Pair.create(parentData, parentVersion + 1)); } - tree.put(path, Pair.create(new String(data), 0)); + tree.put(path, Pair.create(data, 0)); if (!parent.isEmpty()) { final Set toNotifyParent = Sets.newHashSet(); @@ -201,7 +201,7 @@ public void create(final String path, final byte[] data, final List acl, Cr mutex.unlock(); cb.processResult(KeeperException.Code.NONODE.intValue(), path, ctx, null); } else { - tree.put(path, Pair.create(new String(data), 0)); + tree.put(path, Pair.create(data, 0)); mutex.unlock(); cb.processResult(0, path, ctx, null); if (!parent.isEmpty()) { @@ -218,7 +218,7 @@ public byte[] getData(String path, Watcher watcher, Stat stat) throws KeeperExce mutex.lock(); try { checkProgrammedFail(); - Pair value = tree.get(path); + Pair value = tree.get(path); if (value == null) { throw new KeeperException.NoNodeException(path); } else { @@ -228,7 +228,7 @@ public byte[] getData(String path, Watcher watcher, Stat stat) throws KeeperExce if (stat != null) { stat.setVersion(value.second); } - return value.first.getBytes(); + return value.first; } } finally { mutex.unlock(); @@ -247,7 +247,7 @@ public void getData(final String path, boolean watch, final DataCallback cb, fin return; } - Pair value; + Pair value; mutex.lock(); try { value = tree.get(path); @@ -260,7 +260,7 @@ public void getData(final String path, boolean watch, final DataCallback cb, fin } else { Stat stat = new Stat(); stat.setVersion(value.second); - cb.processResult(0, path, ctx, value.first.getBytes(), stat); + cb.processResult(0, path, ctx, value.first, stat); } }); } @@ -280,7 +280,7 @@ public void getData(final String path, final Watcher watcher, final DataCallback return; } - Pair value = tree.get(path); + Pair value = tree.get(path); if (value == null) { mutex.unlock(); cb.processResult(KeeperException.Code.NONODE.intValue(), path, ctx, null, null); @@ -292,7 +292,7 @@ public void getData(final String path, final Watcher watcher, final DataCallback Stat stat = new Stat(); stat.setVersion(value.second); mutex.unlock(); - cb.processResult(0, path, ctx, value.first.getBytes(), stat); + cb.processResult(0, path, ctx, value.first, stat); } }); } @@ -493,6 +493,7 @@ public Stat exists(String path, Watcher watcher) throws KeeperException, Interru } } + @Override public void exists(String path, boolean watch, StatCallback cb, Object ctx) { executor.execute(() -> { mutex.lock(); @@ -559,7 +560,7 @@ public Stat setData(final String path, byte[] data, int version) throws KeeperEx newVersion = currentVersion + 1; log.debug("[{}] Updating -- current version: {}", path, currentVersion); - tree.put(path, Pair.create(new String(data), newVersion)); + tree.put(path, Pair.create(data, newVersion)); toNotify.addAll(watchers.get(path)); watchers.removeAll(path); @@ -617,7 +618,7 @@ public void setData(final String path, final byte[] data, int version, final Sta int newVersion = currentVersion + 1; log.debug("[{}] Updating -- current version: {}", path, currentVersion); - tree.put(path, Pair.create(new String(data), newVersion)); + tree.put(path, Pair.create(data, newVersion)); Stat stat = new Stat(); stat.setVersion(newVersion);