From def4a94ad3e03cf0d032634fb91004bb8c5a5f38 Mon Sep 17 00:00:00 2001 From: Xiaopeng Zhang Date: Sat, 11 Apr 2020 14:20:02 +0800 Subject: [PATCH 1/7] fix filesystem offload oom based on https://github.com/apache/pulsar/pull/6697 --- .../FileSystemManagedLedgerOffloader.java | 91 +++++++++++++------ 1 file changed, 65 insertions(+), 26 deletions(-) diff --git a/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java b/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java index bbee8289b58ae..e0ca37ca77485 100644 --- a/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java +++ b/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java @@ -20,6 +20,7 @@ import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.ImmutableMap; +import io.netty.util.Recycler; import org.apache.bookkeeper.client.api.LedgerEntries; import org.apache.bookkeeper.client.api.LedgerEntry; import org.apache.bookkeeper.client.api.ReadHandle; @@ -42,6 +43,7 @@ import java.util.Map; import java.util.UUID; +import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executors; @@ -62,6 +64,7 @@ public class FileSystemManagedLedgerOffloader implements LedgerOffloader { private final FileSystem fileSystem; private OrderedScheduler scheduler; private static final long ENTRIES_PER_READ = 100; + private static final int PREFETCH_ROUNDS = 100; private OrderedScheduler assignmentScheduler; private OffloadPolicies offloadPolicies; @@ -188,12 +191,15 @@ public void run() { AtomicLong haveOffloadEntryNumber = new AtomicLong(0); long needToOffloadFirstEntryNumber = 0; CountDownLatch countDownLatch; + //avoid prefetch too much data into memory + ArrayBlockingQueue tasks = new ArrayBlockingQueue<>(PREFETCH_ROUNDS); do { long end = Math.min(needToOffloadFirstEntryNumber + ENTRIES_PER_READ - 1, readHandle.getLastAddConfirmed()); log.debug("read ledger entries. start: {}, end: {}", needToOffloadFirstEntryNumber, end); LedgerEntries ledgerEntriesOnce = readHandle.readAsync(needToOffloadFirstEntryNumber, end).get(); + tasks.put(true); countDownLatch = new CountDownLatch(1); - assignmentScheduler.chooseThread(ledgerId).submit(new FileSystemWriter(ledgerEntriesOnce, dataWriter, + assignmentScheduler.chooseThread(ledgerId).submit(FileSystemWriter.create(ledgerEntriesOnce, dataWriter, tasks, countDownLatch, haveOffloadEntryNumber, this)).addListener(() -> {}, Executors.newSingleThreadExecutor()); needToOffloadFirstEntryNumber = end + 1; } while (needToOffloadFirstEntryNumber - 1 != readHandle.getLastAddConfirmed() && fileSystemWriteException == null); @@ -216,45 +222,78 @@ public void run() { private static class FileSystemWriter implements Runnable { - private final LedgerEntries ledgerEntriesOnce; + private LedgerEntries ledgerEntriesOnce; private final LongWritable key = new LongWritable(); private final BytesWritable value = new BytesWritable(); - private final MapFile.Writer dataWriter; - private final CountDownLatch countDownLatch; - private final AtomicLong haveOffloadEntryNumber; - private final LedgerReader ledgerReader; + private MapFile.Writer dataWriter; + private CountDownLatch countDownLatch; + private AtomicLong haveOffloadEntryNumber; + private LedgerReader ledgerReader; + private ArrayBlockingQueue tasks; + private Recycler.Handle recyclerHandle; + private FileSystemWriter(Recycler.Handle recyclerHandle) { + this.recyclerHandle = recyclerHandle; + } + + private static final Recycler RECYCLER = new Recycler() { + @Override + protected FileSystemWriter newObject(Recycler.Handle handle) { + return new FileSystemWriter(handle); + } + }; - private FileSystemWriter(LedgerEntries ledgerEntriesOnce, MapFile.Writer dataWriter, + private void recycle() { + this.dataWriter = null; + this.countDownLatch = null; + this.haveOffloadEntryNumber = null; + this.ledgerReader = null; + this.ledgerEntriesOnce = null; + this.tasks = null; + recyclerHandle.recycle(this); + } + + + public static FileSystemWriter create(LedgerEntries ledgerEntriesOnce, MapFile.Writer dataWriter, ArrayBlockingQueue tasks, CountDownLatch countDownLatch, AtomicLong haveOffloadEntryNumber, LedgerReader ledgerReader) { - this.ledgerEntriesOnce = ledgerEntriesOnce; - this.dataWriter = dataWriter; - this.countDownLatch = countDownLatch; - this.haveOffloadEntryNumber = haveOffloadEntryNumber; - this.ledgerReader = ledgerReader; + FileSystemWriter writer = RECYCLER.get(); + writer.ledgerReader = ledgerReader; + writer.dataWriter = dataWriter; + writer.countDownLatch = countDownLatch; + writer.haveOffloadEntryNumber = haveOffloadEntryNumber; + writer.ledgerEntriesOnce = ledgerEntriesOnce; + writer.tasks = tasks; + return writer; } @Override public void run() { - if (ledgerReader.fileSystemWriteException == null) { - Iterator iterator = ledgerEntriesOnce.iterator(); - while (iterator.hasNext()) { - LedgerEntry entry = iterator.next(); - long entryId = entry.getEntryId(); - key.set(entryId); - try { - value.set(entry.getEntryBytes(), 0, entry.getEntryBytes().length); - dataWriter.append(key, value); - } catch (IOException e) { - ledgerReader.fileSystemWriteException = e; - break; + try { + if (ledgerReader.fileSystemWriteException == null) { + Iterator iterator = ledgerEntriesOnce.iterator(); + while (iterator.hasNext()) { + LedgerEntry entry = iterator.next(); + long entryId = entry.getEntryId(); + key.set(entryId); + try { + value.set(entry.getEntryBytes(), 0, entry.getEntryBytes().length); + dataWriter.append(key, value); + } catch (IOException e) { + ledgerReader.fileSystemWriteException = e; + break; + } + haveOffloadEntryNumber.incrementAndGet(); } - haveOffloadEntryNumber.incrementAndGet(); } + countDownLatch.countDown(); + ledgerEntriesOnce.close(); + tasks.take(); + this.recycle(); + } catch (InterruptedException e) { + ledgerReader.fileSystemWriteException = e; } - countDownLatch.countDown(); } } From c5295836f8c8914a60a05fa059546021f469d94f Mon Sep 17 00:00:00 2001 From: Xiaopeng Zhang Date: Sun, 12 Apr 2020 15:42:08 +0800 Subject: [PATCH 2/7] use semaphore instead of blockqueue --- .../FileSystemManagedLedgerOffloader.java | 72 +++++++++---------- 1 file changed, 36 insertions(+), 36 deletions(-) diff --git a/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java b/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java index e0ca37ca77485..7ecc215186d76 100644 --- a/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java +++ b/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java @@ -6,9 +6,9 @@ * to you under the Apache License, Version 2.0 (the * "License"); you may not use this file except in compliance * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * + *

