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 @@ -56,6 +56,7 @@
import org.apache.accumulo.core.spi.crypto.CryptoService;
import org.apache.accumulo.core.spi.crypto.FileDecrypter;
import org.apache.accumulo.core.spi.crypto.FileEncrypter;
import org.apache.accumulo.core.spi.fs.VolumeChooserEnvironment;
import org.apache.accumulo.core.tabletserver.log.LogEntry;
import org.apache.accumulo.core.util.Pair;
import org.apache.accumulo.core.util.threads.Threads;
Expand All @@ -74,7 +75,6 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import com.google.common.base.Joiner;
import com.google.common.base.Preconditions;

/**
Expand Down Expand Up @@ -116,12 +116,11 @@ public LogHeaderIncompleteException(EOFException cause) {
}
}

private final LinkedBlockingQueue<DfsLogger.LogWork> workQueue = new LinkedBlockingQueue<>();
private final LinkedBlockingQueue<LogWork> workQueue = new LinkedBlockingQueue<>();

private final Object closeLock = new Object();

private static final DfsLogger.LogWork CLOSED_MARKER =
new DfsLogger.LogWork(null, Durability.FLUSH);
private static final LogWork CLOSED_MARKER = new LogWork(null, Durability.FLUSH);

private static final LogFileValue EMPTY = new LogFileValue();

Expand All @@ -130,9 +129,19 @@ public LogHeaderIncompleteException(EOFException cause) {
private class LogSyncingTask implements Runnable {
private int expectedReplication = 0;

private final AtomicLong syncCounter;
private final AtomicLong flushCounter;
private final long slowFlushMillis;

LogSyncingTask(AtomicLong syncCounter, AtomicLong flushCounter, long slowFlushMillis) {
this.syncCounter = syncCounter;
this.flushCounter = flushCounter;
this.slowFlushMillis = slowFlushMillis;
}

@Override
public void run() {
ArrayList<DfsLogger.LogWork> work = new ArrayList<>();
ArrayList<LogWork> work = new ArrayList<>();
boolean sawClosedMarker = false;
while (!sawClosedMarker) {
work.clear();
Expand Down Expand Up @@ -205,7 +214,7 @@ public void run() {
}
}

for (DfsLogger.LogWork logWork : work) {
for (LogWork logWork : work) {
if (logWork == CLOSED_MARKER) {
sawClosedMarker = true;
} else {
Expand All @@ -215,9 +224,9 @@ public void run() {
}
}

private void fail(ArrayList<DfsLogger.LogWork> work, Exception ex, String why) {
private void fail(ArrayList<LogWork> work, Exception ex, String why) {
log.warn("Exception {} {}", why, ex, ex);
for (DfsLogger.LogWork logWork : work) {
for (LogWork logWork : work) {
logWork.exception = ex;
}
}
Expand Down Expand Up @@ -288,32 +297,46 @@ public int hashCode() {
return logEntry.hashCode();
}

private final ServerContext context;
private FSDataOutputStream logFile;
private DataOutputStream encryptingLogFile = null;
private LogEntry logEntry;
private final LogEntry logEntry;
private Thread syncThread;

private AtomicLong syncCounter;
private AtomicLong flushCounter;
private final long slowFlushMillis;
private long writes = 0;

public DfsLogger(ServerContext context, AtomicLong syncCounter, AtomicLong flushCounter) {
this(context, null);
this.syncCounter = syncCounter;
this.flushCounter = flushCounter;
/**
* Create a new DfsLogger with the provided characteristics.
*/
public static DfsLogger createNew(ServerContext context, AtomicLong syncCounter,
AtomicLong flushCounter, String address) throws IOException {

String filename = UUID.randomUUID().toString();
String addressForFilename = address.replace(':', '+');

var chooserEnv =
new VolumeChooserEnvironmentImpl(VolumeChooserEnvironment.Scope.LOGGER, context);
String logPath =
context.getVolumeManager().choose(chooserEnv, context.getBaseUris()) + Path.SEPARATOR
+ Constants.WAL_DIR + Path.SEPARATOR + addressForFilename + Path.SEPARATOR + filename;

LogEntry log = LogEntry.fromPath(logPath);
DfsLogger dfsLogger = new DfsLogger(log);
long slowFlushMillis =
context.getConfiguration().getTimeInMillis(Property.TSERV_SLOW_FLUSH_MILLIS);
dfsLogger.open(context, logPath, filename, address, syncCounter, flushCounter, slowFlushMillis);
return dfsLogger;
}

