Skip to content
Closed
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 @@ -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;
Expand Down Expand Up @@ -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);
Expand All @@ -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<FileSystemWriter> recyclerHandle;

private FileSystemWriter(Recycler.Handle<FileSystemWriter> recyclerHandle) {
this.recyclerHandle = recyclerHandle;
}

private static final Recycler<FileSystemWriter> RECYCLER = new Recycler<FileSystemWriter>() {
@Override
protected FileSystemWriter newObject(Recycler.Handle<FileSystemWriter> 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
Expand All @@ -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;
Expand All @@ -255,6 +280,7 @@ public void run() {
}
}
countDownLatch.countDown();
this.recycle();
}
}

Expand Down