+ * 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 @@ -47,6 +47,7 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executors; +import java.util.concurrent.Semaphore; import java.util.concurrent.atomic.AtomicLong; import static org.apache.bookkeeper.mledger.offload.OffloadUtils.buildLedgerMetadataFormat; @@ -71,6 +72,7 @@ public class FileSystemManagedLedgerOffloader implements LedgerOffloader { public static boolean driverSupported(String driver) { return DRIVER_NAMES.equals(driver); } + @Override public String getOffloadDriverName() { return driverName; @@ -85,7 +87,7 @@ private FileSystemManagedLedgerOffloader(OffloadPolicies conf, OrderedScheduler this.configuration = new Configuration(); if (conf.getFileSystemProfilePath() != null) { String[] paths = conf.getFileSystemProfilePath().split(","); - for (int i =0 ; i < paths.length; i++) { + for (int i = 0; i < paths.length; i++) { configuration.addResource(new Path(paths[i])); } } @@ -109,6 +111,7 @@ private FileSystemManagedLedgerOffloader(OffloadPolicies conf, OrderedScheduler .numThreads(conf.getManagedLedgerOffloadMaxThreads()) .name("offload-assignment").build(); } + @VisibleForTesting public FileSystemManagedLedgerOffloader(OffloadPolicies conf, OrderedScheduler scheduler, String testHDFSPath, String baseDir) throws IOException { this.offloadPolicies = conf; @@ -135,8 +138,8 @@ public Map getOffloadDriverMetadata() { } /* - * ledgerMetadata stored in an index of -1 - * */ + * ledgerMetadata stored in an index of -1 + * */ @Override public CompletableFuture offload(ReadHandle readHandle, UUID uuid, Map extraMetadata) { CompletableFuture promise = new CompletableFuture<>(); @@ -192,15 +195,16 @@ public void run() { long needToOffloadFirstEntryNumber = 0; CountDownLatch countDownLatch; //avoid prefetch too much data into memory - ArrayBlockingQueue tasks = new ArrayBlockingQueue<>(PREFETCH_ROUNDS); + Semaphore semaphore = new Semaphore(PREFETCH_ROUNDS); do { long end = Math.min(needToOffloadFirstEntryNumber + ENTRIES_PER_READ - 1, readHandle.getLastAddConfirmed()); log.debug("read ledger entries. start: {}, end: {}", needToOffloadFirstEntryNumber, end); LedgerEntries ledgerEntriesOnce = readHandle.readAsync(needToOffloadFirstEntryNumber, end).get(); - tasks.put(true); + semaphore.acquire(); countDownLatch = new CountDownLatch(1); - assignmentScheduler.chooseThread(ledgerId).submit(FileSystemWriter.create(ledgerEntriesOnce, dataWriter, tasks, - countDownLatch, haveOffloadEntryNumber, this)).addListener(() -> {}, Executors.newSingleThreadExecutor()); + assignmentScheduler.chooseThread(ledgerId).submit(FileSystemWriter.create(ledgerEntriesOnce, dataWriter, semaphore, + countDownLatch, haveOffloadEntryNumber, this)).addListener(() -> { + }, Executors.newSingleThreadExecutor()); needToOffloadFirstEntryNumber = end + 1; } while (needToOffloadFirstEntryNumber - 1 != readHandle.getLastAddConfirmed() && fileSystemWriteException == null); countDownLatch.await(); @@ -231,7 +235,7 @@ private static class FileSystemWriter implements Runnable { private CountDownLatch countDownLatch; private AtomicLong haveOffloadEntryNumber; private LedgerReader ledgerReader; - private ArrayBlockingQueue tasks; + private Semaphore semaphore; private Recycler.Handle recyclerHandle; private FileSystemWriter(Recycler.Handle recyclerHandle) { @@ -251,49 +255,45 @@ private void recycle() { this.haveOffloadEntryNumber = null; this.ledgerReader = null; this.ledgerEntriesOnce = null; - this.tasks = null; + this.semaphore = null; recyclerHandle.recycle(this); } - public static FileSystemWriter create(LedgerEntries ledgerEntriesOnce, MapFile.Writer dataWriter, ArrayBlockingQueue tasks, - CountDownLatch countDownLatch, AtomicLong haveOffloadEntryNumber, LedgerReader ledgerReader) { + public static FileSystemWriter create(LedgerEntries ledgerEntriesOnce, MapFile.Writer dataWriter, Semaphore semaphore, + CountDownLatch countDownLatch, AtomicLong haveOffloadEntryNumber, LedgerReader ledgerReader) { FileSystemWriter writer = RECYCLER.get(); writer.ledgerReader = ledgerReader; writer.dataWriter = dataWriter; writer.countDownLatch = countDownLatch; writer.haveOffloadEntryNumber = haveOffloadEntryNumber; writer.ledgerEntriesOnce = ledgerEntriesOnce; - writer.tasks = tasks; + writer.semaphore = semaphore; return writer; } @Override public void run() { - try { - if (ledgerReader.fileSystemWriteException == null) { - Iterator iterator = ledgerEntriesOnce.iterator(); - while (iterator.hasNext()) { - LedgerEntry entry = iterator.next(); - long entryId = entry.getEntryId(); - key.set(entryId); - try { - value.set(entry.getEntryBytes(), 0, entry.getEntryBytes().length); - dataWriter.append(key, value); - } catch (IOException e) { - ledgerReader.fileSystemWriteException = e; - break; - } - haveOffloadEntryNumber.incrementAndGet(); + if (ledgerReader.fileSystemWriteException == null) { + Iterator iterator = ledgerEntriesOnce.iterator(); + while (iterator.hasNext()) { + LedgerEntry entry = iterator.next(); + long entryId = entry.getEntryId(); + key.set(entryId); + try { + value.set(entry.getEntryBytes(), 0, entry.getEntryBytes().length); + dataWriter.append(key, value); + } catch (IOException e) { + ledgerReader.fileSystemWriteException = e; + break; } + haveOffloadEntryNumber.incrementAndGet(); } - countDownLatch.countDown(); - ledgerEntriesOnce.close(); - tasks.take(); - this.recycle(); - } catch (InterruptedException e) { - ledgerReader.fileSystemWriteException = e; } + countDownLatch.countDown(); + ledgerEntriesOnce.close(); + semaphore.release(); + this.recycle(); } } From bc7294cf0b3525090e38d7175d4f03a7f0963026 Mon Sep 17 00:00:00 2001 From: Xiaopeng Zhang Date: Sun, 12 Apr 2020 15:47:38 +0800 Subject: [PATCH 3/7] fix format --- .../filesystem/impl/FileSystemManagedLedgerOffloader.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java b/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java index 7ecc215186d76..cd81c95126a13 100644 --- a/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java +++ b/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java @@ -6,9 +6,9 @@ * to you under the Apache License, Version 2.0 (the * "License"); you may not use this file except in compliance * with the License. You may obtain a copy of the License at - *

- * http://www.apache.org/licenses/LICENSE-2.0 - *

+ * + * 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 From 01ae3bca7a1b7fd2629fccb06391822a31cb1e8e Mon Sep 17 00:00:00 2001 From: Xiaopeng Zhang Date: Sun, 12 Apr 2020 22:45:29 +0800 Subject: [PATCH 4/7] more conservative prefetch --- .../filesystem/impl/FileSystemManagedLedgerOffloader.java | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java b/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java index cd81c95126a13..c8ae1d7b3f1c9 100644 --- a/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java +++ b/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java @@ -43,7 +43,6 @@ import java.util.Map; import java.util.UUID; -import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executors; @@ -65,7 +64,6 @@ public class FileSystemManagedLedgerOffloader implements LedgerOffloader { private final FileSystem fileSystem; private OrderedScheduler scheduler; private static final long ENTRIES_PER_READ = 100; - private static final int PREFETCH_ROUNDS = 100; private OrderedScheduler assignmentScheduler; private OffloadPolicies offloadPolicies; @@ -195,7 +193,7 @@ public void run() { long needToOffloadFirstEntryNumber = 0; CountDownLatch countDownLatch; //avoid prefetch too much data into memory - Semaphore semaphore = new Semaphore(PREFETCH_ROUNDS); + Semaphore semaphore = new Semaphore(1); do { long end = Math.min(needToOffloadFirstEntryNumber + ENTRIES_PER_READ - 1, readHandle.getLastAddConfirmed()); log.debug("read ledger entries. start: {}, end: {}", needToOffloadFirstEntryNumber, end); From 4e3a7929de7415d2ac3c1a870087b3440e7156a3 Mon Sep 17 00:00:00 2001 From: Xiaopeng Zhang Date: Sun, 12 Apr 2020 22:50:06 +0800 Subject: [PATCH 5/7] fix format --- .../filesystem/impl/FileSystemManagedLedgerOffloader.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java b/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java index c8ae1d7b3f1c9..deb2254593d81 100644 --- a/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java +++ b/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java @@ -136,8 +136,8 @@ public Map getOffloadDriverMetadata() { } /* - * ledgerMetadata stored in an index of -1 - * */ + * ledgerMetadata stored in an index of -1 + * */ @Override public CompletableFuture offload(ReadHandle readHandle, UUID uuid, Map extraMetadata) { CompletableFuture promise = new CompletableFuture<>(); From 532bf5d5e5cc4c4541dbc474de261e0419954db3 Mon Sep 17 00:00:00 2001 From: Xiaopeng Zhang Date: Tue, 14 Apr 2020 20:33:59 +0800 Subject: [PATCH 6/7] make prefetch for offload configurable --- conf/broker.conf | 3 +++ .../org/apache/pulsar/broker/ServiceConfiguration.java | 6 ++++++ .../pulsar/common/policies/data/OffloadPolicies.java | 5 +++++ .../filesystem/impl/FileSystemManagedLedgerOffloader.java | 8 +++++--- 4 files changed, 19 insertions(+), 3 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index eb4d22d26920b..c49bcadcf9a61 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -837,6 +837,9 @@ managedLedgerOffloadDriver= # Maximum number of thread pool threads for ledger offloading managedLedgerOffloadMaxThreads=2 +# Maximum prefetch rounds for ledger reading for offloading +managedLedgerOffloadPrefetchRounds=1 + # Use Open Range-Set to cache unacked messages managedLedgerUnackedRangesOpenCacheSetEnabled=true 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 f6499107ff382..d3c6a62d828b2 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 @@ -1461,7 +1461,13 @@ public class ServiceConfiguration implements PulsarConfiguration { ) private int managedLedgerOffloadMaxThreads = 2; + @FieldContext( + category = CATEGORY_STORAGE_OFFLOADING, + doc = "Maximum prefetch rounds for ledger reading for offloading" + ) + private int managedLedgerOffloadPrefetchRounds = 1; /**** --- Transaction config variables --- ****/ + @FieldContext( category = CATEGORY_TRANSACTION, doc = "Enable transaction coordinator in broker" diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java index 5ccb75c88957a..4936923dfda6b 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java @@ -37,6 +37,7 @@ public class OffloadPolicies { public final static int DEFAULT_MAX_BLOCK_SIZE_IN_BYTES = 64 * 1024 * 1024; // 64MB public final static int DEFAULT_READ_BUFFER_SIZE_IN_BYTES = 1024 * 1024; // 1MB public final static int DEFAULT_OFFLOAD_MAX_THREADS = 2; + public final static int DEFAULT_OFFLOAD_MAX_PREFETCH_ROUNDS = 1; public final static String[] DRIVER_NAMES = {"S3", "aws-s3", "google-cloud-storage", "filesystem"}; public final static String DEFAULT_OFFLOADER_DIRECTORY = "./offloaders"; public final static long DEFAULT_OFFLOAD_THRESHOLD_IN_BYTES = -1; @@ -46,6 +47,7 @@ public class OffloadPolicies { private String offloadersDirectory = DEFAULT_OFFLOADER_DIRECTORY; private String managedLedgerOffloadDriver = null; private int managedLedgerOffloadMaxThreads = DEFAULT_OFFLOAD_MAX_THREADS; + private int managedLedgerOffloadPrefetchRounds = DEFAULT_OFFLOAD_MAX_PREFETCH_ROUNDS; private long managedLedgerOffloadThresholdInBytes = DEFAULT_OFFLOAD_THRESHOLD_IN_BYTES; private Long managedLedgerOffloadDeletionLagInMillis = DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS; @@ -161,6 +163,7 @@ public int hashCode() { return Objects.hash( managedLedgerOffloadDriver, managedLedgerOffloadMaxThreads, + managedLedgerOffloadPrefetchRounds, managedLedgerOffloadThresholdInBytes, managedLedgerOffloadDeletionLagInMillis, s3ManagedLedgerOffloadRegion, @@ -190,6 +193,7 @@ public boolean equals(Object obj) { OffloadPolicies other = (OffloadPolicies) obj; return Objects.equals(managedLedgerOffloadDriver, other.getManagedLedgerOffloadDriver()) && Objects.equals(managedLedgerOffloadMaxThreads, other.getManagedLedgerOffloadMaxThreads()) + && Objects.equals(managedLedgerOffloadPrefetchRounds, other.getManagedLedgerOffloadPrefetchRounds()) && Objects.equals(managedLedgerOffloadThresholdInBytes, other.getManagedLedgerOffloadThresholdInBytes()) && Objects.equals(managedLedgerOffloadDeletionLagInMillis, @@ -222,6 +226,7 @@ public String toString() { return MoreObjects.toStringHelper(this) .add("managedLedgerOffloadDriver", managedLedgerOffloadDriver) .add("managedLedgerOffloadMaxThreads", managedLedgerOffloadMaxThreads) + .add("managedLedgerOffloadPrefetchRounds", managedLedgerOffloadPrefetchRounds) .add("managedLedgerOffloadThresholdInBytes", managedLedgerOffloadThresholdInBytes) .add("managedLedgerOffloadDeletionLagInMillis", managedLedgerOffloadDeletionLagInMillis) .add("s3ManagedLedgerOffloadRegion", s3ManagedLedgerOffloadRegion) diff --git a/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java b/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java index deb2254593d81..5438459c43a56 100644 --- a/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java +++ b/tiered-storage/file-system/src/main/java/org/apache/bookkeeper/mledger/offload/filesystem/impl/FileSystemManagedLedgerOffloader.java @@ -141,7 +141,7 @@ public Map getOffloadDriverMetadata() { @Override public CompletableFuture offload(ReadHandle readHandle, UUID uuid, Map extraMetadata) { CompletableFuture promise = new CompletableFuture<>(); - scheduler.chooseThread(readHandle.getId()).submit(new LedgerReader(readHandle, uuid, extraMetadata, promise, storageBasePath, configuration, assignmentScheduler)); + scheduler.chooseThread(readHandle.getId()).submit(new LedgerReader(readHandle, uuid, extraMetadata, promise, storageBasePath, configuration, assignmentScheduler, offloadPolicies.getManagedLedgerOffloadPrefetchRounds())); return promise; } @@ -155,9 +155,10 @@ private static class LedgerReader implements Runnable { private final Configuration configuration; volatile Exception fileSystemWriteException = null; private OrderedScheduler assignmentScheduler; + private int managedLedgerOffloadPrefetchRounds = 1; private LedgerReader(ReadHandle readHandle, UUID uuid, Map extraMetadata, CompletableFuture promise, - String storageBasePath, Configuration configuration, OrderedScheduler assignmentScheduler) { + String storageBasePath, Configuration configuration, OrderedScheduler assignmentScheduler, int managedLedgerOffloadPrefetchRounds) { this.readHandle = readHandle; this.uuid = uuid; this.extraMetadata = extraMetadata; @@ -165,6 +166,7 @@ private LedgerReader(ReadHandle readHandle, UUID uuid, Map extra this.storageBasePath = storageBasePath; this.configuration = configuration; this.assignmentScheduler = assignmentScheduler; + this.managedLedgerOffloadPrefetchRounds = managedLedgerOffloadPrefetchRounds; } @Override @@ -193,7 +195,7 @@ public void run() { long needToOffloadFirstEntryNumber = 0; CountDownLatch countDownLatch; //avoid prefetch too much data into memory - Semaphore semaphore = new Semaphore(1); + Semaphore semaphore = new Semaphore(managedLedgerOffloadPrefetchRounds); do { long end = Math.min(needToOffloadFirstEntryNumber + ENTRIES_PER_READ - 1, readHandle.getLastAddConfirmed()); log.debug("read ledger entries. start: {}, end: {}", needToOffloadFirstEntryNumber, end); From 2ac1c38909a1c2a598b4518098a05f557c110b31 Mon Sep 17 00:00:00 2001 From: Xiaopeng Zhang Date: Tue, 14 Apr 2020 21:00:20 +0800 Subject: [PATCH 7/7] fix format --- .../java/org/apache/pulsar/broker/ServiceConfiguration.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 d3c6a62d828b2..9a4879dfbc8df 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 @@ -1466,8 +1466,8 @@ public class ServiceConfiguration implements PulsarConfiguration { doc = "Maximum prefetch rounds for ledger reading for offloading" ) private int managedLedgerOffloadPrefetchRounds = 1; - /**** --- Transaction config variables --- ****/ + /**** --- Transaction config variables --- ****/ @FieldContext( category = CATEGORY_TRANSACTION, doc = "Enable transaction coordinator in broker"