From 61c0f577ee56f15126cad07ca7673af04ebb85c8 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Wed, 8 Mar 2017 09:51:23 -0800 Subject: [PATCH] Enable ML binary format from broker configuration --- .../mledger/ManagedLedgerFactoryConfig.java | 10 + .../mledger/impl/ManagedCursorImpl.java | 25 ++- .../mledger/impl/ManagedLedgerImpl.java | 2 +- .../mledger/impl/MetaStoreImplZookeeper.java | 183 ++++++++++++------ .../mledger/impl/ManagedCursorTest.java | 8 + .../ManagedLedgerBinaryFormatConversion.java | 106 ++++++++++ .../mledger/impl/ManagedLedgerTest.java | 8 + .../test/MockedBookKeeperTestCase.java | 14 +- 8 files changed, 286 insertions(+), 70 deletions(-) create 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 908d0ea14d3dd..b9bf306b5ef4b 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 @@ -21,6 +21,8 @@ public class ManagedLedgerFactoryConfig { private long maxCacheSize = 128 * MB; private double cacheEvictionWatermark = 0.90; + private boolean useProtobufBinaryFormatInZK = false; + public long getMaxCacheSize() { return maxCacheSize; } @@ -50,4 +52,12 @@ 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 6d1addf448349..cd5b60c35e9cd 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 @@ -18,6 +18,7 @@ import static com.google.common.base.Preconditions.checkArgument; import static com.google.common.base.Preconditions.checkNotNull; import static org.apache.bookkeeper.mledger.util.SafeRun.safeRun; +import static org.apache.bookkeeper.mledger.impl.MetaStoreImplZookeeper.ZNodeProtobufFormat; import java.util.ArrayDeque; import java.util.Collections; @@ -111,6 +112,8 @@ public class ManagedCursorImpl implements ManagedCursor { private final ReadWriteLock lock = new ReentrantReadWriteLock(); private final RateLimiter markDeleteLimiter; + + private final ZNodeProtobufFormat protobufFormat; class PendingMarkDeleteEntry { final PositionImpl newPosition; @@ -164,6 +167,9 @@ 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()); @@ -1666,21 +1672,20 @@ private void persistPositionMetaStore(long cursorsLedgerId, PositionImpl positio MetaStoreCallback callback) { // When closing we store the last mark-delete position in the z-node itself, so we won't need the cursor ledger, // hence we write it as -1. The cursor ledger is deleted once the z-node write is confirmed. - ManagedCursorInfo info = ManagedCursorInfo.newBuilder() // + ManagedCursorInfo.Builder info = ManagedCursorInfo.newBuilder() // .setCursorsLedgerId(cursorsLedgerId) // .setMarkDeleteLedgerId(position.getLedgerId()) // - .setMarkDeleteEntryId(position.getEntryId()) // + .setMarkDeleteEntryId(position.getEntryId()); // + + if (protobufFormat == ZNodeProtobufFormat.Binary) { + info.addAllIndividualDeletedMessages(buildIndividualDeletedMessageRanges()); + } - // Do not add individually deleted messages in text format since it would break - // backward compatibility. - // TODO: Add this again, when binary format is enabled - // .addAllIndividualDeletedMessages(buildIndividualDeletedMessageRanges()) // - .build(); if (log.isDebugEnabled()) { log.debug("[{}][{}] Closing cursor at md-position: {}", ledger.getName(), name, markDeletePosition); } - ledger.getStore().asyncUpdateCursorInfo(ledger.getName(), name, info, cursorLedgerVersion, + ledger.getStore().asyncUpdateCursorInfo(ledger.getName(), name, info.build(), cursorLedgerVersion, new MetaStoreCallback() { @Override public void operationComplete(Void result, Version version) { @@ -1705,7 +1710,8 @@ public void asyncClose(final AsyncCallbacks.CloseCallback callback, final Object lock.readLock().lock(); try { - if (!individualDeletedMessages.isEmpty()) { + 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. @@ -1921,6 +1927,7 @@ void persistPosition(final LedgerHandle lh, final PositionImpl position, final V position); } + checkNotNull(lh); lh.asyncAddEntry(pi.toByteArray(), (rc, lh1, entryId, ctx) -> { if (rc == BKException.Code.OK) { if (log.isDebugEnabled()) { 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 1e1e472a63657..0b1b0106bdf7f 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 @@ -170,7 +170,7 @@ enum PositionBound { private final ScheduledExecutorService scheduledExecutor; private final OrderedSafeExecutor executor; - private final ManagedLedgerFactoryImpl factory; + final ManagedLedgerFactoryImpl factory; protected final ManagedLedgerMBeanImpl mbean; /** 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 bb4dc3a9f828f..c878907960a8e 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 @@ -42,7 +42,11 @@ import com.google.protobuf.TextFormat; import com.google.protobuf.TextFormat.ParseException; -class MetaStoreImplZookeeper implements MetaStore { +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; @@ -51,6 +55,7 @@ class MetaStoreImplZookeeper implements MetaStore { private static final String prefix = prefixName + "/"; private final ZooKeeper zk; + private final ZNodeProtobufFormat protobufFormat; private final OrderedSafeExecutor executor; private static class ZKVersion implements Version { @@ -62,7 +67,13 @@ private static class ZKVersion implements Version { } public MetaStoreImplZookeeper(ZooKeeper zk, OrderedSafeExecutor executor) throws Exception { + this(zk, ZNodeProtobufFormat.Text, executor); + } + + public MetaStoreImplZookeeper(ZooKeeper zk, ZNodeProtobufFormat protobufFormat, OrderedSafeExecutor executor) + throws Exception { this.zk = zk; + this.protobufFormat = protobufFormat; this.executor = executor; if (zk.exists(prefixName, false) == null) { @@ -100,12 +111,10 @@ public void getManagedLedgerInfo(final String ledgerName, final MetaStoreCallbac zk.getData(prefix + ledgerName, false, (rc, path, ctx, readData, stat) -> executor.submit(safeRun(() -> { if (rc == Code.OK.intValue()) { try { - ManagedLedgerInfo.Builder builder = ManagedLedgerInfo.newBuilder(); - TextFormat.merge(new String(readData, Encoding), builder); - ManagedLedgerInfo info = builder.build(); + ManagedLedgerInfo info = parseManagedLedgerInfo(readData); info = updateMLInfoTimestamp(info); callback.operationComplete(info, new ZKVersion(stat.getVersion())); - } catch (ParseException e) { + } catch (ParseException | InvalidProtocolBufferException e) { callback.operationFailed(new MetaStoreException(e)); } } else if (rc == Code.NONODE.intValue()) { @@ -116,16 +125,14 @@ public void getManagedLedgerInfo(final String ledgerName, final MetaStoreCallbac ManagedLedgerInfo info = ManagedLedgerInfo.getDefaultInstance(); callback.operationComplete(info, new ZKVersion(0)); } else { - callback.operationFailed( - new MetaStoreException(KeeperException.create(Code.get(rc1)))); + callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc1)))); } }; - ZkUtils.asyncCreateFullPathOptimistic(zk, prefix + ledgerName, new byte[0], Acl, - CreateMode.PERSISTENT, createcb, null); + ZkUtils.asyncCreateFullPathOptimistic(zk, prefix + ledgerName, new byte[0], Acl, CreateMode.PERSISTENT, + createcb, null); } else { - callback.operationFailed( - new MetaStoreException(KeeperException.create(Code.get(rc)))); + callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc)))); } })), null); } @@ -139,23 +146,28 @@ public void asyncUpdateLedgerIds(String ledgerName, ManagedLedgerInfo mlInfo, Ve log.debug("[{}] Updating metadata version={} with content={}", ledgerName, zkVersion.version, mlInfo); } - zk.setData(prefix + ledgerName, mlInfo.toString().getBytes(Encoding), zkVersion.version, (rc, path, zkCtx, stat) -> executor.submit(safeRun(() -> { - if (log.isDebugEnabled()) { - log.debug("[{}] UpdateLedgersIdsCallback.processResult rc={} newVersion={}", ledgerName, - Code.get(rc), stat != null ? stat.getVersion() : "null"); - } - MetaStoreException status = null; - if (rc == Code.BADVERSION.intValue()) { - // Content has been modified on ZK since our last read - status = new BadVersionException(KeeperException.create(Code.get(rc))); - callback.operationFailed(status); - } else if (rc != Code.OK.intValue()) { - status = new MetaStoreException(KeeperException.create(Code.get(rc))); - callback.operationFailed(status); - } else { - callback.operationComplete(null, new ZKVersion(stat.getVersion())); - } - })), null); + byte[] serializedMlInfo = protobufFormat == ZNodeProtobufFormat.Text ? // + mlInfo.toString().getBytes(Encoding) : // Text format + mlInfo.toByteArray(); // Binary format + + zk.setData(prefix + ledgerName, serializedMlInfo, zkVersion.version, + (rc, path, zkCtx, stat) -> executor.submit(safeRun(() -> { + if (log.isDebugEnabled()) { + log.debug("[{}] UpdateLedgersIdsCallback.processResult rc={} newVersion={}", ledgerName, + Code.get(rc), stat != null ? stat.getVersion() : "null"); + } + MetaStoreException status = null; + if (rc == Code.BADVERSION.intValue()) { + // Content has been modified on ZK since our last read + status = new BadVersionException(KeeperException.create(Code.get(rc))); + callback.operationFailed(status); + } else if (rc != Code.OK.intValue()) { + status = new MetaStoreException(KeeperException.create(Code.get(rc))); + callback.operationFailed(status); + } else { + callback.operationComplete(null, new ZKVersion(stat.getVersion())); + } + })), null); } @Override @@ -168,8 +180,7 @@ public void getCursors(final String ledgerName, final MetaStoreCallback executor.submit(safeRun(() -> { if (rc != Code.OK.intValue()) { - callback.operationFailed( - new MetaStoreException(KeeperException.create(Code.get(rc)))); + callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc)))); } else { try { - ManagedCursorInfo.Builder info = ManagedCursorInfo.newBuilder(); - TextFormat.merge(new String(data, Encoding), info); - callback.operationComplete(info.build(), new ZKVersion(stat.getVersion())); - } catch (ParseException e) { + ManagedCursorInfo info = parseManagedCursorInfo(data); + callback.operationComplete(info, new ZKVersion(stat.getVersion())); + } catch (ParseException | InvalidProtocolBufferException e) { callback.operationFailed(new MetaStoreException(e)); } } @@ -216,26 +225,28 @@ public void asyncUpdateCursorInfo(final String ledgerName, final String cursorNa info.getCursorsLedgerId(), info.getMarkDeleteLedgerId(), info.getMarkDeleteEntryId()); String path = prefix + ledgerName + "/" + cursorName; - byte[] content = info.toString().getBytes(Encoding); + byte[] content = protobufFormat == ZNodeProtobufFormat.Text ? // + info.toString().getBytes(Encoding) : // Text format + info.toByteArray(); // Binary format if (version == null) { if (log.isDebugEnabled()) { log.debug("[{}] Creating consumer {} on meta-data store with {}", ledgerName, cursorName, info); } - zk.create(path, content, Acl, CreateMode.PERSISTENT, (rc, path1, ctx, name) -> executor.submit(safeRun(() -> { - if (rc != Code.OK.intValue()) { - log.warn("[{}] Error creating cosumer {} node on meta-data store with {}: ", ledgerName, - cursorName, info, Code.get(rc)); - callback.operationFailed( - new MetaStoreException(KeeperException.create(Code.get(rc)))); - } else { - if (log.isDebugEnabled()) { - log.debug("[{}] Created consumer {} on meta-data store with {}", ledgerName, cursorName, - info); - } - callback.operationComplete(null, new ZKVersion(0)); - } - })), null); + zk.create(path, content, Acl, CreateMode.PERSISTENT, + (rc, path1, ctx, name) -> executor.submit(safeRun(() -> { + if (rc != Code.OK.intValue()) { + log.warn("[{}] Error creating cosumer {} node on meta-data store with {}: ", ledgerName, + cursorName, info, Code.get(rc)); + callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc)))); + } else { + if (log.isDebugEnabled()) { + log.debug("[{}] Created consumer {} on meta-data store with {}", ledgerName, cursorName, + info); + } + callback.operationComplete(null, new ZKVersion(0)); + } + })), null); } else { ZKVersion zkVersion = (ZKVersion) version; if (log.isDebugEnabled()) { @@ -243,11 +254,9 @@ public void asyncUpdateCursorInfo(final String ledgerName, final String cursorNa } zk.setData(path, content, zkVersion.version, (rc, path1, ctx, stat) -> executor.submit(safeRun(() -> { if (rc == Code.BADVERSION.intValue()) { - callback.operationFailed( - new BadVersionException(KeeperException.create(Code.get(rc)))); + callback.operationFailed(new BadVersionException(KeeperException.create(Code.get(rc)))); } else if (rc != Code.OK.intValue()) { - callback.operationFailed( - new MetaStoreException(KeeperException.create(Code.get(rc)))); + callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc)))); } else { callback.operationComplete(null, new ZKVersion(stat.getVersion())); } @@ -266,8 +275,7 @@ public void asyncRemoveCursor(final String ledgerName, final String consumerName if (rc == Code.OK.intValue()) { callback.operationComplete(null, null); } else { - callback.operationFailed( - new MetaStoreException(KeeperException.create(Code.get(rc)))); + callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc)))); } })), null); } @@ -282,8 +290,7 @@ public void removeManagedLedger(String ledgerName, MetaStoreCallback callb if (rc == Code.OK.intValue()) { callback.operationComplete(null, null); } else { - callback.operationFailed( - new MetaStoreException(KeeperException.create(Code.get(rc)))); + callback.operationFailed(new MetaStoreException(KeeperException.create(Code.get(rc)))); } })), null); } @@ -297,5 +304,63 @@ 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); + } + } + } + + 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); + } + } + } + + 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 2572dbe73d75e..0fc94a6d950e3 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 @@ -47,6 +47,7 @@ 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.MetaStoreImplZookeeper.ZNodeProtobufFormat; import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ManagedLedgerException; @@ -56,6 +57,7 @@ import org.apache.bookkeeper.test.MockedBookKeeperTestCase; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.testng.annotations.Factory; import org.testng.annotations.Test; import com.google.common.base.Charsets; @@ -65,6 +67,12 @@ 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 { 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 new file mode 100644 index 0000000000000..8d8d9acee0915 --- /dev/null +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerBinaryFormatConversion.java @@ -0,0 +1,106 @@ +/** + * Copyright 2016 Yahoo Inc. + * + * Licensed 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/ManagedLedgerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java index f5e27b90be237..973058f77f885 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 @@ -59,6 +59,7 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback; import org.apache.bookkeeper.mledger.impl.MetaStore.Version; +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; @@ -69,6 +70,7 @@ 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; @@ -87,6 +89,12 @@ 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"); 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 1e1091416b470..b224197355382 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 @@ -21,7 +21,9 @@ import org.apache.bookkeeper.client.MockBookKeeper; 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; @@ -29,6 +31,7 @@ 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. @@ -50,6 +53,8 @@ public abstract class MockedBookKeeperTestCase { protected OrderedSafeExecutor executor; protected ExecutorService cachedExecutor; + + protected ZNodeProtobufFormat protobufFormat = ZNodeProtobufFormat.Text; public MockedBookKeeperTestCase() { // By default start a 3 bookies cluster @@ -59,6 +64,11 @@ 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 { @@ -73,7 +83,9 @@ public void setUp(Method method) throws Exception { executor = new OrderedSafeExecutor(2, "test"); cachedExecutor = Executors.newCachedThreadPool(); - factory = new ManagedLedgerFactoryImpl(bkc, zkc); + ManagedLedgerFactoryConfig conf = new ManagedLedgerFactoryConfig(); + conf.setUseProtobufBinaryFormatInZK(protobufFormat == ZNodeProtobufFormat.Binary); + factory = new ManagedLedgerFactoryImpl(bkc, zkc, conf); } @AfterMethod