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
52 changes: 26 additions & 26 deletions distribution/server/src/assemble/LICENSE.bin.txt
Original file line number Diff line number Diff line change
Expand Up @@ -358,32 +358,32 @@ The Apache Software License, Version 2.0
- net.java.dev.jna-jna-jpms-5.19.1.jar
- net.java.dev.jna-jna-platform-jpms-5.19.1.jar
* BookKeeper
- org.apache.bookkeeper-bookkeeper-common-4.18.0.jar
- org.apache.bookkeeper-bookkeeper-common-allocator-4.18.0.jar
- org.apache.bookkeeper-bookkeeper-proto-4.18.0.jar
- org.apache.bookkeeper-bookkeeper-server-4.18.0.jar
- org.apache.bookkeeper-bookkeeper-tools-framework-4.18.0.jar
- org.apache.bookkeeper-circe-checksum-4.18.0.jar
- org.apache.bookkeeper-cpu-affinity-4.18.0.jar
- org.apache.bookkeeper-statelib-4.18.0.jar
- org.apache.bookkeeper-stream-storage-api-4.18.0.jar
- org.apache.bookkeeper-stream-storage-common-4.18.0.jar
- org.apache.bookkeeper-stream-storage-java-client-4.18.0.jar
- org.apache.bookkeeper-stream-storage-java-client-base-4.18.0.jar
- org.apache.bookkeeper-stream-storage-proto-4.18.0.jar
- org.apache.bookkeeper-stream-storage-server-4.18.0.jar
- org.apache.bookkeeper-stream-storage-service-api-4.18.0.jar
- org.apache.bookkeeper-stream-storage-service-impl-4.18.0.jar
- org.apache.bookkeeper.http-http-server-4.18.0.jar
- org.apache.bookkeeper.http-vertx-http-server-4.18.0.jar
- org.apache.bookkeeper.stats-bookkeeper-stats-api-4.18.0.jar
- org.apache.bookkeeper.stats-prometheus-metrics-provider-4.18.0.jar
- org.apache.distributedlog-distributedlog-common-4.18.0.jar
- org.apache.distributedlog-distributedlog-core-4.18.0-tests.jar
- org.apache.distributedlog-distributedlog-core-4.18.0.jar
- org.apache.distributedlog-distributedlog-protocol-4.18.0.jar
- org.apache.bookkeeper-native-io-4.18.0.jar
- org.apache.bookkeeper-native-library-common-4.18.0.jar
- org.apache.bookkeeper-bookkeeper-common-4.18.1.jar
- org.apache.bookkeeper-bookkeeper-common-allocator-4.18.1.jar
- org.apache.bookkeeper-bookkeeper-proto-4.18.1.jar
- org.apache.bookkeeper-bookkeeper-server-4.18.1.jar
- org.apache.bookkeeper-bookkeeper-tools-framework-4.18.1.jar
- org.apache.bookkeeper-circe-checksum-4.18.1.jar
- org.apache.bookkeeper-cpu-affinity-4.18.1.jar
- org.apache.bookkeeper-statelib-4.18.1.jar
- org.apache.bookkeeper-stream-storage-api-4.18.1.jar
- org.apache.bookkeeper-stream-storage-common-4.18.1.jar
- org.apache.bookkeeper-stream-storage-java-client-4.18.1.jar
- org.apache.bookkeeper-stream-storage-java-client-base-4.18.1.jar
- org.apache.bookkeeper-stream-storage-proto-4.18.1.jar
- org.apache.bookkeeper-stream-storage-server-4.18.1.jar
- org.apache.bookkeeper-stream-storage-service-api-4.18.1.jar
- org.apache.bookkeeper-stream-storage-service-impl-4.18.1.jar
- org.apache.bookkeeper.http-http-server-4.18.1.jar
- org.apache.bookkeeper.http-vertx-http-server-4.18.1.jar
- org.apache.bookkeeper.stats-bookkeeper-stats-api-4.18.1.jar
- org.apache.bookkeeper.stats-prometheus-metrics-provider-4.18.1.jar
- org.apache.distributedlog-distributedlog-common-4.18.1.jar
- org.apache.distributedlog-distributedlog-core-4.18.1-tests.jar
- org.apache.distributedlog-distributedlog-core-4.18.1.jar
- org.apache.distributedlog-distributedlog-protocol-4.18.1.jar
- org.apache.bookkeeper-native-io-4.18.1.jar
- org.apache.bookkeeper-native-library-common-4.18.1.jar
- at.yawk.lz4-lz4-java-1.11.2.jar
* Apache Felix HTTP Wrappers
- org.apache.felix-org.apache.felix.http.wrappers-1.1.10.jar
Expand Down
8 changes: 4 additions & 4 deletions distribution/shell/src/assemble/LICENSE.bin.txt
Original file line number Diff line number Diff line change
Expand Up @@ -398,10 +398,10 @@ The Apache Software License, Version 2.0
- slog-0.10.0.jar

