Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -57,9 +57,6 @@
private int ensembleSize = 3;
private int writeQuorumSize = 2;
private int ackQuorumSize = 2;
private int metadataEnsembleSize = 3;
private int metadataWriteQuorumSize = 2;
private int metadataAckQuorumSize = 2;
private int metadataMaxEntriesPerLedger = 50000;
private int ledgerRolloverTimeout = 4 * 3600;
private double throttleMarkDelete = 0;
Expand Down Expand Up @@ -301,54 +298,6 @@
return this;
}

/**
* @return the metadataEnsemblesize
*/
public int getMetadataEnsemblesize() {
return metadataEnsembleSize;
}

/**
* @param metadataEnsembleSize
* the metadataEnsembleSize to set
*/
public ManagedLedgerConfig setMetadataEnsembleSize(int metadataEnsembleSize) {
this.metadataEnsembleSize = metadataEnsembleSize;
return this;
}

/**
* @return the metadataAckQuorumSize
*/
public int getMetadataAckQuorumSize() {
return metadataAckQuorumSize;
}

/**
* @return the metadataWriteQuorumSize
*/
public int getMetadataWriteQuorumSize() {
return metadataWriteQuorumSize;
}

/**
* @param metadataAckQuorumSize
* the metadataAckQuorumSize to set
*/
public ManagedLedgerConfig setMetadataAckQuorumSize(int metadataAckQuorumSize) {
this.metadataAckQuorumSize = metadataAckQuorumSize;
return this;
}

/**
* @param metadataWriteQuorumSize
* the metadataWriteQuorumSize to set
*/
public ManagedLedgerConfig setMetadataWriteQuorumSize(int metadataWriteQuorumSize) {
this.metadataWriteQuorumSize = metadataWriteQuorumSize;
return this;
}

/**
* @return the metadataMaxEntriesPerLedger
*/
Expand Down Expand Up @@ -655,7 +604,7 @@

