From 87fd4eea624727d3911492aa9e7f4ffa961d5fb9 Mon Sep 17 00:00:00 2001 From: arbaazkhan1 Date: Wed, 24 Jan 2024 17:03:16 -0500 Subject: [PATCH 1/4] Refactored DfsLogger constructors to static methods that now call a private constructors --- .../accumulo/tserver/log/DfsLogger.java | 49 +++++++++++++------ .../tserver/log/TabletServerLogger.java | 3 +- .../accumulo/tserver/tablet/Tablet.java | 2 +- .../accumulo/tserver/WalRemovalOrderTest.java | 2 +- 4 files changed, 38 insertions(+), 18 deletions(-) diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java index 91e9382b71f..366b38d9bd3 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java @@ -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; @@ -116,12 +117,11 @@ public LogHeaderIncompleteException(EOFException cause) { } } - private final LinkedBlockingQueue workQueue = new LinkedBlockingQueue<>(); + private final LinkedBlockingQueue 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(); @@ -132,7 +132,7 @@ private class LogSyncingTask implements Runnable { @Override public void run() { - ArrayList work = new ArrayList<>(); + ArrayList work = new ArrayList<>(); boolean sawClosedMarker = false; while (!sawClosedMarker) { work.clear(); @@ -205,7 +205,7 @@ public void run() { } } - for (DfsLogger.LogWork logWork : work) { + for (LogWork logWork : work) { if (logWork == CLOSED_MARKER) { sawClosedMarker = true; } else { @@ -215,9 +215,9 @@ public void run() { } } - private void fail(ArrayList work, Exception ex, String why) { + private void fail(ArrayList work, Exception ex, String why) { log.warn("Exception {} {}", why, ex, ex); - for (DfsLogger.LogWork logWork : work) { + for (LogWork logWork : work) { logWork.exception = ex; } } @@ -291,7 +291,7 @@ public int 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; @@ -299,8 +299,28 @@ public int hashCode() { private final long slowFlushMillis; private long writes = 0; - public DfsLogger(ServerContext context, AtomicLong syncCounter, AtomicLong flushCounter) { - this(context, null); + public static DfsLogger fromCounters(ServerContext context, AtomicLong syncCounter, + AtomicLong flushCounter, String address) { + + String filename = UUID.randomUUID().toString(); + String logger = Joiner.on("+").join(address.split(":")); + VolumeManager fs = context.getVolumeManager(); + + var chooserEnv = + new VolumeChooserEnvironmentImpl(VolumeChooserEnvironment.Scope.LOGGER, context); + String logPath = fs.choose(chooserEnv, context.getBaseUris()) + Path.SEPARATOR + + Constants.WAL_DIR + Path.SEPARATOR + logger + Path.SEPARATOR + filename; + + return new DfsLogger(context, syncCounter, flushCounter, LogEntry.fromPath(logPath)); + } + + public static DfsLogger fromExistingLogEntry(ServerContext context, LogEntry logEntry) { + return new DfsLogger(context, logEntry); + } + + private DfsLogger(ServerContext context, AtomicLong syncCounter, AtomicLong flushCounter, + LogEntry logEntry) { + this(context, logEntry); this.syncCounter = syncCounter; this.flushCounter = flushCounter; } @@ -310,7 +330,7 @@ public DfsLogger(ServerContext context, AtomicLong syncCounter, AtomicLong flush * * @param logEntry the "log" entry in +r/!0 */ - public DfsLogger(ServerContext context, LogEntry logEntry) { + private DfsLogger(ServerContext context, LogEntry logEntry) { this.context = context; this.slowFlushMillis = context.getConfiguration().getTimeInMillis(Property.TSERV_SLOW_FLUSH_MILLIS); @@ -380,11 +400,10 @@ public synchronized void open(String address) throws IOException { log.debug("DfsLogger.open() begin"); VolumeManager fs = context.getVolumeManager(); - var chooserEnv = new VolumeChooserEnvironmentImpl( - org.apache.accumulo.core.spi.fs.VolumeChooserEnvironment.Scope.LOGGER, context); + var chooserEnv = + new VolumeChooserEnvironmentImpl(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); LoggerOperation op; var serverConf = context.getConfiguration(); @@ -548,7 +567,7 @@ private LoggerOperation logKeyData(LogFileKey key, Durability d) throws IOExcept private LoggerOperation logFileData(List> 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 pair : keys) { write(pair.getFirst(), pair.getSecond()); diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/TabletServerLogger.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/TabletServerLogger.java index 3103d010441..807e8ce0042 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/TabletServerLogger.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/TabletServerLogger.java @@ -266,7 +266,8 @@ private synchronized void startLogMaker() { DfsLogger alog = null; try { - alog = new DfsLogger(tserver.getContext(), syncCounter, flushCounter); + alog = DfsLogger.fromCounters(tserver.getContext(), syncCounter, flushCounter, + tserver.getClientAddressString()); alog.open(tserver.getClientAddressString()); } catch (Exception t) { log.error("Failed to open WAL", t); diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Tablet.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Tablet.java index 9d6043ef3ed..80433ed5876 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Tablet.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Tablet.java @@ -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.fromExistingLogEntry(tabletServer.getContext(), logEntry)); } rebuildReferencedLogs(); diff --git a/server/tserver/src/test/java/org/apache/accumulo/tserver/WalRemovalOrderTest.java b/server/tserver/src/test/java/org/apache/accumulo/tserver/WalRemovalOrderTest.java index 69d7253a725..5ce9634aef5 100644 --- a/server/tserver/src/test/java/org/apache/accumulo/tserver/WalRemovalOrderTest.java +++ b/server/tserver/src/test/java/org/apache/accumulo/tserver/WalRemovalOrderTest.java @@ -57,7 +57,7 @@ private void verifyMocks() { private DfsLogger mockLogger(String filename) { var mockLogEntry = LogEntry.fromPath(filename + "+1234/11111111-1111-1111-1111-111111111111"); - return new DfsLogger(context, mockLogEntry); + return DfsLogger.fromExistingLogEntry(context, mockLogEntry); } private LinkedHashSet mockLoggers(String... logs) { From 3ce1ec1c76d670003d75a284c424e6d721ed8418 Mon Sep 17 00:00:00 2001 From: arbaazkhan1 Date: Fri, 26 Jan 2024 11:28:59 -0500 Subject: [PATCH 2/4] Open method is now private and is being called from a new Dfslogger variable located in fromCounters. --- .../accumulo/tserver/log/DfsLogger.java | 21 +++++++------------ .../tserver/log/TabletServerLogger.java | 1 - 2 files changed, 8 insertions(+), 14 deletions(-) diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java index 366b38d9bd3..9d77b38380a 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java @@ -300,7 +300,7 @@ public int hashCode() { private long writes = 0; public static DfsLogger fromCounters(ServerContext context, AtomicLong syncCounter, - AtomicLong flushCounter, String address) { + AtomicLong flushCounter, String address) throws IOException { String filename = UUID.randomUUID().toString(); String logger = Joiner.on("+").join(address.split(":")); @@ -311,15 +311,18 @@ public static DfsLogger fromCounters(ServerContext context, AtomicLong syncCount String logPath = fs.choose(chooserEnv, context.getBaseUris()) + Path.SEPARATOR + Constants.WAL_DIR + Path.SEPARATOR + logger + Path.SEPARATOR + filename; - return new DfsLogger(context, syncCounter, flushCounter, LogEntry.fromPath(logPath)); + LogEntry log = LogEntry.fromPath(logPath); + DfsLogger dfsLogger = new DfsLogger(context, log, syncCounter, flushCounter); + dfsLogger.open(fs, logPath, filename, address); + + return dfsLogger; } public static DfsLogger fromExistingLogEntry(ServerContext context, LogEntry logEntry) { return new DfsLogger(context, logEntry); } - private DfsLogger(ServerContext context, AtomicLong syncCounter, AtomicLong flushCounter, - LogEntry logEntry) { + private DfsLogger(ServerContext context, LogEntry logEntry, AtomicLong syncCounter, AtomicLong flushCounter) { this(context, logEntry); this.syncCounter = syncCounter; this.flushCounter = flushCounter; @@ -392,18 +395,10 @@ 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(VolumeManager fs, String logPath, String filename, String address) 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(VolumeChooserEnvironment.Scope.LOGGER, context); - String logPath = fs.choose(chooserEnv, context.getBaseUris()) + Path.SEPARATOR - + Constants.WAL_DIR + Path.SEPARATOR + logger + Path.SEPARATOR + filename; LoggerOperation op; var serverConf = context.getConfiguration(); diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/TabletServerLogger.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/TabletServerLogger.java index 807e8ce0042..45578da41b6 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/TabletServerLogger.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/TabletServerLogger.java @@ -268,7 +268,6 @@ private synchronized void startLogMaker() { try { alog = DfsLogger.fromCounters(tserver.getContext(), syncCounter, flushCounter, tserver.getClientAddressString()); - alog.open(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 From 10a7b59d080ad0563e7232116b9275b52344d773 Mon Sep 17 00:00:00 2001 From: arbaazkhan1 Date: Mon, 29 Jan 2024 10:18:40 -0500 Subject: [PATCH 3/4] Open method is now private and is being called from a new Dfslogger variable located in fromCounters. fix bug --- .../java/org/apache/accumulo/tserver/log/DfsLogger.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java index 9d77b38380a..5e0a3f42daa 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java @@ -322,7 +322,8 @@ public static DfsLogger fromExistingLogEntry(ServerContext context, LogEntry log return new DfsLogger(context, logEntry); } - private DfsLogger(ServerContext context, LogEntry logEntry, AtomicLong syncCounter, AtomicLong flushCounter) { + private DfsLogger(ServerContext context, LogEntry logEntry, AtomicLong syncCounter, + AtomicLong flushCounter) { this(context, logEntry); this.syncCounter = syncCounter; this.flushCounter = flushCounter; @@ -395,7 +396,8 @@ public static DataInputStream getDecryptingStream(FSDataInputStream input, * * @param address The address of the host using this WAL */ - private synchronized void open(VolumeManager fs, String logPath, String filename, String address) throws IOException { + private synchronized void open(VolumeManager fs, String logPath, String filename, String address) + throws IOException { log.debug("Address is {}", address); log.debug("DfsLogger.open() begin"); From 0ec45b5b30ed576911f7f363c1baf44ca6fc431a Mon Sep 17 00:00:00 2001 From: Christopher Tubbs Date: Wed, 7 Feb 2024 16:14:32 -0500 Subject: [PATCH 4/4] Code review feedback * Push params used only for creating a new file down into the open method and remove from DfsLogger fields * Remove unneeded second private constructor * Tweak method names * Remove unnecessary mock object from test --- .../accumulo/tserver/log/DfsLogger.java | 63 ++++++++++--------- .../tserver/log/TabletServerLogger.java | 2 +- .../accumulo/tserver/tablet/Tablet.java | 2 +- .../accumulo/tserver/WalRemovalOrderTest.java | 24 +------ 4 files changed, 36 insertions(+), 55 deletions(-) diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java index 5e0a3f42daa..c0ed0aed396 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/DfsLogger.java @@ -75,7 +75,6 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import com.google.common.base.Joiner; import com.google.common.base.Preconditions; /** @@ -130,6 +129,16 @@ 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 work = new ArrayList<>(); @@ -288,56 +297,46 @@ public int hashCode() { return logEntry.hashCode(); } - private final ServerContext context; private FSDataOutputStream logFile; private DataOutputStream encryptingLogFile = null; private final LogEntry logEntry; private Thread syncThread; - private AtomicLong syncCounter; - private AtomicLong flushCounter; - private final long slowFlushMillis; private long writes = 0; - public static DfsLogger fromCounters(ServerContext context, AtomicLong syncCounter, + /** + * 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 logger = Joiner.on("+").join(address.split(":")); - VolumeManager fs = context.getVolumeManager(); + String addressForFilename = address.replace(':', '+'); var chooserEnv = new VolumeChooserEnvironmentImpl(VolumeChooserEnvironment.Scope.LOGGER, context); - String logPath = fs.choose(chooserEnv, context.getBaseUris()) + Path.SEPARATOR - + Constants.WAL_DIR + Path.SEPARATOR + logger + Path.SEPARATOR + filename; + 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(context, log, syncCounter, flushCounter); - dfsLogger.open(fs, logPath, filename, address); - + 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; } - public static DfsLogger fromExistingLogEntry(ServerContext context, LogEntry logEntry) { - return new DfsLogger(context, logEntry); - } - - private DfsLogger(ServerContext context, LogEntry logEntry, AtomicLong syncCounter, - AtomicLong flushCounter) { - this(context, logEntry); - this.syncCounter = syncCounter; - this.flushCounter = flushCounter; - } - /** * Reference a pre-existing log file. * * @param logEntry the "log" entry in +r/!0 */ - private 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; } @@ -396,12 +395,15 @@ public static DataInputStream getDecryptingStream(FSDataInputStream input, * * @param address The address of the host using this WAL */ - private synchronized void open(VolumeManager fs, String logPath, String filename, String address) + 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); log.debug("DfsLogger.open() begin"); + VolumeManager fs = context.getVolumeManager(); + LoggerOperation op; var serverConf = context.getConfiguration(); try { @@ -469,7 +471,8 @@ private synchronized void open(VolumeManager fs, String logPath, String filename 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); diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/TabletServerLogger.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/TabletServerLogger.java index 45578da41b6..dc66adc26af 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/log/TabletServerLogger.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/log/TabletServerLogger.java @@ -266,7 +266,7 @@ private synchronized void startLogMaker() { DfsLogger alog = null; try { - alog = DfsLogger.fromCounters(tserver.getContext(), syncCounter, flushCounter, + alog = DfsLogger.createNew(tserver.getContext(), syncCounter, flushCounter, tserver.getClientAddressString()); } catch (Exception t) { log.error("Failed to open WAL", t); diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Tablet.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Tablet.java index 9d6571e5aa7..90af6d43268 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Tablet.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Tablet.java @@ -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(DfsLogger.fromExistingLogEntry(tabletServer.getContext(), logEntry)); + currentLogs.add(DfsLogger.fromLogEntry(logEntry)); } rebuildReferencedLogs(); diff --git a/server/tserver/src/test/java/org/apache/accumulo/tserver/WalRemovalOrderTest.java b/server/tserver/src/test/java/org/apache/accumulo/tserver/WalRemovalOrderTest.java index 5ce9634aef5..2e682b4231a 100644 --- a/server/tserver/src/test/java/org/apache/accumulo/tserver/WalRemovalOrderTest.java +++ b/server/tserver/src/test/java/org/apache/accumulo/tserver/WalRemovalOrderTest.java @@ -18,10 +18,6 @@ */ 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; @@ -29,35 +25,17 @@ 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 DfsLogger.fromExistingLogEntry(context, mockLogEntry); + return DfsLogger.fromLogEntry(mockLogEntry); } private LinkedHashSet mockLoggers(String... logs) {