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..0146b25c57050 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; @@ -193,7 +194,7 @@ public void run() { log.debug("read ledger entries. start: {}, end: {}", needToOffloadFirstEntryNumber, end); LedgerEntries ledgerEntriesOnce = readHandle.readAsync(needToOffloadFirstEntryNumber, end).get(); countDownLatch = new CountDownLatch(1); - assignmentScheduler.chooseThread(ledgerId).submit(new FileSystemWriter(ledgerEntriesOnce, dataWriter, + assignmentScheduler.chooseThread(ledgerId).submit(FileSystemWriter.create(ledgerEntriesOnce, dataWriter, countDownLatch, haveOffloadEntryNumber, this)).addListener(() -> {}, Executors.newSingleThreadExecutor()); needToOffloadFirstEntryNumber = end + 1; } while (needToOffloadFirstEntryNumber - 1 != readHandle.getLastAddConfirmed() && fileSystemWriteException == null); @@ -216,24 +217,47 @@ 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 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 void recycle() { + this.dataWriter = null; + this.countDownLatch = null; + this.haveOffloadEntryNumber = null; + this.ledgerReader = null; + this.ledgerEntriesOnce = null; + recyclerHandle.recycle(this); + } - private FileSystemWriter(LedgerEntries ledgerEntriesOnce, MapFile.Writer dataWriter, + public static FileSystemWriter create(LedgerEntries ledgerEntriesOnce, MapFile.Writer dataWriter, 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; + return writer; } @Override @@ -247,6 +271,7 @@ public void run() { try { value.set(entry.getEntryBytes(), 0, entry.getEntryBytes().length); dataWriter.append(key, value); + entry.close(); } catch (IOException e) { ledgerReader.fileSystemWriteException = e; break; @@ -255,6 +280,7 @@ public void run() { } } countDownLatch.countDown(); + this.recycle(); } }