/**
* Ledger read entry timeout after which callback will be completed with failure. (disable timeout by setting.
* readTimeoutSeconds <= 0)

Check warning on line 607 in managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerConfig.java

View workflow job for this annotation

GitHub Actions / Build and License check

invalid input: '<'
*
* @param readEntryTimeoutSeconds
* @return
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -84,8 +84,7 @@ public void testSimpleRead() throws Exception {
@Cleanup("shutdown")
ManagedLedgerFactory factory = new ManagedLedgerFactoryImpl(metadataStore, bkc, factoryConf);
ManagedLedgerConfig config = new ManagedLedgerConfig();
config.setEnsembleSize(1).setWriteQuorumSize(1).setAckQuorumSize(1).setMetadataEnsembleSize(1)
.setMetadataAckQuorumSize(1);
config.setEnsembleSize(1).setWriteQuorumSize(1).setAckQuorumSize(1);
ManagedLedger ledger = factory.open("my-ledger" + testName, config);
ManagedCursor cursor = ledger.openCursor("c1");

Expand All @@ -105,7 +104,7 @@ public void testSimpleRead() throws Exception {
public void testBookieFailure() throws Exception {
ManagedLedgerFactory factory = new ManagedLedgerFactoryImpl(metadataStore, bkc);
ManagedLedgerConfig config = new ManagedLedgerConfig();
config.setEnsembleSize(2).setAckQuorumSize(2).setMetadataEnsembleSize(2);
config.setEnsembleSize(2).setAckQuorumSize(2);
ManagedLedger ledger = factory.open("my-ledger" + testName, config);
ManagedCursor cursor = ledger.openCursor("my-cursor");
ledger.addEntry("entry-0".getBytes());
Expand Down Expand Up @@ -176,7 +175,7 @@ public void verifyConcurrentUsage() throws Exception {

EntryCacheManager cacheManager = factory.getEntryCacheManager();
ManagedLedgerConfig conf = new ManagedLedgerConfig();
conf.setEnsembleSize(2).setAckQuorumSize(2).setMetadataEnsembleSize(2);
conf.setEnsembleSize(2).setAckQuorumSize(2);
final ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("my-ledger" + testName, conf);

int numProducers = 1;
Expand Down Expand Up @@ -254,7 +253,7 @@ public void verifyAsyncReadEntryUsingCache() throws Exception {
ManagedLedgerFactoryImpl factory = new ManagedLedgerFactoryImpl(metadataStore, bkc, config);

ManagedLedgerConfig conf = new ManagedLedgerConfig();
conf.setEnsembleSize(2).setAckQuorumSize(2).setMetadataEnsembleSize(2)
conf.setEnsembleSize(2).setAckQuorumSize(2)
.setRetentionSizeInMB(-1).setRetentionTime(-1, TimeUnit.MILLISECONDS);
final ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("my-ledger" + testName, conf);

Expand Down Expand Up @@ -349,7 +348,7 @@ public void testSimple() throws Exception {
@Cleanup("shutdown")
ManagedLedgerFactory factory = new ManagedLedgerFactoryImpl(metadataStore, bkc);
ManagedLedgerConfig mlConfig = new ManagedLedgerConfig();
mlConfig.setEnsembleSize(1).setAckQuorumSize(1).setMetadataEnsembleSize(1).setWriteQuorumSize(1);
mlConfig.setEnsembleSize(1).setAckQuorumSize(1).setWriteQuorumSize(1);
// set the data ledger size
mlConfig.setMaxEntriesPerLedger(100);
// set the metadata ledger size to 1 to kick off many ledger switching cases
Expand All @@ -365,8 +364,7 @@ public void testConcurrentMarkDelete() throws Exception {
ManagedLedgerFactory factory = new ManagedLedgerFactoryImpl(metadataStore, bkc);
ManagedLedgerConfig mlConfig = new ManagedLedgerConfig();
mlConfig.setEnsembleSize(1).setWriteQuorumSize(1)
.setAckQuorumSize(1).setMetadataEnsembleSize(1).setMetadataWriteQuorumSize(1)
.setMetadataAckQuorumSize(1);
.setAckQuorumSize(1);
// set the data ledger size
mlConfig.setMaxEntriesPerLedger(100);
// set the metadata ledger size to 1 to kick off many ledger switching cases
Expand Down Expand Up @@ -417,8 +415,7 @@ public void asyncMarkDeleteAndClose() throws Exception {
ManagedLedgerFactory factory = new ManagedLedgerFactoryImpl(metadataStore, bkc);

ManagedLedgerConfig config = new ManagedLedgerConfig().setEnsembleSize(1).setWriteQuorumSize(1)
.setAckQuorumSize(1).setMetadataEnsembleSize(1).setMetadataWriteQuorumSize(1)
.setMetadataAckQuorumSize(1);
.setAckQuorumSize(1);
ManagedLedger ledger = factory.open("my_test_ledger" + testName, config);
ManagedCursor cursor = ledger.openCursor("c1");

Expand Down Expand Up @@ -469,7 +466,7 @@ public void ledgerFencedByAutoReplication() throws Exception {
@Cleanup("shutdown")
ManagedLedgerFactoryImpl factory = new ManagedLedgerFactoryImpl(metadataStore, bkc);
ManagedLedgerConfig config = new ManagedLedgerConfig();
config.setEnsembleSize(2).setAckQuorumSize(2).setMetadataEnsembleSize(2);
config.setEnsembleSize(2).setAckQuorumSize(2);
ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("my_test_ledger" + testName, config);
ManagedCursor c1 = ledger.openCursor("c1");

Expand Down Expand Up @@ -499,7 +496,7 @@ public void ledgerFencedByFailover() throws Exception {
@Cleanup("shutdown")
ManagedLedgerFactoryImpl factory1 = new ManagedLedgerFactoryImpl(metadataStore, bkc);
ManagedLedgerConfig config = new ManagedLedgerConfig();
config.setEnsembleSize(2).setAckQuorumSize(2).setMetadataEnsembleSize(2);
config.setEnsembleSize(2).setAckQuorumSize(2);
ManagedLedgerImpl ledger1 = (ManagedLedgerImpl) factory1.open("my_test_ledger" + testName, config);
ledger1.openCursor("c");

Expand Down Expand Up @@ -538,8 +535,7 @@ public void testOfflineTopicBacklog() throws Exception {
@Cleanup("shutdown")
ManagedLedgerFactory factory = new ManagedLedgerFactoryImpl(metadataStore, bkc, factoryConf);
ManagedLedgerConfig config = new ManagedLedgerConfig();
config.setEnsembleSize(1).setWriteQuorumSize(1).setAckQuorumSize(1).setMetadataEnsembleSize(1)
.setMetadataAckQuorumSize(1);
config.setEnsembleSize(1).setWriteQuorumSize(1).setAckQuorumSize(1);
ManagedLedger ledger = factory.open("property/namespace/my-ledger", config);
ManagedCursor cursor = ledger.openCursor("c1");

Expand Down Expand Up @@ -567,8 +563,7 @@ void testResetCursorAfterRecovery() throws Exception {
@Cleanup("shutdown")
ManagedLedgerFactory factory = new ManagedLedgerFactoryImpl(metadataStore, bkc);
ManagedLedgerConfig conf = new ManagedLedgerConfig().setMaxEntriesPerLedger(10).setEnsembleSize(1)
.setWriteQuorumSize(1).setAckQuorumSize(1).setMetadataEnsembleSize(1).setMetadataWriteQuorumSize(1)
.setMetadataAckQuorumSize(1);
.setWriteQuorumSize(1).setAckQuorumSize(1);
ManagedLedger ledger = factory.open("my_test_move_cursor_ledger", conf);
ManagedCursor cursor = ledger.openCursor("trc1");
Position p1 = ledger.addEntry("dummy-entry-1".getBytes());
Expand Down Expand Up @@ -598,7 +593,7 @@ public void managedLedgerClosed() throws Exception {
@Cleanup("shutdown")
ManagedLedgerFactoryImpl factory = new ManagedLedgerFactoryImpl(metadataStore, bkc);
ManagedLedgerConfig config = new ManagedLedgerConfig();
config.setEnsembleSize(2).setAckQuorumSize(2).setMetadataEnsembleSize(2);
config.setEnsembleSize(2).setAckQuorumSize(2);
ManagedLedgerImpl ledger1 = (ManagedLedgerImpl) factory.open("my_test_ledger" + testName, config);

int num = 100;
Expand Down Expand Up @@ -637,7 +632,7 @@ public void testChangeCrcType() throws Exception {
@Cleanup("shutdown")
ManagedLedgerFactoryImpl factory = new ManagedLedgerFactoryImpl(metadataStore, bkc);
ManagedLedgerConfig config = new ManagedLedgerConfig();
config.setEnsembleSize(2).setAckQuorumSize(2).setMetadataEnsembleSize(2);
config.setEnsembleSize(2).setAckQuorumSize(2);
config.setDigestType(DigestType.CRC32);
ManagedLedger ledger = factory.open("my_test_ledger" + testName, config);
ManagedCursor c1 = ledger.openCursor("c1");
Expand Down Expand Up @@ -674,8 +669,7 @@ public void testPeriodicRollover() throws Exception {
@Cleanup("shutdown")
ManagedLedgerFactory factory = new ManagedLedgerFactoryImpl(metadataStore, bkc, factoryConf);
ManagedLedgerConfig config = new ManagedLedgerConfig();
config.setEnsembleSize(1).setWriteQuorumSize(1).setAckQuorumSize(1).setMetadataEnsembleSize(1)
.setMetadataAckQuorumSize(1)
config.setEnsembleSize(1).setWriteQuorumSize(1).setAckQuorumSize(1)
.setLedgerRolloverTimeout(rolloverTimeForCursorInSeconds);
ManagedLedger ledger = factory.open("my-ledger" + testName, config);
ManagedCursor cursor = ledger.openCursor("c1");
Expand Down Expand Up @@ -718,7 +712,6 @@ public void testConfigPersistIndividualAckAsLongArray(boolean enable) throws Exc
ManagedLedgerFactory factory = new ManagedLedgerFactoryImpl(metadataStore, bkc, factoryConf);
final ManagedLedgerConfig config = new ManagedLedgerConfig()
.setEnsembleSize(1).setWriteQuorumSize(1).setAckQuorumSize(1)
.setMetadataEnsembleSize(1).setMetadataWriteQuorumSize(1).setMetadataAckQuorumSize(1)
.setMaxUnackedRangesToPersistInMetadataStore(1)
.setPersistIndividualAckAsLongArray(enable);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,9 +59,7 @@ public void testChangeZKPath() throws Exception {
ManagedLedgerConfig config = new ManagedLedgerConfig();
config.setEnsembleSize(1)
.setWriteQuorumSize(1)
.setAckQuorumSize(1)
.setMetadataAckQuorumSize(1)
.setMetadataAckQuorumSize(1);
.setAckQuorumSize(1);
ManagedLedger ledger = factory.open("test-ledger" + testName, config);
ManagedCursor cursor = ledger.openCursor("test-c1" + testName);

Expand Down Expand Up @@ -96,9 +94,7 @@ public void testChangeZKPath2() throws Exception {
ManagedLedgerConfig config = new ManagedLedgerConfig();
config.setEnsembleSize(1)
.setWriteQuorumSize(1)
.setAckQuorumSize(1)
.setMetadataAckQuorumSize(1)
.setMetadataAckQuorumSize(1);
.setAckQuorumSize(1);
ManagedLedger ledger = factory.open("test-ledger" + testName, config);
ManagedCursor cursor = ledger.openCursor("test-c1" + testName);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,7 @@ public ManagedLedgerSingleBookieTest() {
@Test // (timeOut = 20000)
public void simple() throws Exception {
ManagedLedgerConfig config = new ManagedLedgerConfig().setEnsembleSize(1).setWriteQuorumSize(1)
.setAckQuorumSize(1).setMetadataEnsembleSize(1).setMetadataWriteQuorumSize(1)
.setMetadataAckQuorumSize(1);
.setAckQuorumSize(1);
ManagedLedger ledger = factory.open("my_test_ledger", config);

assertEquals(ledger.getNumberOfEntries(), 0);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -732,7 +732,6 @@
config.setMaxSizePerLedgerMb(1);
config.setEnsembleSize(1);
config.setWriteQuorumSize(1).setAckQuorumSize(1);
config.setMetadataWriteQuorumSize(1).setMetadataAckQuorumSize(1);
ManagedLedger ledger = factory.open("my_test_ledger", config);

assertEquals(ledger.getNumberOfEntries(), 0);
Expand Down Expand Up @@ -1867,7 +1866,7 @@
public void testOpenRaceCondition() throws Exception {
ManagedLedgerConfig config = new ManagedLedgerConfig();
initManagedLedgerConfig(config);
config.setEnsembleSize(2).setAckQuorumSize(2).setMetadataEnsembleSize(2);
config.setEnsembleSize(2).setAckQuorumSize(2);
final ManagedLedger ledger = factory.open("my-ledger", config);
final ManagedCursor c1 = ledger.openCursor("c1");

Expand Down Expand Up @@ -4864,7 +4863,7 @@
doAnswer(inv -> {
if (mlPath.equals(inv.getArgument(0)) && interceptNextPut.compareAndSet(true, false)) {
putIntercepted.countDown();
CompletableFuture<Stat> real = (CompletableFuture<Stat>) inv.callRealMethod();

Check warning on line 4866 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / Flaky tests suite

[unchecked] unchecked cast

Check warning on line 4866 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Proxy

[unchecked] unchecked cast

Check warning on line 4866 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Other

[unchecked] unchecked cast

Check warning on line 4866 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Broker Group 4

[unchecked] unchecked cast

Check warning on line 4866 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Pulsar Metadata

[unchecked] unchecked cast

Check warning on line 4866 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Client Impl

[unchecked] unchecked cast

Check warning on line 4866 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Client Api

[unchecked] unchecked cast

Check warning on line 4866 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Broker Group 3

[unchecked] unchecked cast

Check warning on line 4866 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Broker Group 2

[unchecked] unchecked cast

Check warning on line 4866 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Broker Group 5

[unchecked] unchecked cast

Check warning on line 4866 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Broker Group 1

[unchecked] unchecked cast
CompletableFuture<Stat> gated = new CompletableFuture<>();
// Forward the real result to `gated` only once `releaseGate` is completed.
real.whenComplete((stat, ex) -> releaseGate.whenComplete((ignored, ignoredEx) -> {
Expand Down Expand Up @@ -4971,7 +4970,7 @@
doAnswer(inv -> {
if (mlPath.equals(inv.getArgument(0)) && interceptNextPut.compareAndSet(true, false)) {
putIntercepted.countDown();
CompletableFuture<Stat> real = (CompletableFuture<Stat>) inv.callRealMethod();

Check warning on line 4973 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Proxy

[unchecked] unchecked cast

Check warning on line 4973 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Other

[unchecked] unchecked cast

Check warning on line 4973 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Broker Group 4

[unchecked] unchecked cast

Check warning on line 4973 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Pulsar Metadata

[unchecked] unchecked cast

Check warning on line 4973 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Client Impl

[unchecked] unchecked cast

Check warning on line 4973 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Client Api

[unchecked] unchecked cast

Check warning on line 4973 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Broker Group 3

[unchecked] unchecked cast

Check warning on line 4973 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Broker Group 2

[unchecked] unchecked cast

Check warning on line 4973 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Broker Group 5

[unchecked] unchecked cast

Check warning on line 4973 in managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Brokers - Broker Group 1

[unchecked] unchecked cast
CompletableFuture<Stat> gated = new CompletableFuture<>();
real.whenComplete((stat, ex) -> releaseGate.whenComplete((ignored, ignoredEx) -> {
if (ex != null) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -88,7 +88,6 @@ public void testUpgrade_CrossLedgerMessageRange() throws Exception {

ManagedLedgerConfig config = new ManagedLedgerConfig()
.setEnsembleSize(1).setWriteQuorumSize(1).setAckQuorumSize(1)
.setMetadataEnsembleSize(1).setMetadataWriteQuorumSize(1).setMetadataAckQuorumSize(1)
.setMaxEntriesPerLedger(5);

ManagedLedger ledger = factory.open(mlName, config);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2303,9 +2303,6 @@ public CompletableFuture<ManagedLedgerConfig> getManagedLedgerConfig(@NonNull To
.setReadEntryTimeoutSeconds(serviceConfig.getManagedLedgerReadEntryTimeoutSeconds());
managedLedgerConfig
.setAddEntryTimeoutSeconds(serviceConfig.getManagedLedgerAddEntryTimeoutSeconds());
managedLedgerConfig.setMetadataEnsembleSize(serviceConfig.getManagedLedgerDefaultEnsembleSize());
managedLedgerConfig.setMetadataWriteQuorumSize(serviceConfig.getManagedLedgerDefaultWriteQuorum());
managedLedgerConfig.setMetadataAckQuorumSize(serviceConfig.getManagedLedgerDefaultAckQuorum());
managedLedgerConfig
.setMetadataMaxEntriesPerLedger(serviceConfig.getManagedLedgerCursorMaxEntriesPerLedger());

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -188,9 +188,6 @@ public void run() throws Exception {
mlConf.setWriteQuorumSize(this.writeQuorum);
mlConf.setAckQuorumSize(this.ackQuorum);
mlConf.setMinimumRolloverTime(10, TimeUnit.MINUTES);
mlConf.setMetadataEnsembleSize(this.ensembleSize);
mlConf.setMetadataWriteQuorumSize(this.writeQuorum);
mlConf.setMetadataAckQuorumSize(this.ackQuorum);
mlConf.setDigestType(this.digestType);
mlConf.setMaxSizePerLedgerMb(2048);

Expand Down
Loading