From 979c37a5a0a5562b20e93f269f3282bfc20d5dbb Mon Sep 17 00:00:00 2001 From: coderzc Date: Wed, 7 Dec 2022 20:43:58 +0800 Subject: [PATCH 1/5] add config `metadataFsyncEnabled` --- conf/standalone.conf | 2 ++ .../pulsar/broker/ServiceConfiguration.java | 8 ++++++++ .../org/apache/pulsar/broker/PulsarService.java | 1 + .../metadata/api/MetadataStoreConfig.java | 6 ++++++ .../metadata/impl/RocksdbMetadataStore.java | 17 +++++++---------- .../pulsar/metadata/MetadataBenchmark.java | 3 ++- .../bookkeeper/PulsarLedgerIdGeneratorTest.java | 4 ++-- 7 files changed, 28 insertions(+), 13 deletions(-) diff --git a/conf/standalone.conf b/conf/standalone.conf index 9f22bb8dbe460..511f045e6fa70 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -29,6 +29,8 @@ metadataStoreUrl= # Configuration file path for metadata store. It's supported by RocksdbMetadataStore and EtcdMetadataStore for now metadataStoreConfigPath= +# Whether we should enable fsync for local metadata store. It's supported by RocksdbMetadataStore for now +metadataFsyncEnabled=true # The metadata store URL for the configuration data. If empty, we fall back to use metadataStoreUrl configurationMetadataStoreUrl= diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 1b6bdc9986dd0..5fab0cd6e20e8 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -513,6 +513,14 @@ The delayed message index bucket time step(in seconds) in per bucket snapshot se private String metadataStoreConfigPath = null; + @FieldContext( + category = CATEGORY_SERVER, + doc = "Whether we should enable fsync for local metadata store. It's supported by RocksdbMetadataStore " + + "for now." + ) + private boolean metadataFsyncEnabled = true; + + @FieldContext( dynamic = true, category = CATEGORY_SERVER, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index 0a49d1092d3dc..37c740f55302a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -362,6 +362,7 @@ public MetadataStore createConfigurationMetadataStore(PulsarMetadataEventSynchro .batchingMaxDelayMillis(config.getMetadataStoreBatchingMaxDelayMillis()) .batchingMaxOperations(config.getMetadataStoreBatchingMaxOperations()) .batchingMaxSizeKb(config.getMetadataStoreBatchingMaxSizeKb()) + .fsyncEnable(config.isMetadataFsyncEnabled()) .metadataStoreName(MetadataStoreConfig.CONFIGURATION_METADATA_STORE) .synchronizer(synchronizer) .build()); diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreConfig.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreConfig.java index b742d03f9dbd5..5ddfe33c3912a 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreConfig.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreConfig.java @@ -81,6 +81,12 @@ public class MetadataStoreConfig { @Builder.Default private final String metadataStoreName = ""; + /** + * Whether we should enable fsync for local metadata store, It's supported by RocksdbMetadataStore for now. + */ + @Builder.Default + private final boolean fsyncEnable = true; + /** * Pluggable MetadataEventSynchronizer to sync metadata events across the * separate clusters. diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/RocksdbMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/RocksdbMetadataStore.java index d06d35db7acc4..42f807e1c3b6a 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/RocksdbMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/RocksdbMetadataStore.java @@ -87,8 +87,7 @@ public class RocksdbMetadataStore extends AbstractMetadataStore { private final ReentrantReadWriteLock dbStateLock; private volatile State state; - private final WriteOptions optionSync; - private final WriteOptions optionDontSync; + private final WriteOptions writeOptions; private final ReadOptions optionCache; private final ReadOptions optionDontCache; private MetadataEventSynchronizer synchronizer; @@ -235,8 +234,7 @@ private RocksdbMetadataStore(String metadataURL, MetadataStoreConfig metadataSto db = openDB(dataPath.toString(), metadataStoreConfig.getConfigFilePath()); - this.optionSync = new WriteOptions().setSync(true); - this.optionDontSync = new WriteOptions().setSync(false); + this.writeOptions = new WriteOptions().setSync(metadataStoreConfig.isFsyncEnable()); this.optionCache = new ReadOptions().setFillCache(true); this.optionDontCache = new ReadOptions().setFillCache(false); @@ -261,7 +259,7 @@ private long loadInstanceId() throws RocksDBException { } else { instanceId = 0; } - db.put(optionSync, INSTANCE_ID_KEY, toBytes(instanceId)); + db.put(writeOptions, INSTANCE_ID_KEY, toBytes(instanceId)); return instanceId; } @@ -271,7 +269,7 @@ private AtomicLong loadSequentialIdGenerator() throws RocksDBException { if (value != null) { generator.set(toLong(value)); } else { - db.put(optionSync, INSTANCE_ID_KEY, toBytes(generator.get())); + db.put(writeOptions, INSTANCE_ID_KEY, toBytes(generator.get())); } return generator; } @@ -369,8 +367,7 @@ public synchronized void close() throws MetadataStoreException { state = State.CLOSED; log.info("close.instanceId={}", instanceId); db.close(); - optionSync.close(); - optionDontSync.close(); + writeOptions.close(); optionCache.close(); optionDontCache.close(); super.close(); @@ -496,7 +493,7 @@ protected CompletableFuture storeDelete(String path, Optional expect if (state == State.CLOSED) { throw new MetadataStoreException.AlreadyClosedException(""); } - try (Transaction transaction = db.beginTransaction(optionSync)) { + try (Transaction transaction = db.beginTransaction(writeOptions)) { byte[] pathBytes = toBytes(path); byte[] oldValueData = transaction.getForUpdate(optionDontCache, pathBytes, true); MetaValue metaValue = MetaValue.parse(oldValueData); @@ -535,7 +532,7 @@ protected CompletableFuture storePut(String path, byte[] data, Optional urlSupplier) throw @Test(dataProvider = "impl", enabled = false) public void testPut(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); final int N_KEYS = 10_000; final int N_PUTS = 100_000; diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/PulsarLedgerIdGeneratorTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/PulsarLedgerIdGeneratorTest.java index 09bb85b140c6b..73d5f451c1ff1 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/PulsarLedgerIdGeneratorTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/PulsarLedgerIdGeneratorTest.java @@ -49,8 +49,8 @@ public class PulsarLedgerIdGeneratorTest extends BaseMetadataStoreTest { @Test(dataProvider = "impl") public void testGenerateLedgerId(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStoreExtended store = - MetadataStoreExtended.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); @Cleanup PulsarLedgerIdGenerator ledgerIdGenerator = new PulsarLedgerIdGenerator(store, "/ledgers"); From 8ce26824403803cc2613503ccc838169c7ce3b6d Mon Sep 17 00:00:00 2001 From: coderzc Date: Thu, 8 Dec 2022 17:20:53 +0800 Subject: [PATCH 2/5] improve doc --- conf/standalone.conf | 3 +++ .../java/org/apache/pulsar/broker/ServiceConfiguration.java | 5 ++++- 2 files changed, 7 insertions(+), 1 deletion(-) diff --git a/conf/standalone.conf b/conf/standalone.conf index 511f045e6fa70..ccf0acbd50b1c 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -30,6 +30,9 @@ metadataStoreUrl= metadataStoreConfigPath= # Whether we should enable fsync for local metadata store. It's supported by RocksdbMetadataStore for now +# If this flag is true, metadata writes will be slower. +# If this flag is false, and the machine crashes, some recent metadata writes may be lost. +# Note that if it is just the process that crashes (i.e., the machine does not reboot), no writes will be lost even if it is false. metadataFsyncEnabled=true # The metadata store URL for the configuration data. If empty, we fall back to use metadataStoreUrl diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 5fab0cd6e20e8..88b9d77181df0 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -516,7 +516,10 @@ The delayed message index bucket time step(in seconds) in per bucket snapshot se @FieldContext( category = CATEGORY_SERVER, doc = "Whether we should enable fsync for local metadata store. It's supported by RocksdbMetadataStore " - + "for now." + + "for now. If this flag is true, metadata writes will be slower. " + + "If this flag is false, and the machine crashes, some recent metadata writes may be lost. " + + "Note that if it is just the process that crashes (i.e., the machine does not reboot), " + + "no writes will be lost even if it is false." ) private boolean metadataFsyncEnabled = true; From c63ab68a38901fcb5b3eba6106522a7550b7d680 Mon Sep 17 00:00:00 2001 From: coderzc Date: Fri, 9 Dec 2022 16:25:43 +0800 Subject: [PATCH 3/5] remove metadataFsyncEnabled from ServiceConfiguration --- conf/standalone.conf | 6 ------ .../apache/pulsar/broker/ServiceConfiguration.java | 12 ------------ 2 files changed, 18 deletions(-) diff --git a/conf/standalone.conf b/conf/standalone.conf index ccf0acbd50b1c..03e9562657c82 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -29,12 +29,6 @@ metadataStoreUrl= # Configuration file path for metadata store. It's supported by RocksdbMetadataStore and EtcdMetadataStore for now metadataStoreConfigPath= -# Whether we should enable fsync for local metadata store. It's supported by RocksdbMetadataStore for now -# If this flag is true, metadata writes will be slower. -# If this flag is false, and the machine crashes, some recent metadata writes may be lost. -# Note that if it is just the process that crashes (i.e., the machine does not reboot), no writes will be lost even if it is false. -metadataFsyncEnabled=true - # The metadata store URL for the configuration data. If empty, we fall back to use metadataStoreUrl configurationMetadataStoreUrl= diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 88b9d77181df0..63ad1bb73525a 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -512,18 +512,6 @@ The delayed message index bucket time step(in seconds) in per bucket snapshot se ) private String metadataStoreConfigPath = null; - - @FieldContext( - category = CATEGORY_SERVER, - doc = "Whether we should enable fsync for local metadata store. It's supported by RocksdbMetadataStore " - + "for now. If this flag is true, metadata writes will be slower. " - + "If this flag is false, and the machine crashes, some recent metadata writes may be lost. " - + "Note that if it is just the process that crashes (i.e., the machine does not reboot), " - + "no writes will be lost even if it is false." - ) - private boolean metadataFsyncEnabled = true; - - @FieldContext( dynamic = true, category = CATEGORY_SERVER, From 2b5ab947762be1c11cb43528c260e47602b97ef2 Mon Sep 17 00:00:00 2001 From: coderzc Date: Fri, 9 Dec 2022 16:37:26 +0800 Subject: [PATCH 4/5] Make more test use `fsync=false` --- .../apache/pulsar/metadata/CounterTest.java | 5 +-- .../pulsar/metadata/LeaderElectionTest.java | 8 ++-- .../pulsar/metadata/LockManagerTest.java | 12 +++--- .../pulsar/metadata/MetadataBenchmark.java | 3 +- .../pulsar/metadata/MetadataCacheTest.java | 3 +- .../pulsar/metadata/MetadataStoreTest.java | 42 ++++++++++++------- .../LedgerUnderreplicationManagerTest.java | 3 +- .../PulsarRegistrationClientTest.java | 19 ++++----- .../PulsarRegistrationManagerTest.java | 3 +- 9 files changed, 56 insertions(+), 42 deletions(-) diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/CounterTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/CounterTest.java index e642bb4fea7ac..ead80a0287348 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/CounterTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/CounterTest.java @@ -25,7 +25,6 @@ import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotEquals; import static org.testng.Assert.fail; - import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.atomic.AtomicInteger; @@ -45,7 +44,7 @@ public class CounterTest extends BaseMetadataStoreTest { public void basicTest(String provider, Supplier urlSupplier) throws Exception { @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), - MetadataStoreConfig.builder().build()); + MetadataStoreConfig.builder().fsyncEnable(false).build()); @Cleanup CoordinationService cs1 = new CoordinationServiceImpl(store); @@ -102,7 +101,7 @@ public void testCounterDoesNotAutoReset(String provider, Supplier urlSup public void testGetNextCounterRetry(String provider, Supplier urlSupplier) throws Exception { @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), - MetadataStoreConfig.builder().build()); + MetadataStoreConfig.builder().fsyncEnable(false).build()); MetadataStoreExtended spy = spy(store); diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LeaderElectionTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LeaderElectionTest.java index c187d20ed2d40..632c0c0fedf56 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LeaderElectionTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LeaderElectionTest.java @@ -43,7 +43,7 @@ public class LeaderElectionTest extends BaseMetadataStoreTest { public void basicTest(String provider, Supplier urlSupplier) throws Exception { @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), - MetadataStoreConfig.builder().build()); + MetadataStoreConfig.builder().fsyncEnable(false).build()); @Cleanup CoordinationService coordinationService = new CoordinationServiceImpl(store); @@ -133,7 +133,7 @@ public void multipleMembers(String provider, Supplier urlSupplier) throw public void leaderNodeIsDeletedExternally(String provider, Supplier urlSupplier) throws Exception { @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), - MetadataStoreConfig.builder().build()); + MetadataStoreConfig.builder().fsyncEnable(false).build()); @Cleanup CoordinationService coordinationService = new CoordinationServiceImpl(store); @@ -161,7 +161,7 @@ public void leaderNodeIsDeletedExternally(String provider, Supplier urlS public void closeAll(String provider, Supplier urlSupplier) throws Exception { @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), - MetadataStoreConfig.builder().build()); + MetadataStoreConfig.builder().fsyncEnable(false).build()); MetadataCache cache = store.getMetadataCache(String.class); CoordinationService cs = new CoordinationServiceImpl(store); @@ -191,7 +191,7 @@ public void closeAll(String provider, Supplier urlSupplier) throws Excep public void revalidateLeaderWithinSameSession(String provider, Supplier urlSupplier) throws Exception { @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), - MetadataStoreConfig.builder().build()); + MetadataStoreConfig.builder().fsyncEnable(false).build()); String path = newKey(); diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java index be0b0448e9f36..250e9f02dd5df 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java @@ -53,7 +53,7 @@ public class LockManagerTest extends BaseMetadataStoreTest { public void acquireLocks(String provider, Supplier urlSupplier) throws Exception { @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), - MetadataStoreConfig.builder().build()); + MetadataStoreConfig.builder().fsyncEnable(false).build()); @Cleanup CoordinationService coordinationService = new CoordinationServiceImpl(store); @@ -104,7 +104,7 @@ public void acquireLocks(String provider, Supplier urlSupplier) throws E public void cleanupOnClose(String provider, Supplier urlSupplier) throws Exception { @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), - MetadataStoreConfig.builder().build()); + MetadataStoreConfig.builder().fsyncEnable(false).build()); @Cleanup CoordinationService coordinationService = new CoordinationServiceImpl(store); @@ -135,7 +135,7 @@ public void cleanupOnClose(String provider, Supplier urlSupplier) throws public void updateValue(String provider, Supplier urlSupplier) throws Exception { @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), - MetadataStoreConfig.builder().build()); + MetadataStoreConfig.builder().fsyncEnable(false).build()); MetadataCache cache = store.getMetadataCache(String.class); @@ -159,7 +159,7 @@ public void updateValue(String provider, Supplier urlSupplier) throws Ex public void updateValueWhenVersionIsOutOfSync(String provider, Supplier urlSupplier) throws Exception { @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), - MetadataStoreConfig.builder().build()); + MetadataStoreConfig.builder().fsyncEnable(false).build()); MetadataCache cache = store.getMetadataCache(String.class); @@ -187,7 +187,7 @@ public void updateValueWhenVersionIsOutOfSync(String provider, Supplier public void updateValueWhenKeyDisappears(String provider, Supplier urlSupplier) throws Exception { @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), - MetadataStoreConfig.builder().build()); + MetadataStoreConfig.builder().fsyncEnable(false).build()); MetadataCache cache = store.getMetadataCache(String.class); @@ -213,7 +213,7 @@ public void updateValueWhenKeyDisappears(String provider, Supplier urlSu public void revalidateLockWithinSameSession(String provider, Supplier urlSupplier) throws Exception { @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), - MetadataStoreConfig.builder().build()); + MetadataStoreConfig.builder().fsyncEnable(false).build()); @Cleanup CoordinationService cs2 = new CoordinationServiceImpl(store); diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataBenchmark.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataBenchmark.java index 9b7f8daca57e9..227c0e2c9dc35 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataBenchmark.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataBenchmark.java @@ -108,8 +108,7 @@ public void testGetChildren(String provider, Supplier urlSupplier) throw @Test(dataProvider = "impl", enabled = false) public void testPut(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), - MetadataStoreConfig.builder().fsyncEnable(false).build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); final int N_KEYS = 10_000; final int N_PUTS = 100_000; diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java index b3a670dc43d14..e7ccd35571261 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java @@ -245,7 +245,8 @@ public void insertionDeletionWitGenericType(String provider, Supplier ur @Test(dataProvider = "impl") public void insertionDeletion(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); MetadataCache objCache = store.getMetadataCache(MyClass.class); String key1 = newKey(); diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataStoreTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataStoreTest.java index ece23583c4d46..b8a383da58ddf 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataStoreTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataStoreTest.java @@ -60,7 +60,8 @@ public class MetadataStoreTest extends BaseMetadataStoreTest { @Test(dataProvider = "impl") public void emptyStoreTest(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); assertFalse(store.exists("/non-existing-key").join()); assertFalse(store.exists("/non-existing-key/child").join()); @@ -89,7 +90,8 @@ public void emptyStoreTest(String provider, Supplier urlSupplier) throws @Test(dataProvider = "impl") public void concurrentPutTest(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); String data = "data"; String path = "/non-existing-key"; @@ -109,7 +111,8 @@ public void concurrentPutTest(String provider, Supplier urlSupplier) thr @Test(dataProvider = "impl") public void insertionTestWithExpectedVersion(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); String key1 = newKey(); @@ -167,7 +170,8 @@ public void insertionTestWithExpectedVersion(String provider, Supplier u @Test(dataProvider = "impl") public void getChildrenTest(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); String key = newKey(); int n = 10; @@ -218,7 +222,8 @@ public void navigateChildrenTest(String provider, Supplier urlSupplier) @Test(dataProvider = "impl") public void deletionTest(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); String key = newKey(); int n = 10; @@ -252,7 +257,8 @@ public void deletionTest(String provider, Supplier urlSupplier) throws E @Test(dataProvider = "impl") public void emptyKeyTest(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); try { store.delete("", Optional.empty()).join(); @@ -293,7 +299,8 @@ public void emptyKeyTest(String provider, Supplier urlSupplier) throws E @Test(dataProvider = "impl") public void notificationListeners(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); BlockingQueue notifications = new LinkedBlockingDeque<>(); store.registerListener(n -> { @@ -359,7 +366,8 @@ public void notificationListeners(String provider, Supplier urlSupplier) @Test(dataProvider = "impl") public void testDeleteRecursive(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); String prefix = newKey(); @@ -384,7 +392,8 @@ public void testDeleteRecursive(String provider, Supplier urlSupplier) t @Test(dataProvider = "impl") public void testDeleteUnusedDirectories(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); String prefix = newKey(); @@ -413,7 +422,8 @@ public void testDeleteUnusedDirectories(String provider, Supplier urlSup @Test(dataProvider = "impl") public void testPersistent(String provider, Supplier urlSupplier) throws Exception { String metadataUrl = urlSupplier.get(); - MetadataStore store = MetadataStoreFactory.create(metadataUrl, MetadataStoreConfig.builder().build()); + MetadataStore store = + MetadataStoreFactory.create(metadataUrl, MetadataStoreConfig.builder().fsyncEnable(false).build()); byte[] data = "testPersistent".getBytes(StandardCharsets.UTF_8); String key = newKey() + "/a/b/c"; @@ -429,7 +439,8 @@ public void testPersistent(String provider, Supplier urlSupplier) throws @Test(dataProvider = "impl") public void testConcurrentPutGetOneKey(String provider, Supplier urlSupplier) throws Exception { - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); byte[] data = new byte[]{0}; String path = newKey(); int maxValue = 100; @@ -474,7 +485,8 @@ public void run() { @Test(dataProvider = "impl") public void testConcurrentPut(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); String k = newKey(); CompletableFuture f1 = @@ -489,7 +501,8 @@ public void testConcurrentPut(String provider, Supplier urlSupplier) thr @Test(dataProvider = "impl") public void testConcurrentDelete(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); String k = newKey(); store.put(k, new byte[0], Optional.of(-1L)).join(); @@ -505,7 +518,8 @@ public void testConcurrentDelete(String provider, Supplier urlSupplier) @Test(dataProvider = "impl") public void testGetChildren(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStore store = MetadataStoreFactory.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); store.put("/a/a-1", "value1".getBytes(StandardCharsets.UTF_8), Optional.empty()).join(); store.put("/a/a-2", "value1".getBytes(StandardCharsets.UTF_8), Optional.empty()).join(); diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/LedgerUnderreplicationManagerTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/LedgerUnderreplicationManagerTest.java index 0dd40c4779406..0df325b3c57a0 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/LedgerUnderreplicationManagerTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/LedgerUnderreplicationManagerTest.java @@ -93,7 +93,8 @@ private Future getLedgerToReplicate(LedgerUnderreplicationManager m) { private void methodSetup(Supplier urlSupplier) throws Exception { this.executor = Executors.newSingleThreadExecutor(); String ledgersRoot = "/ledgers-" + UUID.randomUUID(); - this.store = MetadataStoreExtended.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + this.store = MetadataStoreExtended.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); this.layoutManager = new PulsarLayoutManager(store, ledgersRoot); this.lmf = new PulsarLedgerManagerFactory(); diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/PulsarRegistrationClientTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/PulsarRegistrationClientTest.java index d2a21a2249df3..f599451c00710 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/PulsarRegistrationClientTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/PulsarRegistrationClientTest.java @@ -19,12 +19,11 @@ package org.apache.pulsar.metadata.bookkeeper; import static org.apache.bookkeeper.common.concurrent.FutureUtils.result; +import static org.mockito.Mockito.mock; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; -import static org.mockito.Mockito.mock; import static org.testng.Assert.assertTrue; import static org.testng.Assert.expectThrows; - import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -73,8 +72,8 @@ private static Set prepareNBookies(int num) { @Test(dataProvider = "impl") public void testGetWritableBookies(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStoreExtended store = - MetadataStoreExtended.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); String ledgersRoot = "/test/ledgers-" + UUID.randomUUID(); @@ -100,8 +99,8 @@ public void testGetWritableBookies(String provider, Supplier urlSupplier @Test(dataProvider = "impl") public void testGetReadonlyBookies(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStoreExtended store = - MetadataStoreExtended.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); String ledgersRoot = "/test/ledgers-" + UUID.randomUUID(); @@ -126,8 +125,8 @@ public void testGetReadonlyBookies(String provider, Supplier urlSupplier @Test(dataProvider = "impl") public void testGetBookieServiceInfo(String provider, Supplier urlSupplier) throws Exception { @Cleanup - MetadataStoreExtended store = - MetadataStoreExtended.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); String ledgersRoot = "/test/ledgers-" + UUID.randomUUID(); @@ -318,8 +317,8 @@ private void testWatchBookiesSuccess(String provider, Supplier urlSuppli throws Exception { @Cleanup - MetadataStoreExtended store = - MetadataStoreExtended.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); String ledgersRoot = "/test/ledgers-" + UUID.randomUUID(); diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/PulsarRegistrationManagerTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/PulsarRegistrationManagerTest.java index 1a677914f8fa1..a61a66bf2c947 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/PulsarRegistrationManagerTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/PulsarRegistrationManagerTest.java @@ -51,7 +51,8 @@ public class PulsarRegistrationManagerTest extends BaseMetadataStoreTest { private void methodSetup(Supplier urlSupplier) throws Exception { this.ledgersRootPath = "/ledgers-" + UUID.randomUUID(); - this.store = MetadataStoreExtended.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + this.store = MetadataStoreExtended.create(urlSupplier.get(), + MetadataStoreConfig.builder().fsyncEnable(false).build()); this.registrationManager = new PulsarRegistrationManager(store, ledgersRootPath, new ServerConfiguration()); } From f08bf96a5024a25b682e5a208a705a111873f22a Mon Sep 17 00:00:00 2001 From: coderzc Date: Fri, 9 Dec 2022 16:41:55 +0800 Subject: [PATCH 5/5] fix code --- .../src/main/java/org/apache/pulsar/broker/PulsarService.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index 37c740f55302a..0a49d1092d3dc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -362,7 +362,6 @@ public MetadataStore createConfigurationMetadataStore(PulsarMetadataEventSynchro .batchingMaxDelayMillis(config.getMetadataStoreBatchingMaxDelayMillis()) .batchingMaxOperations(config.getMetadataStoreBatchingMaxOperations()) .batchingMaxSizeKb(config.getMetadataStoreBatchingMaxSizeKb()) - .fsyncEnable(config.isMetadataFsyncEnabled()) .metadataStoreName(MetadataStoreConfig.CONFIGURATION_METADATA_STORE) .synchronizer(synchronizer) .build());