/**
* Reference a pre-existing log file.
*
* @param logEntry the "log" entry in +r/!0
*/
public DfsLogger(ServerContext context, LogEntry logEntry) {
this.context = context;
this.slowFlushMillis =
context.getConfiguration().getTimeInMillis(Property.TSERV_SLOW_FLUSH_MILLIS);
public static DfsLogger fromLogEntry(LogEntry logEntry) {
return new DfsLogger(logEntry);
}

private DfsLogger(LogEntry logEntry) {
this.logEntry = logEntry;
}

Expand Down Expand Up @@ -372,19 +395,14 @@ public static DataInputStream getDecryptingStream(FSDataInputStream input,
*
* @param address The address of the host using this WAL
*/
public synchronized void open(String address) throws IOException {
String filename = UUID.randomUUID().toString();
private synchronized void open(ServerContext context, String logPath, String filename,
String address, AtomicLong syncCounter, AtomicLong flushCounter, long slowFlushMillis)
throws IOException {
log.debug("Address is {}", address);
String logger = Joiner.on("+").join(address.split(":"));

log.debug("DfsLogger.open() begin");
VolumeManager fs = context.getVolumeManager();

var chooserEnv = new VolumeChooserEnvironmentImpl(
org.apache.accumulo.core.spi.fs.VolumeChooserEnvironment.Scope.LOGGER, context);
String logPath = fs.choose(chooserEnv, context.getBaseUris()) + Path.SEPARATOR
+ Constants.WAL_DIR + Path.SEPARATOR + logger + Path.SEPARATOR + filename;
this.logEntry = LogEntry.fromPath(logPath);
VolumeManager fs = context.getVolumeManager();

LoggerOperation op;
var serverConf = context.getConfiguration();
Expand Down Expand Up @@ -453,7 +471,8 @@ public synchronized void open(String address) throws IOException {
throw new IOException(ex);
}

syncThread = Threads.createThread("Accumulo WALog thread " + this, new LogSyncingTask());
syncThread = Threads.createThread("Accumulo WALog thread " + this,
new LogSyncingTask(syncCounter, flushCounter, slowFlushMillis));
syncThread.start();
op.await();
log.debug("Got new write-ahead log: {}", this);
Expand Down Expand Up @@ -548,7 +567,7 @@ private LoggerOperation logKeyData(LogFileKey key, Durability d) throws IOExcept

private LoggerOperation logFileData(List<Pair<LogFileKey,LogFileValue>> keys,
Durability durability) throws IOException {
DfsLogger.LogWork work = new DfsLogger.LogWork(new CountDownLatch(1), durability);
LogWork work = new LogWork(new CountDownLatch(1), durability);
try {
for (Pair<LogFileKey,LogFileValue> pair : keys) {
write(pair.getFirst(), pair.getSecond());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -266,8 +266,8 @@ private synchronized void startLogMaker() {
DfsLogger alog = null;

try {
alog = new DfsLogger(tserver.getContext(), syncCounter, flushCounter);
alog.open(tserver.getClientAddressString());
alog = DfsLogger.createNew(tserver.getContext(), syncCounter, flushCounter,
tserver.getClientAddressString());
} catch (Exception t) {
log.error("Failed to open WAL", t);
// the log is not advertised in ZK yet, so we can just delete it if it exists
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -369,7 +369,7 @@ public Tablet(final TabletServer tabletServer, final KeyExtent extent,
// make some closed references that represent the recovered logs
currentLogs = new HashSet<>();
for (LogEntry logEntry : logEntries) {
currentLogs.add(new DfsLogger(tabletServer.getContext(), logEntry));
currentLogs.add(DfsLogger.fromLogEntry(logEntry));
}

rebuildReferencedLogs();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,46 +18,24 @@
*/
package org.apache.accumulo.tserver;

import static org.easymock.EasyMock.createMock;
import static org.easymock.EasyMock.expect;
import static org.easymock.EasyMock.replay;
import static org.easymock.EasyMock.verify;
import static org.junit.jupiter.api.Assertions.assertEquals;

import java.util.Collections;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Set;

import org.apache.accumulo.core.conf.DefaultConfiguration;
import org.apache.accumulo.core.tabletserver.log.LogEntry;
import org.apache.accumulo.server.ServerContext;
import org.apache.accumulo.tserver.log.DfsLogger;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;

import com.google.common.collect.Sets;

public class WalRemovalOrderTest {

private ServerContext context;

@BeforeEach
private void createMocks() {
context = createMock(ServerContext.class);
expect(context.getConfiguration()).andReturn(DefaultConfiguration.getInstance()).anyTimes();
replay(context);
}

@AfterEach
private void verifyMocks() {
verify(context);
}

private DfsLogger mockLogger(String filename) {
var mockLogEntry = LogEntry.fromPath(filename + "+1234/11111111-1111-1111-1111-111111111111");
return new DfsLogger(context, mockLogEntry);
return DfsLogger.fromLogEntry(mockLogEntry);
}

private LinkedHashSet<DfsLogger> mockLoggers(String... logs) {
Expand Down