* BookKeeper
- bookkeeper-common-allocator-4.18.0.jar
- cpu-affinity-4.18.0.jar
- circe-checksum-4.18.0.jar
- native-library-common-4.18.0.jar
- bookkeeper-common-allocator-4.18.1.jar
- cpu-affinity-4.18.1.jar
- circe-checksum-4.18.1.jar
- native-library-common-4.18.1.jar
* AirCompressor
- aircompressor-2.0.3.jar
* AsyncHttpClient
Expand Down
2 changes: 1 addition & 1 deletion gradle/libs.versions.toml
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ async-profiler = "4.5"
# Code quality
checkstyle = "13.3.0"
# Major frameworks
bookkeeper = "4.18.0"
bookkeeper = "4.18.1"
zookeeper = "3.9.5"
netty = "4.2.18.Final"
jetty = "12.1.12"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6794,7 +6794,7 @@ public void deleteFailed(ManagedLedgerException exception, Object ctx) {

@Test
void testForceCursorRecovery() throws Exception {
TestPulsarMockBookKeeper bk = new TestPulsarMockBookKeeper(executor);
TestPulsarMockBookKeeper bk = new TestPulsarMockBookKeeper(bkExecutor);
factory.shutdown();
factory = new ManagedLedgerFactoryImpl(metadataStore, bk);
ManagedLedgerConfig config = new ManagedLedgerConfig();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
import lombok.CustomLog;
import lombok.SneakyThrows;
import org.apache.bookkeeper.client.PulsarMockBookKeeper;
import org.apache.bookkeeper.common.util.OrderedExecutor;
import org.apache.bookkeeper.common.util.OrderedScheduler;
import org.apache.bookkeeper.mledger.ManagedLedgerConfig;
import org.apache.bookkeeper.mledger.ManagedLedgerException;
Expand Down Expand Up @@ -53,6 +54,7 @@ public abstract class MockedBookKeeperTestCase {
protected ManagedLedgerFactoryImpl factory;

protected OrderedScheduler executor;
protected OrderedExecutor bkExecutor;
protected ExecutorService cachedExecutor;

protected FaultInjectionMetadataStore metadataStore;
Expand Down Expand Up @@ -141,6 +143,8 @@ protected void cleanUpTestCase() throws Exception {
@BeforeClass(alwaysRun = true)
public final void setUpClass() {
executor = OrderedScheduler.newSchedulerBuilder().numThreads(2).name("test").build();
// The mock BookKeeper client needs an OrderedExecutor (not an OrderedScheduler) as its main worker pool.
bkExecutor = OrderedExecutor.newBuilder().numThreads(2).name("test-bk").build();
cachedExecutor = Executors.newCachedThreadPool();
}

Expand All @@ -149,6 +153,9 @@ public final void tearDownClass() {
if (executor != null) {
executor.shutdownNow();
}
if (bkExecutor != null) {
bkExecutor.shutdownNow();
}
if (cachedExecutor != null) {
cachedExecutor.shutdownNow();
}
Expand All @@ -166,7 +173,7 @@ protected void startBookKeeper() throws Exception {

metadataStore.put("/ledgers/LAYOUT", "1\nflat:1".getBytes(), Optional.empty()).join();

bkc = new PulsarMockBookKeeper(executor);
bkc = new PulsarMockBookKeeper(bkExecutor);
}

protected void stopBookKeeper() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import java.util.concurrent.Executors;
import lombok.CustomLog;
import org.apache.bookkeeper.client.PulsarMockBookKeeper;
import org.apache.bookkeeper.common.util.OrderedExecutor;
import org.apache.bookkeeper.common.util.OrderedScheduler;
import org.apache.bookkeeper.conf.ClientConfiguration;
import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig;
Expand Down Expand Up @@ -53,6 +54,7 @@ public abstract class MockedBookKeeperTestCase {
protected ClientConfiguration baseClientConf = new ClientConfiguration();

protected OrderedScheduler executor;
protected OrderedExecutor bkExecutor;
protected ExecutorService cachedExecutor;

public MockedBookKeeperTestCase() {
Expand Down Expand Up @@ -98,12 +100,15 @@ public void tearDown(Method method) {
@BeforeClass(alwaysRun = true)
public void setUpClass() {
executor = OrderedScheduler.newSchedulerBuilder().numThreads(2).name("test").build();
// The mock BookKeeper client needs an OrderedExecutor (not an OrderedScheduler) as its main worker pool.
bkExecutor = OrderedExecutor.newBuilder().numThreads(2).name("test-bk").build();
cachedExecutor = Executors.newCachedThreadPool();
}

@AfterClass(alwaysRun = true)
public void tearDownClass() {
executor.shutdownNow();
bkExecutor.shutdownNow();
cachedExecutor.shutdownNow();
}

Expand All @@ -120,7 +125,7 @@ protected void startBookKeeper() throws Exception {

metadataStore.put("/ledgers/LAYOUT", "1\nflat:1".getBytes(), Optional.empty()).join();

bkc = new PulsarMockBookKeeper(executor);
bkc = new PulsarMockBookKeeper(bkExecutor);
}

protected void stopBookKeeper() throws Exception {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import java.util.concurrent.Executors;
import lombok.CustomLog;
import org.apache.bookkeeper.client.PulsarMockBookKeeper;
import org.apache.bookkeeper.common.util.OrderedExecutor;
import org.apache.bookkeeper.common.util.OrderedScheduler;
import org.apache.bookkeeper.conf.ClientConfiguration;
import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig;
Expand Down Expand Up @@ -53,6 +54,7 @@ public abstract class MockedBookKeeperTestCase {
protected ClientConfiguration baseClientConf = new ClientConfiguration();

protected OrderedScheduler executor;
protected OrderedExecutor bkExecutor;
protected ExecutorService cachedExecutor;

public MockedBookKeeperTestCase() {
Expand Down Expand Up @@ -105,12 +107,15 @@ public void tearDown(Method method) {
@BeforeClass(alwaysRun = true)
public void setUpClass() {
executor = OrderedScheduler.newSchedulerBuilder().numThreads(2).name("test").build();
// The mock BookKeeper client needs an OrderedExecutor (not an OrderedScheduler) as its main worker pool.
bkExecutor = OrderedExecutor.newBuilder().numThreads(2).name("test-bk").build();
cachedExecutor = Executors.newCachedThreadPool();
}

@AfterClass(alwaysRun = true)
public void tearDownClass() {
executor.shutdownNow();
bkExecutor.shutdownNow();
cachedExecutor.shutdownNow();
}

Expand All @@ -126,7 +131,7 @@ protected void startBookKeeper() throws Exception {

metadataStore.put("/ledgers/LAYOUT", "1\nflat:1".getBytes(), Optional.empty());

bkc = new PulsarMockBookKeeper(executor);
bkc = new PulsarMockBookKeeper(bkExecutor);
}

protected void stopBookKeeper() throws Exception {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import java.util.Properties;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import org.apache.bookkeeper.common.util.OrderedExecutor;
import org.apache.bookkeeper.common.util.OrderedScheduler;
import org.apache.bookkeeper.mledger.LedgerOffloaderStats;
import org.apache.bookkeeper.mledger.offload.filesystem.impl.FileSystemManagedLedgerOffloader;
Expand All @@ -38,6 +39,7 @@
public abstract class FileStoreTestBase {
protected FileSystemManagedLedgerOffloader fileSystemManagedLedgerOffloader;
protected OrderedScheduler scheduler;
protected OrderedExecutor bkExecutor;
protected final String basePath = "pulsar";
private MiniDFSCluster hdfsCluster;
private String hdfsURI;
Expand All @@ -51,6 +53,8 @@ public final void beforeClass() throws Exception {

public void init() throws Exception {
scheduler = OrderedScheduler.newSchedulerBuilder().numThreads(1).name("offloader").build();
// The mock BookKeeper client needs an OrderedExecutor (not an OrderedScheduler) as its main worker pool.
bkExecutor = OrderedExecutor.newBuilder().numThreads(1).name("offloader-bk").build();
}

@AfterClass(alwaysRun = true)
Expand All @@ -63,6 +67,10 @@ public void cleanup() throws IOException {
scheduler.shutdownNow();
scheduler = null;
}
if (bkExecutor != null) {
bkExecutor.shutdownNow();
bkExecutor = null;
}
}

@BeforeMethod(alwaysRun = true)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@ public class FileSystemManagedLedgerOffloaderTest extends FileStoreTestBase {
@Override
public void init() throws Exception {
super.init();
this.bk = new PulsarMockBookKeeper(scheduler);
this.bk = new PulsarMockBookKeeper(bkExecutor);
this.toWrite = buildReadHandle();
map.put("ManagedLedgerName", managedLedgerName);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
import org.apache.bookkeeper.client.api.LedgerEntries;
import org.apache.bookkeeper.client.api.LedgerEntry;
import org.apache.bookkeeper.client.api.ReadHandle;
import org.apache.bookkeeper.common.util.OrderedExecutor;
import org.apache.bookkeeper.common.util.OrderedScheduler;
import org.apache.bookkeeper.mledger.LedgerOffloaderStats;
import org.apache.pulsar.common.naming.TopicName;
Expand All @@ -44,11 +45,14 @@

public class FileSystemOffloaderLocalFileTest {
private OrderedScheduler scheduler;
private OrderedExecutor bkExecutor;
private LedgerOffloaderStats offloaderStats;

@BeforeClass
public void setup() throws Exception {
scheduler = OrderedScheduler.newSchedulerBuilder().numThreads(1).name("offloader").build();
// The mock BookKeeper client needs an OrderedExecutor (not an OrderedScheduler) as its main worker pool.
bkExecutor = OrderedExecutor.newBuilder().numThreads(1).name("offloader-bk").build();
offloaderStats = LedgerOffloaderStats.create(true, true, scheduler, 60);
}

Expand All @@ -57,6 +61,9 @@ public void cleanup() throws Exception {
if (scheduler != null) {
scheduler.shutdown();
}
if (bkExecutor != null) {
bkExecutor.shutdownNow();
}
if (offloaderStats != null) {
offloaderStats.close();
}
Expand All @@ -83,7 +90,7 @@ public void testReadWriteWithLocalFileUsingFileSystemURI() throws Exception {

// prepare the data in bookkeeper
@Cleanup
BookKeeper bk = new PulsarMockBookKeeper(scheduler);
BookKeeper bk = new PulsarMockBookKeeper(bkExecutor);
LedgerHandle lh = bk.createLedger(1, 1, 1, BookKeeper.DigestType.CRC32, "".getBytes());
for (int i = 0; i < numberOfEntries; i++) {
byte[] entry = ("foobar" + i).getBytes();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
import org.apache.bookkeeper.client.PulsarMockBookKeeper;
import org.apache.bookkeeper.client.api.DigestType;
import org.apache.bookkeeper.client.api.ReadHandle;
import org.apache.bookkeeper.common.util.OrderedExecutor;
import org.apache.bookkeeper.common.util.OrderedScheduler;
import org.apache.bookkeeper.mledger.offload.jcloud.provider.JCloudBlobStoreProvider;
import org.apache.bookkeeper.mledger.offload.jcloud.provider.TieredStorageConfiguration;
Expand All @@ -43,6 +44,7 @@ public abstract class BlobStoreManagedLedgerOffloaderBase {
protected static final int DEFAULT_READ_BUFFER_SIZE = 1 * 1024 * 1024;

protected final OrderedScheduler scheduler;
protected final OrderedExecutor bkExecutor;
protected final PulsarMockBookKeeper bk;
protected final JCloudBlobStoreProvider provider;
protected TieredStorageConfiguration config;
Expand All @@ -51,7 +53,9 @@ public abstract class BlobStoreManagedLedgerOffloaderBase {

protected BlobStoreManagedLedgerOffloaderBase() throws Exception {
scheduler = OrderedScheduler.newSchedulerBuilder().numThreads(5).name("offloader").build();
bk = new PulsarMockBookKeeper(scheduler);
// The mock BookKeeper client needs an OrderedExecutor (not an OrderedScheduler) as its main worker pool.
bkExecutor = OrderedExecutor.newBuilder().numThreads(1).name("offloader-bk").build();
bk = new PulsarMockBookKeeper(bkExecutor);
provider = getBlobStoreProvider();
}

Expand All @@ -65,6 +69,7 @@ public void cleanupMockBookKeeper() {
public void cleanup() throws Exception {
entryOffsetsCache.close();
scheduler.shutdownNow();
bkExecutor.shutdownNow();
}

protected static MockManagedLedger createMockManagedLedger() {
Expand Down
Loading