From 4775c6ad4deaec91e6ed35aebdd3a73f169631f9 Mon Sep 17 00:00:00 2001 From: Daniel Roberts ddanielr Date: Tue, 21 Jul 2026 16:27:08 +0000 Subject: [PATCH 01/11] Update to the latest log4j and spotbugs versions Bumps versions and fixes deprecated builder pattern --- .../apache/accumulo/core/logging/EscalatingLoggerTest.java | 4 ++-- pom.xml | 6 +++--- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/core/src/test/java/org/apache/accumulo/core/logging/EscalatingLoggerTest.java b/core/src/test/java/org/apache/accumulo/core/logging/EscalatingLoggerTest.java index a119601cf0a..33f340a45fb 100644 --- a/core/src/test/java/org/apache/accumulo/core/logging/EscalatingLoggerTest.java +++ b/core/src/test/java/org/apache/accumulo/core/logging/EscalatingLoggerTest.java @@ -48,8 +48,8 @@ public void test() throws InterruptedException { // Programatically modify the Log4j2 Logging configuration to add an appender LoggerContext ctx = LoggerContext.getContext(false); Configuration cfg = ctx.getConfiguration(); - PatternLayout layout = PatternLayout.newBuilder().withConfiguration(cfg) - .withPattern(PatternLayout.SIMPLE_CONVERSION_PATTERN).build(); + PatternLayout layout = PatternLayout.newBuilder().setConfiguration(cfg) + .setPattern(PatternLayout.SIMPLE_CONVERSION_PATTERN).build(); Appender appender = WriterAppender.createAppender(layout, null, writer, "EscalatingLoggerTestAppender", false, true); appender.start(); diff --git a/pom.xml b/pom.xml index 9df56e7e19b..470c5be634e 100644 --- a/pom.xml +++ b/pom.xml @@ -171,7 +171,7 @@ under the License. 5.8.0 2.41.0 3.3.6 - 2.25.4 + 2.26.1 1.48.0 2.0.9 2.0.17 @@ -277,7 +277,7 @@ under the License. com.github.spotbugs spotbugs-annotations - 4.9.3 + 4.10.3 com.google.auto.service @@ -780,7 +780,7 @@ under the License. com.github.spotbugs spotbugs-maven-plugin - 4.9.8.2 + 4.10.3.0 true Max From f75dffc0c5c625821ac7e8f98acc520aea8f1f75 Mon Sep 17 00:00:00 2001 From: Daniel Roberts ddanielr Date: Tue, 21 Jul 2026 17:09:34 +0000 Subject: [PATCH 02/11] Add exclusion filter for spotbugs --- core/src/main/spotbugs/exclude-filter.xml | 4 ++++ start/src/main/spotbugs/exclude-filter.xml | 4 ++++ 2 files changed, 8 insertions(+) diff --git a/core/src/main/spotbugs/exclude-filter.xml b/core/src/main/spotbugs/exclude-filter.xml index f1776b30607..dd1d0d7b6d8 100644 --- a/core/src/main/spotbugs/exclude-filter.xml +++ b/core/src/main/spotbugs/exclude-filter.xml @@ -43,4 +43,8 @@ + + + + diff --git a/start/src/main/spotbugs/exclude-filter.xml b/start/src/main/spotbugs/exclude-filter.xml index c5b24cd502c..df55f1f613d 100644 --- a/start/src/main/spotbugs/exclude-filter.xml +++ b/start/src/main/spotbugs/exclude-filter.xml @@ -27,4 +27,8 @@ + + + + From ce8dc928ac92ce46ee417512547f6535a96bf2d0 Mon Sep 17 00:00:00 2001 From: Daniel Roberts ddanielr Date: Wed, 22 Jul 2026 03:46:08 +0000 Subject: [PATCH 03/11] Added private locks and removed exclusion filter Added private locks for methods that were synchronized to fix spotbugs findings --- .../java/org/apache/accumulo/start/Main.java | 48 +++++++++------- .../classloader/AccumuloClassLoader.java | 55 ++++++++++--------- .../vfs/AccumuloVFSClassLoader.java | 34 ++++++------ start/src/main/spotbugs/exclude-filter.xml | 4 -- 4 files changed, 74 insertions(+), 67 deletions(-) diff --git a/start/src/main/java/org/apache/accumulo/start/Main.java b/start/src/main/java/org/apache/accumulo/start/Main.java index e736821e215..2d58790377d 100644 --- a/start/src/main/java/org/apache/accumulo/start/Main.java +++ b/start/src/main/java/org/apache/accumulo/start/Main.java @@ -40,6 +40,7 @@ public class Main { private static ClassLoader classLoader; private static Class vfsClassLoader; private static Map servicesMap; + private static final Object lock = new Object(); public static void main(final String[] args) throws Exception { // Preload classes that cause a deadlock between the ServiceLoader and the DFSClient when @@ -91,29 +92,32 @@ public static void main(final String[] args) throws Exception { } - public static synchronized ClassLoader getClassLoader() { - if (classLoader == null) { - try { - classLoader = (ClassLoader) getVFSClassLoader().getMethod("getClassLoader").invoke(null); - Thread.currentThread().setContextClassLoader(classLoader); - } catch (IOException | IllegalArgumentException | ReflectiveOperationException - | SecurityException e) { - die(e, "Problem initializing the class loader"); + public static ClassLoader getClassLoader() { + synchronized (lock) { + if (classLoader == null) { + try { + classLoader = (ClassLoader) getVFSClassLoader().getMethod("getClassLoader").invoke(null); + Thread.currentThread().setContextClassLoader(classLoader); + } catch (IOException | IllegalArgumentException | ReflectiveOperationException + | SecurityException e) { + die(e, "Problem initializing the class loader"); + } } + return classLoader; } - return classLoader; } @Deprecated - private static synchronized Class getVFSClassLoader() - throws IOException, ClassNotFoundException { - if (vfsClassLoader == null) { - Thread.currentThread().setContextClassLoader( - org.apache.accumulo.start.classloader.AccumuloClassLoader.getClassLoader()); - vfsClassLoader = org.apache.accumulo.start.classloader.AccumuloClassLoader.getClassLoader() - .loadClass("org.apache.accumulo.start.classloader.vfs.AccumuloVFSClassLoader"); + private static Class getVFSClassLoader() throws IOException, ClassNotFoundException { + synchronized (lock) { + if (vfsClassLoader == null) { + Thread.currentThread().setContextClassLoader( + org.apache.accumulo.start.classloader.AccumuloClassLoader.getClassLoader()); + vfsClassLoader = org.apache.accumulo.start.classloader.AccumuloClassLoader.getClassLoader() + .loadClass("org.apache.accumulo.start.classloader.vfs.AccumuloVFSClassLoader"); + } + return vfsClassLoader; } - return vfsClassLoader; } private static void execKeyword(final KeywordExecutable keywordExec, final String[] args) { @@ -226,11 +230,13 @@ public static void printUsage() { System.out.println(); } - public static synchronized Map getExecutables(final ClassLoader cl) { - if (servicesMap == null) { - servicesMap = checkDuplicates(ServiceLoader.load(KeywordExecutable.class, cl)); + public static Map getExecutables(final ClassLoader cl) { + synchronized (lock) { + if (servicesMap == null) { + servicesMap = checkDuplicates(ServiceLoader.load(KeywordExecutable.class, cl)); + } + return servicesMap; } - return servicesMap; } public static Map diff --git a/start/src/main/java/org/apache/accumulo/start/classloader/AccumuloClassLoader.java b/start/src/main/java/org/apache/accumulo/start/classloader/AccumuloClassLoader.java index 72c76eaa106..c124e2aa72f 100644 --- a/start/src/main/java/org/apache/accumulo/start/classloader/AccumuloClassLoader.java +++ b/start/src/main/java/org/apache/accumulo/start/classloader/AccumuloClassLoader.java @@ -48,6 +48,7 @@ public class AccumuloClassLoader { private static URL accumuloConfigUrl; private static URLClassLoader classloader; private static final Logger log = LoggerFactory.getLogger(AccumuloClassLoader.class); + private static final Object lock = new Object(); static { String configFile = System.getProperty("accumulo.properties", "accumulo.properties"); @@ -193,34 +194,36 @@ private static ArrayList findAccumuloURLs() throws IOException { return urls; } - public static synchronized ClassLoader getClassLoader() throws IOException { - if (classloader == null) { - ArrayList urls = findAccumuloURLs(); - - ClassLoader parentClassLoader = ClassLoader.getSystemClassLoader(); - - log.debug("Create 2nd tier ClassLoader using URLs: {}", urls); - classloader = - new URLClassLoader("AccumuloClassLoader (loads everything defined by general.classpaths)", - urls.toArray(new URL[urls.size()]), parentClassLoader) { - @Override - protected synchronized Class loadClass(String name, boolean resolve) - throws ClassNotFoundException { - - if (name.startsWith("org.apache.accumulo.start.classloader.vfs")) { - Class c = findLoadedClass(name); - if (c == null) { - try { - // try finding this class here instead of parent - findClass(name); - } catch (ClassNotFoundException e) {} - } + public static ClassLoader getClassLoader() throws IOException { + synchronized (lock) { + if (classloader == null) { + ArrayList urls = findAccumuloURLs(); + + ClassLoader parentClassLoader = ClassLoader.getSystemClassLoader(); + + log.debug("Create 2nd tier ClassLoader using URLs: {}", urls); + classloader = new URLClassLoader( + "AccumuloClassLoader (loads everything defined by general.classpaths)", + urls.toArray(new URL[urls.size()]), parentClassLoader) { + @Override + protected synchronized Class loadClass(String name, boolean resolve) + throws ClassNotFoundException { + + if (name.startsWith("org.apache.accumulo.start.classloader.vfs")) { + Class c = findLoadedClass(name); + if (c == null) { + try { + // try finding this class here instead of parent + findClass(name); + } catch (ClassNotFoundException e) {} } - return super.loadClass(name, resolve); } - }; - } + return super.loadClass(name, resolve); + } + }; + } - return classloader; + return classloader; + } } } diff --git a/start/src/main/java/org/apache/accumulo/start/classloader/vfs/AccumuloVFSClassLoader.java b/start/src/main/java/org/apache/accumulo/start/classloader/vfs/AccumuloVFSClassLoader.java index 942fffcb2c2..c9f51968464 100644 --- a/start/src/main/java/org/apache/accumulo/start/classloader/vfs/AccumuloVFSClassLoader.java +++ b/start/src/main/java/org/apache/accumulo/start/classloader/vfs/AccumuloVFSClassLoader.java @@ -432,24 +432,26 @@ public static ClassLoader getContextClassLoader(String contextName) throws IOExc return getContextManager().getClassLoader(contextName); } - public static synchronized ContextManager getContextManager() throws IOException { - if (contextManager == null) { - getClassLoader(); - try { - contextManager = new ContextManager(generateVfs(), () -> { - try { - return getClassLoader(); - } catch (IOException e) { - // throw runtime, then unwrap it. - throw new UncheckedIOException(e); - } - }); - } catch (UncheckedIOException uioe) { - throw uioe.getCause(); + public static ContextManager getContextManager() throws IOException { + synchronized (lock) { + if (contextManager == null) { + + getClassLoader(); + try { + contextManager = new ContextManager(generateVfs(), () -> { + try { + return getClassLoader(); + } catch (IOException e) { + // throw runtime, then unwrap it. + throw new UncheckedIOException(e); + } + }); + } catch (UncheckedIOException uioe) { + throw uioe.getCause(); + } } + return contextManager; } - - return contextManager; } public static void close() { diff --git a/start/src/main/spotbugs/exclude-filter.xml b/start/src/main/spotbugs/exclude-filter.xml index df55f1f613d..c5b24cd502c 100644 --- a/start/src/main/spotbugs/exclude-filter.xml +++ b/start/src/main/spotbugs/exclude-filter.xml @@ -27,8 +27,4 @@ - - - - From 66d14a0da8fef47c4e8611c3dd5a0443079dfd1e Mon Sep 17 00:00:00 2001 From: Daniel Roberts ddanielr Date: Wed, 22 Jul 2026 04:19:14 +0000 Subject: [PATCH 04/11] Suppress warnings for run statements Runnable requires a public run() method so these can't be fixed --- .../java/org/apache/accumulo/tserver/AssignmentHandler.java | 3 +++ .../java/org/apache/accumulo/tserver/TabletClientHandler.java | 2 ++ .../main/java/org/apache/accumulo/tserver/TabletServer.java | 2 ++ .../java/org/apache/accumulo/tserver/UnloadTabletHandler.java | 3 +++ 4 files changed, 10 insertions(+) diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/AssignmentHandler.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/AssignmentHandler.java index e6154e1392e..33f8d0e59a5 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/AssignmentHandler.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/AssignmentHandler.java @@ -49,6 +49,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; + class AssignmentHandler implements Runnable { private static final Logger log = LoggerFactory.getLogger(AssignmentHandler.class); private static final String METADATA_ISSUE = "Saw metadata issue when loading tablet : "; @@ -67,6 +69,7 @@ public AssignmentHandler(TabletServer server, KeyExtent extent, int retryAttempt } @Override + @SuppressFBWarnings("USO_UNSAFE_OBJECT_SYNCHRONIZATION") public void run() { synchronized (server.unopenedTablets) { synchronized (server.openingTablets) { diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/TabletClientHandler.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/TabletClientHandler.java index 608291db7af..c6d8da6a220 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/TabletClientHandler.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/TabletClientHandler.java @@ -133,6 +133,7 @@ import com.google.common.cache.Cache; +import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; import io.opentelemetry.api.trace.Span; import io.opentelemetry.context.Scope; @@ -1183,6 +1184,7 @@ static void checkPermission(ServerContext context, TabletHostingServer server, } @Override + @SuppressFBWarnings("USO_UNSAFE_OBJECT_SYNCHRONIZATION") public void loadTablet(TInfo tinfo, TCredentials credentials, String lock, final TKeyExtent textent) { diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/TabletServer.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/TabletServer.java index b6fc4695f4f..23a28dab426 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/TabletServer.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/TabletServer.java @@ -168,6 +168,7 @@ import com.google.common.annotations.VisibleForTesting; +import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; import io.opentelemetry.api.trace.Span; import io.opentelemetry.context.Scope; @@ -451,6 +452,7 @@ public MajorCompactor(ServerContext context) { } @Override + @SuppressFBWarnings("USO_UNSAFE_OBJECT_SYNCHRONIZATION") public void run() { while (true) { try { diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/UnloadTabletHandler.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/UnloadTabletHandler.java index 2000f40b99c..516fb364e8a 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/UnloadTabletHandler.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/UnloadTabletHandler.java @@ -35,6 +35,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; + class UnloadTabletHandler implements Runnable { private static final Logger log = LoggerFactory.getLogger(UnloadTabletHandler.class); private final KeyExtent extent; @@ -51,6 +53,7 @@ public UnloadTabletHandler(TabletServer server, KeyExtent extent, TUnloadTabletG } @Override + @SuppressFBWarnings("USO_UNSAFE_OBJECT_SYNCHRONIZATION") public void run() { Tablet t = null; From 0cb99ce7f746db1c2f0a867ca63b4c26b2d89636 Mon Sep 17 00:00:00 2001 From: Daniel Roberts ddanielr Date: Wed, 22 Jul 2026 04:21:31 +0000 Subject: [PATCH 05/11] Added private final locks for sync methods Fixed these spotbug findings --- .../manager/recovery/RecoveryManager.java | 3 ++- .../tserver/TabletServerResourceManager.java | 6 +++-- .../tserver/memory/NativeMapLoader.java | 27 ++++++++++--------- 3 files changed, 21 insertions(+), 15 deletions(-) diff --git a/server/manager/src/main/java/org/apache/accumulo/manager/recovery/RecoveryManager.java b/server/manager/src/main/java/org/apache/accumulo/manager/recovery/RecoveryManager.java index 666970cc09f..1ca47daf4b1 100644 --- a/server/manager/src/main/java/org/apache/accumulo/manager/recovery/RecoveryManager.java +++ b/server/manager/src/main/java/org/apache/accumulo/manager/recovery/RecoveryManager.java @@ -86,6 +86,7 @@ public RecoveryManager(Manager manager, long timeToCacheExistsInMillis) { } private class LogSortTask implements Runnable { + private final Object lock = new Object(); private final String source; private final String destination; private final String sortId; @@ -118,7 +119,7 @@ public void run() { log.warn("Failed to initiate log sort " + source, e); } finally { if (!rescheduled) { - synchronized (RecoveryManager.this) { + synchronized (lock) { closeTasksQueued.remove(sortId); } } diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/TabletServerResourceManager.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/TabletServerResourceManager.java index 7b3ec10e86a..005f3051d3e 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/TabletServerResourceManager.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/TabletServerResourceManager.java @@ -135,6 +135,7 @@ public class TabletServerResourceManager { private final ServerContext context; private final Cache fileLenCache; + private final Object tabletServerResourceManagerLock = new Object(); /** * This method creates a task that changes the number of core and maximum threads on the thread @@ -664,6 +665,7 @@ public class TabletResourceManager { private final KeyExtent extent; private final AccumuloConfiguration tableConf; + private final Object tabletResourceManagerLock = new Object(); TabletResourceManager(KeyExtent extent, AccumuloConfiguration tableConf) { requireNonNull(extent, "extent is null"); @@ -751,8 +753,8 @@ public void executeMinorCompaction(final Runnable r) { public void close() throws IOException { // always obtain locks in same order to avoid deadlock - synchronized (TabletServerResourceManager.this) { - synchronized (this) { + synchronized (tabletServerResourceManagerLock) { + synchronized (tabletResourceManagerLock) { if (closed) { throw new IOException("closed"); } diff --git a/server/tserver/src/main/java/org/apache/accumulo/tserver/memory/NativeMapLoader.java b/server/tserver/src/main/java/org/apache/accumulo/tserver/memory/NativeMapLoader.java index 14ff551895f..704757596b7 100644 --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/memory/NativeMapLoader.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/memory/NativeMapLoader.java @@ -37,23 +37,26 @@ public class NativeMapLoader { private static final Pattern dotSuffix = Pattern.compile("[.][^.]*$"); private static final String PROP_NAME = "accumulo.native.lib.path"; private static final AtomicBoolean loaded = new AtomicBoolean(false); + private static final Object lock = new Object(); // don't allow instantiation private NativeMapLoader() {} - public synchronized static void load() { - // load at most once; System.exit if loading fails - if (loaded.compareAndSet(false, true)) { - if (loadFromSearchPath(System.getProperty(PROP_NAME)) || loadFromSystemLinker()) { - return; + public static void load() { + synchronized (lock) { + // load at most once; System.exit if loading fails + if (loaded.compareAndSet(false, true)) { + if (loadFromSearchPath(System.getProperty(PROP_NAME)) || loadFromSystemLinker()) { + return; + } + log.error( + "FATAL! Accumulo native libraries were requested but could not" + + " be be loaded. Either set '{}' to false in accumulo.properties or make" + + " sure native libraries are created in directories set by the JVM" + + " system property '{}' in accumulo-env.sh!", + Property.TSERV_NATIVEMAP_ENABLED, PROP_NAME); + System.exit(1); } - log.error( - "FATAL! Accumulo native libraries were requested but could not" - + " be be loaded. Either set '{}' to false in accumulo.properties or make" - + " sure native libraries are created in directories set by the JVM" - + " system property '{}' in accumulo-env.sh!", - Property.TSERV_NATIVEMAP_ENABLED, PROP_NAME); - System.exit(1); } } From 1350f37ad10b588a20c30a27bbb456ed150771cb Mon Sep 17 00:00:00 2001 From: Daniel Roberts ddanielr Date: Wed, 22 Jul 2026 05:42:21 +0000 Subject: [PATCH 06/11] Add filter for ExternalDoNothingCompactor --- test/src/main/spotbugs/exclude-filter.xml | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/test/src/main/spotbugs/exclude-filter.xml b/test/src/main/spotbugs/exclude-filter.xml index ea7b5384050..3861f32eb19 100644 --- a/test/src/main/spotbugs/exclude-filter.xml +++ b/test/src/main/spotbugs/exclude-filter.xml @@ -35,4 +35,10 @@ + + + + + + From 92a13e3e205c27a0851d5005565397b72fe67ddf Mon Sep 17 00:00:00 2001 From: Daniel Roberts ddanielr Date: Wed, 22 Jul 2026 07:08:34 +0000 Subject: [PATCH 07/11] Add private locks to fix spotbug errors --- core/src/main/spotbugs/exclude-filter.xml | 10 +- .../accumulo/server/client/BulkImporter.java | 2 +- .../server/compaction/CompactionWatcher.java | 14 ++- .../accumulo/server/fs/FileManager.java | 3 +- .../server/problems/ProblemReports.java | 13 ++- .../server/replication/StatusUtil.java | 13 ++- .../server/util/ReplicationTableUtil.java | 95 ++++++++++--------- .../base/src/main/spotbugs/exclude-filter.xml | 32 +++++++ 8 files changed, 118 insertions(+), 64 deletions(-) create mode 100644 server/base/src/main/spotbugs/exclude-filter.xml diff --git a/core/src/main/spotbugs/exclude-filter.xml b/core/src/main/spotbugs/exclude-filter.xml index dd1d0d7b6d8..c503d6f344f 100644 --- a/core/src/main/spotbugs/exclude-filter.xml +++ b/core/src/main/spotbugs/exclude-filter.xml @@ -44,7 +44,13 @@ - - + + + + + + + + diff --git a/server/base/src/main/java/org/apache/accumulo/server/client/BulkImporter.java b/server/base/src/main/java/org/apache/accumulo/server/client/BulkImporter.java index 61457d907e0..ce5d8ee7546 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/client/BulkImporter.java +++ b/server/base/src/main/java/org/apache/accumulo/server/client/BulkImporter.java @@ -475,7 +475,7 @@ private Map> assignMapFiles(VolumeManager fs, } private class AssignmentTask implements Runnable { - final Map> assignmentFailures; + private final Map> assignmentFailures; HostAndPort location; private Map> assignmentsPerTablet; diff --git a/server/base/src/main/java/org/apache/accumulo/server/compaction/CompactionWatcher.java b/server/base/src/main/java/org/apache/accumulo/server/compaction/CompactionWatcher.java index 47e75fd42d4..a6134ed1a17 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/compaction/CompactionWatcher.java +++ b/server/base/src/main/java/org/apache/accumulo/server/compaction/CompactionWatcher.java @@ -45,6 +45,7 @@ public static void setTimer(LongTaskTimer ltt) { private final AccumuloConfiguration config; private static boolean watching = false; private static LongTaskTimer timer = null; + private static final Object lock = new Object(); private static class ObservedCompactionInfo { final CompactionInfo compactionInfo; @@ -121,11 +122,14 @@ public void run() { } } - public static synchronized void startWatching(ServerContext context) { - if (!watching) { - ThreadPools.watchCriticalScheduledTask(context.getScheduledExecutor().scheduleWithFixedDelay( - new CompactionWatcher(context.getConfiguration()), 10000, 10000, TimeUnit.MILLISECONDS)); - watching = true; + public static void startWatching(ServerContext context) { + synchronized (lock) { + if (!watching) { + ThreadPools.watchCriticalScheduledTask(context.getScheduledExecutor() + .scheduleWithFixedDelay(new CompactionWatcher(context.getConfiguration()), 10000, 10000, + TimeUnit.MILLISECONDS)); + watching = true; + } } } diff --git a/server/base/src/main/java/org/apache/accumulo/server/fs/FileManager.java b/server/base/src/main/java/org/apache/accumulo/server/fs/FileManager.java index 4e81899ed2f..ecfde27baf2 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/fs/FileManager.java +++ b/server/base/src/main/java/org/apache/accumulo/server/fs/FileManager.java @@ -111,6 +111,7 @@ public int hashCode() { private final long slowFilePermitMillis; private final ServerContext context; + private final Object lock = new Object(); private class IdleFileCloser implements Runnable { @@ -123,7 +124,7 @@ public void run() { // determine which files to close in a sync block, and then close the // files outside of the sync block - synchronized (FileManager.this) { + synchronized (lock) { Iterator>> iter = openFiles.entrySet().iterator(); while (iter.hasNext()) { Entry> entry = iter.next(); diff --git a/server/base/src/main/java/org/apache/accumulo/server/problems/ProblemReports.java b/server/base/src/main/java/org/apache/accumulo/server/problems/ProblemReports.java index c08d470d276..527454ecbfd 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/problems/ProblemReports.java +++ b/server/base/src/main/java/org/apache/accumulo/server/problems/ProblemReports.java @@ -59,6 +59,7 @@ public class ProblemReports implements Iterable { private static final Logger log = LoggerFactory.getLogger(ProblemReports.class); private final LRUMap problemReports = new LRUMap<>(1000); + private static final Object lock = new Object(); /** * use a thread pool so that reporting a problem never blocks @@ -283,12 +284,14 @@ public Iterator iterator() { return iterator(null); } - public static synchronized ProblemReports getInstance(ServerContext context) { - if (instance == null) { - instance = new ProblemReports(context); - } + public static ProblemReports getInstance(ServerContext context) { + synchronized (lock) { + if (instance == null) { + instance = new ProblemReports(context); + } - return instance; + return instance; + } } public static void main(String[] args) { diff --git a/server/base/src/main/java/org/apache/accumulo/server/replication/StatusUtil.java b/server/base/src/main/java/org/apache/accumulo/server/replication/StatusUtil.java index 7c40d0ef56b..081d5b34dea 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/replication/StatusUtil.java +++ b/server/base/src/main/java/org/apache/accumulo/server/replication/StatusUtil.java @@ -35,6 +35,7 @@ public class StatusUtil { private static final Value INF_END_REPLICATION_STATUS_VALUE, CLOSED_STATUS_VALUE; private static final Status.Builder CREATED_STATUS_BUILDER; + private static final Object lock = new Object(); static { CREATED_STATUS_BUILDER = Status.newBuilder(); @@ -119,11 +120,13 @@ public static Status replicatedAndIngested(Status.Builder builder, long recordsR /** * @return A {@link Status} for a new file that was just created */ - public static synchronized Status fileCreated(long timeCreated) { - // We're using a shared builder, so we need to synchronize access on it until we make a Status - // (which is then immutable) - CREATED_STATUS_BUILDER.setCreatedTime(timeCreated); - return CREATED_STATUS_BUILDER.build(); + public static Status fileCreated(long timeCreated) { + synchronized (lock) { + // We're using a shared builder, so we need to synchronize access on it until we make a Status + // (which is then immutable) + CREATED_STATUS_BUILDER.setCreatedTime(timeCreated); + return CREATED_STATUS_BUILDER.build(); + } } /** diff --git a/server/base/src/main/java/org/apache/accumulo/server/util/ReplicationTableUtil.java b/server/base/src/main/java/org/apache/accumulo/server/util/ReplicationTableUtil.java index 70f78451f8e..4450ba46b7f 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/util/ReplicationTableUtil.java +++ b/server/base/src/main/java/org/apache/accumulo/server/util/ReplicationTableUtil.java @@ -65,6 +65,8 @@ public class ReplicationTableUtil { public static final String STATUS_FORMATTER_CLASS_NAME = org.apache.accumulo.server.replication.StatusFormatter.class.getName(); + private static final Object lock = new Object(); + private ReplicationTableUtil() {} /** @@ -89,62 +91,65 @@ static synchronized Writer getWriter(ClientContext context) { return replicationTable; } - public synchronized static void configureMetadataTable(AccumuloClient client, String tableName) { - TableOperations tops = client.tableOperations(); - Map> iterators = null; - try { - iterators = tops.listIterators(tableName); - } catch (AccumuloSecurityException | AccumuloException | TableNotFoundException e) { - throw new RuntimeException(e); - } - - if (!iterators.containsKey(COMBINER_NAME)) { - // Set our combiner and combine all columns - // Need to set the combiner beneath versioning since we don't want to turn it off - @SuppressWarnings("deprecation") - var statusCombinerClass = org.apache.accumulo.server.replication.StatusCombiner.class; - IteratorSetting setting = new IteratorSetting(9, COMBINER_NAME, statusCombinerClass); - Combiner.setColumns(setting, Collections.singletonList(new Column(ReplicationSection.COLF))); + public static void configureMetadataTable(AccumuloClient client, String tableName) { + synchronized (lock) { + TableOperations tops = client.tableOperations(); + Map> iterators = null; try { - tops.attachIterator(tableName, setting); + iterators = tops.listIterators(tableName); } catch (AccumuloSecurityException | AccumuloException | TableNotFoundException e) { throw new RuntimeException(e); } - } - // Make sure the StatusFormatter is set on the metadata table - Map properties; - try { - properties = tops.getConfiguration(tableName); - } catch (AccumuloException | TableNotFoundException e) { - throw new RuntimeException(e); - } + if (!iterators.containsKey(COMBINER_NAME)) { + // Set our combiner and combine all columns + // Need to set the combiner beneath versioning since we don't want to turn it off + @SuppressWarnings("deprecation") + var statusCombinerClass = org.apache.accumulo.server.replication.StatusCombiner.class; + IteratorSetting setting = new IteratorSetting(9, COMBINER_NAME, statusCombinerClass); + Combiner.setColumns(setting, + Collections.singletonList(new Column(ReplicationSection.COLF))); + try { + tops.attachIterator(tableName, setting); + } catch (AccumuloSecurityException | AccumuloException | TableNotFoundException e) { + throw new RuntimeException(e); + } + } - for (Entry property : properties.entrySet()) { - if (Property.TABLE_FORMATTER_CLASS.getKey().equals(property.getKey())) { - if (!STATUS_FORMATTER_CLASS_NAME.equals(property.getValue())) { - log.info("Setting formatter for {} from {} to {}", tableName, property.getValue(), - STATUS_FORMATTER_CLASS_NAME); - try { - tops.setProperty(tableName, Property.TABLE_FORMATTER_CLASS.getKey(), + // Make sure the StatusFormatter is set on the metadata table + Map properties; + try { + properties = tops.getConfiguration(tableName); + } catch (AccumuloException | TableNotFoundException e) { + throw new RuntimeException(e); + } + + for (Entry property : properties.entrySet()) { + if (Property.TABLE_FORMATTER_CLASS.getKey().equals(property.getKey())) { + if (!STATUS_FORMATTER_CLASS_NAME.equals(property.getValue())) { + log.info("Setting formatter for {} from {} to {}", tableName, property.getValue(), STATUS_FORMATTER_CLASS_NAME); - } catch (AccumuloException | AccumuloSecurityException e) { - throw new RuntimeException(e); + try { + tops.setProperty(tableName, Property.TABLE_FORMATTER_CLASS.getKey(), + STATUS_FORMATTER_CLASS_NAME); + } catch (AccumuloException | AccumuloSecurityException e) { + throw new RuntimeException(e); + } } - } - // Don't need to keep iterating over the properties after we found the one we were looking - // for - return; + // Don't need to keep iterating over the properties after we found the one we were looking + // for + return; + } } - } - // Set the formatter on the table because it wasn't already there - try { - tops.setProperty(tableName, Property.TABLE_FORMATTER_CLASS.getKey(), - STATUS_FORMATTER_CLASS_NAME); - } catch (AccumuloException | AccumuloSecurityException e) { - throw new RuntimeException(e); + // Set the formatter on the table because it wasn't already there + try { + tops.setProperty(tableName, Property.TABLE_FORMATTER_CLASS.getKey(), + STATUS_FORMATTER_CLASS_NAME); + } catch (AccumuloException | AccumuloSecurityException e) { + throw new RuntimeException(e); + } } } diff --git a/server/base/src/main/spotbugs/exclude-filter.xml b/server/base/src/main/spotbugs/exclude-filter.xml new file mode 100644 index 00000000000..0b7dee2f0c0 --- /dev/null +++ b/server/base/src/main/spotbugs/exclude-filter.xml @@ -0,0 +1,32 @@ + + + + + + + + + + From 6ea6b9a2086156550a89fa5da4a7a3722383c28a Mon Sep 17 00:00:00 2001 From: Daniel Roberts ddanielr Date: Wed, 22 Jul 2026 07:10:04 +0000 Subject: [PATCH 08/11] Use zooCacheLock instead of synch on zooCache --- .../handler/KerberosAuthenticator.java | 5 +++-- .../security/handler/ZKAuthenticator.java | 9 +++++---- .../server/security/handler/ZKAuthorizor.java | 5 +++-- .../security/handler/ZKPermHandler.java | 19 ++++++++++--------- 4 files changed, 21 insertions(+), 17 deletions(-) diff --git a/server/base/src/main/java/org/apache/accumulo/server/security/handler/KerberosAuthenticator.java b/server/base/src/main/java/org/apache/accumulo/server/security/handler/KerberosAuthenticator.java index f998f148839..7d82beb1259 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/security/handler/KerberosAuthenticator.java +++ b/server/base/src/main/java/org/apache/accumulo/server/security/handler/KerberosAuthenticator.java @@ -56,6 +56,7 @@ public class KerberosAuthenticator implements Authenticator { private ServerContext context; private String zkUserPath; private UserImpersonation impersonation; + private final Object zooCacheLock = new Object(); @Override public void initialize(ServerContext context) { @@ -72,7 +73,7 @@ public boolean validSecurityHandlers() { } private void createUserNodeInZk(String principal) throws KeeperException, InterruptedException { - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(); ZooReaderWriter zoo = context.getZooReaderWriter(); zoo.putPrivatePersistentData(zkUserPath + "/" + principal, new byte[0], @@ -85,7 +86,7 @@ public void initializeSecurity(String principal, byte[] token) { try { // remove old settings from zookeeper first, if any ZooReaderWriter zoo = context.getZooReaderWriter(); - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(); if (zoo.exists(zkUserPath)) { zoo.recursiveDelete(zkUserPath, NodeMissingPolicy.SKIP); diff --git a/server/base/src/main/java/org/apache/accumulo/server/security/handler/ZKAuthenticator.java b/server/base/src/main/java/org/apache/accumulo/server/security/handler/ZKAuthenticator.java index bdd41135b31..e830917782a 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/security/handler/ZKAuthenticator.java +++ b/server/base/src/main/java/org/apache/accumulo/server/security/handler/ZKAuthenticator.java @@ -48,6 +48,7 @@ public final class ZKAuthenticator implements Authenticator { private ServerContext context; private String ZKUserPath; private ZooCache zooCache; + private final Object zooCacheLock = new Object(); @Override public void initialize(ServerContext context) { @@ -88,7 +89,7 @@ public void initializeSecurity(String principal, byte[] token) { try { // remove old settings from zookeeper first, if any ZooReaderWriter zoo = context.getZooReaderWriter(); - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(); if (zoo.exists(ZKUserPath)) { zoo.recursiveDelete(ZKUserPath, NodeMissingPolicy.SKIP); @@ -115,7 +116,7 @@ public void initializeSecurity(String principal, byte[] token) { */ private void constructUser(String user, byte[] pass) throws KeeperException, InterruptedException { - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(); ZooReaderWriter zoo = context.getZooReaderWriter(); zoo.putPrivatePersistentData(ZKUserPath + "/" + user, pass, NodeExistsPolicy.FAIL); @@ -154,7 +155,7 @@ public void createUser(String principal, AuthenticationToken token) @Override public void dropUser(String user) throws AccumuloSecurityException { try { - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(); context.getZooReaderWriter().recursiveDelete(ZKUserPath + "/" + user, NodeMissingPolicy.FAIL); @@ -181,7 +182,7 @@ public void changePassword(String principal, AuthenticationToken token) PasswordToken pt = (PasswordToken) token; if (userExists(principal)) { try { - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(ZKUserPath + "/" + principal); context.getZooReaderWriter().putPrivatePersistentData(ZKUserPath + "/" + principal, ZKSecurityTool.createPass(pt.getPassword()), NodeExistsPolicy.OVERWRITE); diff --git a/server/base/src/main/java/org/apache/accumulo/server/security/handler/ZKAuthorizor.java b/server/base/src/main/java/org/apache/accumulo/server/security/handler/ZKAuthorizor.java index 8f964e1d4b2..34cb5e3bb95 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/security/handler/ZKAuthorizor.java +++ b/server/base/src/main/java/org/apache/accumulo/server/security/handler/ZKAuthorizor.java @@ -45,6 +45,7 @@ public class ZKAuthorizor implements Authorizor { private ServerContext context; private String ZKUserPath; private ZooCache zooCache; + private final Object zooCacheLock = new Object(); @Override public void initialize(ServerContext context) { @@ -109,7 +110,7 @@ public void initUser(String user) throws AccumuloSecurityException { @Override public void dropUser(String user) throws AccumuloSecurityException { try { - synchronized (zooCache) { + synchronized (zooCacheLock) { ZooReaderWriter zoo = context.getZooReaderWriter(); zoo.recursiveDelete(ZKUserPath + "/" + user + ZKUserAuths, NodeMissingPolicy.SKIP); zooCache.clear(ZKUserPath + "/" + user); @@ -132,7 +133,7 @@ public void dropUser(String user) throws AccumuloSecurityException { public void changeAuthorizations(String user, Authorizations authorizations) throws AccumuloSecurityException { try { - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(); context.getZooReaderWriter().putPersistentData(ZKUserPath + "/" + user + ZKUserAuths, ZKSecurityTool.convertAuthorizations(authorizations), NodeExistsPolicy.OVERWRITE); diff --git a/server/base/src/main/java/org/apache/accumulo/server/security/handler/ZKPermHandler.java b/server/base/src/main/java/org/apache/accumulo/server/security/handler/ZKPermHandler.java index 96bb92079e2..f04b8431515 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/security/handler/ZKPermHandler.java +++ b/server/base/src/main/java/org/apache/accumulo/server/security/handler/ZKPermHandler.java @@ -59,6 +59,7 @@ public class ZKPermHandler implements PermissionHandler { private String ZKTablePath; private String ZKNamespacePath; private ZooCache zooCache; + private final Object zooCacheLock = new Object(); private final String ZKUserSysPerms = "/System"; private final String ZKUserTablePerms = "/Tables"; private final String ZKUserNamespacePerms = "/Namespaces"; @@ -191,7 +192,7 @@ public void grantSystemPermission(String user, SystemPermission permission) } if (perms.add(permission)) { - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(); zoo.putPersistentData(ZKUserPath + "/" + user + ZKUserSysPerms, ZKSecurityTool.convertSystemPermissions(perms), NodeExistsPolicy.OVERWRITE); @@ -220,7 +221,7 @@ public void grantTablePermission(String user, String table, TablePermission perm try { if (tablePerms.add(permission)) { - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(ZKUserPath + "/" + user + ZKUserTablePerms + "/" + table); zoo.putPersistentData(ZKUserPath + "/" + user + ZKUserTablePerms + "/" + table, ZKSecurityTool.convertTablePermissions(tablePerms), NodeExistsPolicy.OVERWRITE); @@ -250,7 +251,7 @@ public void grantNamespacePermission(String user, String namespace, try { if (namespacePerms.add(permission)) { - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(ZKUserPath + "/" + user + ZKUserNamespacePerms + "/" + namespace); zoo.putPersistentData(ZKUserPath + "/" + user + ZKUserNamespacePerms + "/" + namespace, ZKSecurityTool.convertNamespacePermissions(namespacePerms), @@ -281,7 +282,7 @@ public void revokeSystemPermission(String user, SystemPermission permission) try { if (sysPerms.remove(permission)) { - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(); zoo.putPersistentData(ZKUserPath + "/" + user + ZKUserSysPerms, ZKSecurityTool.convertSystemPermissions(sysPerms), NodeExistsPolicy.OVERWRITE); @@ -367,7 +368,7 @@ public void revokeNamespacePermission(String user, String namespace, @Override public void cleanTablePermissions(String table) throws AccumuloSecurityException { try { - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(); for (String user : zooCache.getChildren(ZKUserPath)) { zoo.recursiveDelete(ZKUserPath + "/" + user + ZKUserTablePerms + "/" + table, @@ -387,7 +388,7 @@ public void cleanTablePermissions(String table) throws AccumuloSecurityException @Override public void cleanNamespacePermissions(String namespace) throws AccumuloSecurityException { try { - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(); for (String user : zooCache.getChildren(ZKUserPath)) { zoo.recursiveDelete(ZKUserPath + "/" + user + ZKUserNamespacePerms + "/" + namespace, @@ -471,7 +472,7 @@ public void initUser(String user) throws AccumuloSecurityException { */ private void createTablePerm(String user, TableId table, Set perms) throws KeeperException, InterruptedException { - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(); zoo.putPersistentData(ZKUserPath + "/" + user + ZKUserTablePerms + "/" + table, ZKSecurityTool.convertTablePermissions(perms), NodeExistsPolicy.FAIL); @@ -484,7 +485,7 @@ private void createTablePerm(String user, TableId table, Set pe */ private void createNamespacePerm(String user, NamespaceId namespace, Set perms) throws KeeperException, InterruptedException { - synchronized (zooCache) { + synchronized (zooCacheLock) { zooCache.clear(); zoo.putPersistentData(ZKUserPath + "/" + user + ZKUserNamespacePerms + "/" + namespace, ZKSecurityTool.convertNamespacePermissions(perms), NodeExistsPolicy.FAIL); @@ -494,7 +495,7 @@ private void createNamespacePerm(String user, NamespaceId namespace, @Override public void cleanUser(String user) throws AccumuloSecurityException { try { - synchronized (zooCache) { + synchronized (zooCacheLock) { zoo.recursiveDelete(ZKUserPath + "/" + user + ZKUserSysPerms, NodeMissingPolicy.SKIP); zoo.recursiveDelete(ZKUserPath + "/" + user + ZKUserTablePerms, NodeMissingPolicy.SKIP); zoo.recursiveDelete(ZKUserPath + "/" + user + ZKUserNamespacePerms, NodeMissingPolicy.SKIP); From a75b59c84a3b84cca57242107f06fcd75ea11cb5 Mon Sep 17 00:00:00 2001 From: Daniel Roberts ddanielr Date: Wed, 22 Jul 2026 07:10:36 +0000 Subject: [PATCH 09/11] Moved lock outside of method call The intent of this code is not clear. This may not be the correct way to do this but will review it at the PR stage. --- .../accumulo/server/zookeeper/DistributedWorkQueue.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/server/base/src/main/java/org/apache/accumulo/server/zookeeper/DistributedWorkQueue.java b/server/base/src/main/java/org/apache/accumulo/server/zookeeper/DistributedWorkQueue.java index d16948aed80..ced88dd1c3e 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/zookeeper/DistributedWorkQueue.java +++ b/server/base/src/main/java/org/apache/accumulo/server/zookeeper/DistributedWorkQueue.java @@ -278,9 +278,9 @@ public List getWorkQueued() throws KeeperException, InterruptedException return children; } - public void waitUntilDone(Set workIDs) throws KeeperException, InterruptedException { + private final Object condVar = new Object(); - final Object condVar = new Object(); + public void waitUntilDone(Set workIDs) throws KeeperException, InterruptedException { Watcher watcher = new Watcher() { @SuppressFBWarnings(value = "NN_NAKED_NOTIFY", From 6f6f2943eb9904364145181794ca0fc1c1fc23ad Mon Sep 17 00:00:00 2001 From: Daniel Roberts ddanielr Date: Wed, 22 Jul 2026 07:28:27 +0000 Subject: [PATCH 10/11] Suppress warnings in test cases that pass null --- .../base/src/main/spotbugs/exclude-filter.xml | 32 ------------------- .../delegation/AuthenticationKeyTest.java | 4 +++ .../server/zookeeper/ZooAclUtilTest.java | 4 +++ 3 files changed, 8 insertions(+), 32 deletions(-) delete mode 100644 server/base/src/main/spotbugs/exclude-filter.xml diff --git a/server/base/src/main/spotbugs/exclude-filter.xml b/server/base/src/main/spotbugs/exclude-filter.xml deleted file mode 100644 index 0b7dee2f0c0..00000000000 --- a/server/base/src/main/spotbugs/exclude-filter.xml +++ /dev/null @@ -1,32 +0,0 @@ - - - - - - - - - - diff --git a/server/base/src/test/java/org/apache/accumulo/server/security/delegation/AuthenticationKeyTest.java b/server/base/src/test/java/org/apache/accumulo/server/security/delegation/AuthenticationKeyTest.java index c8d9e175481..f174729ffb5 100644 --- a/server/base/src/test/java/org/apache/accumulo/server/security/delegation/AuthenticationKeyTest.java +++ b/server/base/src/test/java/org/apache/accumulo/server/security/delegation/AuthenticationKeyTest.java @@ -34,6 +34,10 @@ import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; +import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; + +@SuppressFBWarnings(value = "NP_NULL_PARAM_DEREF_NONVIRTUAL", + justification = "Test code is testing that a NPE is thrown") public class AuthenticationKeyTest { // From org.apache.hadoop.security.token.SecretManager private static final String DEFAULT_HMAC_ALGORITHM = "HmacSHA1"; diff --git a/server/base/src/test/java/org/apache/accumulo/server/zookeeper/ZooAclUtilTest.java b/server/base/src/test/java/org/apache/accumulo/server/zookeeper/ZooAclUtilTest.java index 31d86a3056e..3b7dbceeb09 100644 --- a/server/base/src/test/java/org/apache/accumulo/server/zookeeper/ZooAclUtilTest.java +++ b/server/base/src/test/java/org/apache/accumulo/server/zookeeper/ZooAclUtilTest.java @@ -30,6 +30,10 @@ import org.apache.zookeeper.data.Id; import org.junit.jupiter.api.Test; +import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; + +@SuppressFBWarnings(value = "NP_NULL_PARAM_DEREF_NONVIRTUAL", + justification = "Test code is testing that a NPE is thrown") public class ZooAclUtilTest { @Test From 4a1931207d4ec3519437f26f8c3f5cd45c3e51f7 Mon Sep 17 00:00:00 2001 From: Daniel Roberts ddanielr Date: Mon, 10 Aug 2026 19:51:44 +0000 Subject: [PATCH 11/11] Address the rest of the lock spotbugs findings --- .../core/classloader/ClassLoaderUtil.java | 53 ++++++----- .../accumulo/core/clientImpl/ScannerImpl.java | 5 +- .../core/clientImpl/TabletLocator.java | 92 +++++++++++-------- .../TabletServerBatchReaderIterator.java | 3 +- .../core/fate/zookeeper/ServiceLock.java | 10 +- .../cache/impl/BlockCacheManagerFactory.java | 43 +++++---- .../accumulo/core/file/rfile/BlockIndex.java | 22 +++-- .../core/file/rfile/bcfile/BCFile.java | 3 +- .../core/metadata/schema/TabletMetadata.java | 22 +++-- .../apache/accumulo/core/rpc/ThriftUtil.java | 23 +++-- .../core/singletons/SingletonManager.java | 87 ++++++++++-------- .../ratelimit/SharedRateLimiterFactory.java | 26 +++--- core/src/main/spotbugs/exclude-filter.xml | 7 +- .../DistributedReadWriteLockTest.java | 6 +- 14 files changed, 224 insertions(+), 178 deletions(-) diff --git a/core/src/main/java/org/apache/accumulo/core/classloader/ClassLoaderUtil.java b/core/src/main/java/org/apache/accumulo/core/classloader/ClassLoaderUtil.java index 31e69e0a20e..6cdfed2cc4e 100644 --- a/core/src/main/java/org/apache/accumulo/core/classloader/ClassLoaderUtil.java +++ b/core/src/main/java/org/apache/accumulo/core/classloader/ClassLoaderUtil.java @@ -32,6 +32,7 @@ public class ClassLoaderUtil { private static final Logger LOG = LoggerFactory.getLogger(ClassLoaderUtil.class); private static ContextClassLoaderFactory FACTORY; + private static final Object FACTORY_LOCK = new Object(); private ClassLoaderUtil() { // cannot construct; static utilities only @@ -40,29 +41,33 @@ private ClassLoaderUtil() { /** * Initialize the ContextClassLoaderFactory */ - public static synchronized void initContextFactory(AccumuloConfiguration conf) { - if (FACTORY == null) { - LOG.debug("Creating {}", ContextClassLoaderFactory.class.getName()); - String factoryName = conf.get(Property.GENERAL_CONTEXT_CLASSLOADER_FACTORY); - if (factoryName == null || factoryName.isEmpty()) { - // load the default implementation - LOG.info("Using default {}, which is subject to change in a future release", - ContextClassLoaderFactory.class.getName()); - FACTORY = new DefaultContextClassLoaderFactory(conf); - } else { - // load user's selected implementation and provide it with the service environment - try { - var factoryClass = Class.forName(factoryName).asSubclass(ContextClassLoaderFactory.class); - LOG.info("Creating {}: {}", ContextClassLoaderFactory.class.getName(), factoryName); - FACTORY = factoryClass.getDeclaredConstructor().newInstance(); - FACTORY.init(() -> new ConfigurationImpl(conf)); - } catch (ReflectiveOperationException e) { - throw new IllegalStateException("Unable to load and initialize class: " + factoryName, e); + public static void initContextFactory(AccumuloConfiguration conf) { + synchronized (FACTORY_LOCK) { + if (FACTORY == null) { + LOG.debug("Creating {}", ContextClassLoaderFactory.class.getName()); + String factoryName = conf.get(Property.GENERAL_CONTEXT_CLASSLOADER_FACTORY); + if (factoryName == null || factoryName.isEmpty()) { + // load the default implementation + LOG.info("Using default {}, which is subject to change in a future release", + ContextClassLoaderFactory.class.getName()); + FACTORY = new DefaultContextClassLoaderFactory(conf); + } else { + // load user's selected implementation and provide it with the service environment + try { + var factoryClass = + Class.forName(factoryName).asSubclass(ContextClassLoaderFactory.class); + LOG.info("Creating {}: {}", ContextClassLoaderFactory.class.getName(), factoryName); + FACTORY = factoryClass.getDeclaredConstructor().newInstance(); + FACTORY.init(() -> new ConfigurationImpl(conf)); + } catch (ReflectiveOperationException e) { + throw new IllegalStateException("Unable to load and initialize class: " + factoryName, + e); + } } + } else { + LOG.debug("{} already initialized with {}.", ContextClassLoaderFactory.class.getName(), + FACTORY.getClass().getName()); } - } else { - LOG.debug("{} already initialized with {}.", ContextClassLoaderFactory.class.getName(), - FACTORY.getClass().getName()); } } @@ -72,8 +77,10 @@ static ContextClassLoaderFactory getContextFactory() { } // for testing - public static synchronized void resetContextFactoryForTests() { - FACTORY = null; + public static void resetContextFactoryForTests() { + synchronized (FACTORY_LOCK) { + FACTORY = null; + } } @SuppressWarnings("deprecation") diff --git a/core/src/main/java/org/apache/accumulo/core/clientImpl/ScannerImpl.java b/core/src/main/java/org/apache/accumulo/core/clientImpl/ScannerImpl.java index 0fef7ecaab0..4e08b3c374a 100644 --- a/core/src/main/java/org/apache/accumulo/core/clientImpl/ScannerImpl.java +++ b/core/src/main/java/org/apache/accumulo/core/clientImpl/ScannerImpl.java @@ -82,6 +82,7 @@ public boolean removeEldestEntry(Map.Entry eldest) { return size() > MAX_ENTRIES; } }; + private final Object scannerLock = new Object(); /** * This is used for ScannerIterators to report their activity back to the scanner that created @@ -90,7 +91,7 @@ public boolean removeEldestEntry(Map.Entry eldest) { class Reporter { void readBatch(ScannerIterator iter) { - synchronized (ScannerImpl.this) { + synchronized (scannerLock) { // This iter just had some activity, so access it in map so it becomes the most recently // used. iters.get(iter); @@ -98,7 +99,7 @@ void readBatch(ScannerIterator iter) { } void finished(ScannerIterator iter) { - synchronized (ScannerImpl.this) { + synchronized (scannerLock) { iters.remove(iter); } } diff --git a/core/src/main/java/org/apache/accumulo/core/clientImpl/TabletLocator.java b/core/src/main/java/org/apache/accumulo/core/clientImpl/TabletLocator.java index 6375b9cf841..ac40017664a 100644 --- a/core/src/main/java/org/apache/accumulo/core/clientImpl/TabletLocator.java +++ b/core/src/main/java/org/apache/accumulo/core/clientImpl/TabletLocator.java @@ -60,6 +60,8 @@ boolean isValid() { return isValid; } + private static final Object lock = new Object(); + public abstract TabletLocation locateTablet(ClientContext context, Text row, boolean skipRow, boolean retry) throws AccumuloException, AccumuloSecurityException, TableNotFoundException; @@ -118,12 +120,14 @@ public boolean equals(LocatorKey lk) { new HashMap<>(); private static boolean enabled = true; - public static synchronized void clearLocators() { - for (TabletLocator locator : locators.values()) { - locator.isValid = false; + public static void clearLocators() { + synchronized (lock) { + for (TabletLocator locator : locators.values()) { + locator.isValid = false; + } + locators.clear(); + offlineLocators.clear(); } - locators.clear(); - offlineLocators.clear(); } static synchronized boolean isEnabled() { @@ -139,36 +143,39 @@ static synchronized void enable() { enabled = true; } - public static synchronized TabletLocator getLocator(ClientContext context, TableId tableId) { - Preconditions.checkState(enabled, "The Accumulo singleton that that tracks tablet locations is " - + "disabled. This is likely caused by all AccumuloClients being closed or garbage collected"); - - clearUnusedTables(context); - - TableState state = context.getTableState(tableId); - LocatorKey key = new LocatorKey(context.getInstanceID(), tableId); - if (state == TableState.OFFLINE) { - locators.remove(key); - return offlineLocators.computeIfAbsent(key, - f -> new OfflineTabletLocatorImpl(context, tableId)); - } else { - offlineLocators.remove(key); - TabletLocator tl = locators.get(key); - if (tl == null) { - MetadataLocationObtainer mlo = new MetadataLocationObtainer(); - - if (RootTable.ID.equals(tableId)) { - tl = new RootTabletLocator(context.getTServerLockChecker()); - } else if (MetadataTable.ID.equals(tableId)) { - tl = new TabletLocatorImpl(MetadataTable.ID, getLocator(context, RootTable.ID), mlo, - context.getTServerLockChecker()); - } else { - tl = new TabletLocatorImpl(tableId, getLocator(context, MetadataTable.ID), mlo, - context.getTServerLockChecker()); + public static TabletLocator getLocator(ClientContext context, TableId tableId) { + synchronized (lock) { + Preconditions.checkState(enabled, + "The Accumulo singleton that that tracks tablet locations is " + + "disabled. This is likely caused by all AccumuloClients being closed or garbage collected"); + + clearUnusedTables(context); + + TableState state = context.getTableState(tableId); + LocatorKey key = new LocatorKey(context.getInstanceID(), tableId); + if (state == TableState.OFFLINE) { + locators.remove(key); + return offlineLocators.computeIfAbsent(key, + f -> new OfflineTabletLocatorImpl(context, tableId)); + } else { + offlineLocators.remove(key); + TabletLocator tl = locators.get(key); + if (tl == null) { + MetadataLocationObtainer mlo = new MetadataLocationObtainer(); + + if (RootTable.ID.equals(tableId)) { + tl = new RootTabletLocator(context.getTServerLockChecker()); + } else if (MetadataTable.ID.equals(tableId)) { + tl = new TabletLocatorImpl(MetadataTable.ID, getLocator(context, RootTable.ID), mlo, + context.getTServerLockChecker()); + } else { + tl = new TabletLocatorImpl(tableId, getLocator(context, MetadataTable.ID), mlo, + context.getTServerLockChecker()); + } + locators.put(key, tl); } - locators.put(key, tl); + return tl; } - return tl; } } @@ -177,9 +184,11 @@ public static synchronized TabletLocator getLocator(ClientContext context, Table * Checks if a table id is present in the cache w/o creating it. */ @VisibleForTesting - public static synchronized boolean isPresent(ClientContext context, TableId tableId) { - LocatorKey key = new LocatorKey(context.getInstanceID(), tableId); - return locators.containsKey(key) || offlineLocators.containsKey(key); + public static boolean isPresent(ClientContext context, TableId tableId) { + synchronized (lock) { + LocatorKey key = new LocatorKey(context.getInstanceID(), tableId); + return locators.containsKey(key) || offlineLocators.containsKey(key); + } } private static Duration clearFrequency = Duration.ofMinutes(10); @@ -188,10 +197,13 @@ public static synchronized boolean isPresent(ClientContext context, TableId tabl * Sets how often checks for unused tables are done */ @VisibleForTesting - public static synchronized void setClearFrequency(Duration frequency) { - Preconditions.checkArgument(frequency != null && !frequency.isNegative() && !frequency.isZero(), - "frequency:%s", frequency); - clearFrequency = frequency; + public static void setClearFrequency(Duration frequency) { + synchronized (lock) { + Preconditions.checkArgument( + frequency != null && !frequency.isNegative() && !frequency.isZero(), "frequency:%s", + frequency); + clearFrequency = frequency; + } } private static final Timer lastClearTimer = Timer.startNew(); diff --git a/core/src/main/java/org/apache/accumulo/core/clientImpl/TabletServerBatchReaderIterator.java b/core/src/main/java/org/apache/accumulo/core/clientImpl/TabletServerBatchReaderIterator.java index ccd5ec6e089..39d3c345448 100644 --- a/core/src/main/java/org/apache/accumulo/core/clientImpl/TabletServerBatchReaderIterator.java +++ b/core/src/main/java/org/apache/accumulo/core/clientImpl/TabletServerBatchReaderIterator.java @@ -740,6 +740,7 @@ private static class TimeoutTracker { final String server; final Set badServers; final long timeOut; + private final Object lock = new Object(); // When failures happen, rpc task to scan a server may be requeued in a thread pool. These two // variables track failures across task running in those thread pools. @@ -771,7 +772,7 @@ void check() throws IOException { void madeProgress() { activityTime = System.currentTimeMillis(); - synchronized (TimeoutTracker.this) { + synchronized (lock) { firstErrorTime = null; firstAllFailureTime = null; } diff --git a/core/src/main/java/org/apache/accumulo/core/fate/zookeeper/ServiceLock.java b/core/src/main/java/org/apache/accumulo/core/fate/zookeeper/ServiceLock.java index 95ff12ce0ce..d03697cf0fd 100644 --- a/core/src/main/java/org/apache/accumulo/core/fate/zookeeper/ServiceLock.java +++ b/core/src/main/java/org/apache/accumulo/core/fate/zookeeper/ServiceLock.java @@ -69,6 +69,8 @@ public class ServiceLock implements Watcher { private static final String ZLOCK_PREFIX = "zlock#"; + private static final Object lock = new Object(); + private static class Prefix { private final String prefix; @@ -370,7 +372,7 @@ public void process(WatchedEvent event) { if (event.getType() == EventType.NodeDeleted && event.getPath().equals(nodeToWatch)) { LOG.debug("[{}] Detected deletion of prior node {}, attempting to acquire lock; {}", vmLockPrefix, nodeToWatch, event); - synchronized (ServiceLock.this) { + synchronized (lock) { try { if (createdNodeName != null) { determineLockOwnership(lw); @@ -390,7 +392,7 @@ public void process(WatchedEvent event) { if (event.getState() == KeeperState.Expired || event.getState() == KeeperState.Disconnected) { - synchronized (ServiceLock.this) { + synchronized (lock) { if (lockNodeName == null) { LOG.info("Zookeeper Session expired / disconnected; {}", event); lw.failedToAcquireLock( @@ -400,7 +402,7 @@ public void process(WatchedEvent event) { renew = false; } if (renew) { - synchronized (ServiceLock.this) { + synchronized (lock) { if (createdNodeName != null) { try { Stat restat = zooKeeper.exists(nodeToWatch, this); @@ -524,7 +526,7 @@ private void failedToAcquireLock() { @Override public void process(WatchedEvent event) { - synchronized (ServiceLock.this) { + synchronized (lock) { if (lockNodeName != null && event.getType() == EventType.NodeDeleted && event.getPath().equals(path + "/" + lockNodeName)) { LOG.debug("[{}] {} was deleted; {}", vmLockPrefix, lockNodeName, event); diff --git a/core/src/main/java/org/apache/accumulo/core/file/blockfile/cache/impl/BlockCacheManagerFactory.java b/core/src/main/java/org/apache/accumulo/core/file/blockfile/cache/impl/BlockCacheManagerFactory.java index d62a5e8d468..5fbd54c225c 100644 --- a/core/src/main/java/org/apache/accumulo/core/file/blockfile/cache/impl/BlockCacheManagerFactory.java +++ b/core/src/main/java/org/apache/accumulo/core/file/blockfile/cache/impl/BlockCacheManagerFactory.java @@ -28,6 +28,7 @@ public class BlockCacheManagerFactory { private static final Logger LOG = LoggerFactory.getLogger(BlockCacheManager.class); + private static final Object lock = new Object(); /** * Get the BlockCacheFactory specified by the property 'tserver.cache.factory.class' using the @@ -37,16 +38,17 @@ public class BlockCacheManagerFactory { * @return block cache manager instance * @throws Exception error loading block cache manager implementation class */ - public static synchronized BlockCacheManager getInstance(AccumuloConfiguration conf) - throws Exception { - @SuppressWarnings("deprecation") - var cacheManagerProp = - conf.resolve(Property.GENERAL_CACHE_MANAGER_IMPL, Property.TSERV_CACHE_MANAGER_IMPL); - String impl = conf.get(cacheManagerProp); - Class clazz = - ClassLoaderUtil.loadClass(impl, BlockCacheManager.class); - LOG.info("Created new block cache manager of type: {}", clazz.getSimpleName()); - return clazz.getDeclaredConstructor().newInstance(); + public static BlockCacheManager getInstance(AccumuloConfiguration conf) throws Exception { + synchronized (lock) { + @SuppressWarnings("deprecation") + var cacheManagerProp = + conf.resolve(Property.GENERAL_CACHE_MANAGER_IMPL, Property.TSERV_CACHE_MANAGER_IMPL); + String impl = conf.get(cacheManagerProp); + Class clazz = + ClassLoaderUtil.loadClass(impl, BlockCacheManager.class); + LOG.info("Created new block cache manager of type: {}", clazz.getSimpleName()); + return clazz.getDeclaredConstructor().newInstance(); + } } /** @@ -56,15 +58,16 @@ public static synchronized BlockCacheManager getInstance(AccumuloConfiguration c * @return block cache manager instance * @throws Exception error loading block cache manager implementation class */ - public static synchronized BlockCacheManager getClientInstance(AccumuloConfiguration conf) - throws Exception { - @SuppressWarnings("deprecation") - var cacheManagerProp = - conf.resolve(Property.GENERAL_CACHE_MANAGER_IMPL, Property.TSERV_CACHE_MANAGER_IMPL); - String impl = conf.get(cacheManagerProp); - Class clazz = - Class.forName(impl).asSubclass(BlockCacheManager.class); - LOG.info("Created new block cache factory of type: {}", clazz.getSimpleName()); - return clazz.getDeclaredConstructor().newInstance(); + public static BlockCacheManager getClientInstance(AccumuloConfiguration conf) throws Exception { + synchronized (lock) { + @SuppressWarnings("deprecation") + var cacheManagerProp = + conf.resolve(Property.GENERAL_CACHE_MANAGER_IMPL, Property.TSERV_CACHE_MANAGER_IMPL); + String impl = conf.get(cacheManagerProp); + Class clazz = + Class.forName(impl).asSubclass(BlockCacheManager.class); + LOG.info("Created new block cache factory of type: {}", clazz.getSimpleName()); + return clazz.getDeclaredConstructor().newInstance(); + } } } diff --git a/core/src/main/java/org/apache/accumulo/core/file/rfile/BlockIndex.java b/core/src/main/java/org/apache/accumulo/core/file/rfile/BlockIndex.java index a2ae4049d10..1d685f51dca 100644 --- a/core/src/main/java/org/apache/accumulo/core/file/rfile/BlockIndex.java +++ b/core/src/main/java/org/apache/accumulo/core/file/rfile/BlockIndex.java @@ -33,6 +33,8 @@ public class BlockIndex implements Weighable { + private final Object lock = new Object(); + private BlockIndex() {} public static BlockIndex getIndex(CachedBlockRead cacheBlock, IndexEntry indexEntry) @@ -214,16 +216,18 @@ BlockIndexEntry[] getIndexEntries() { } @Override - public synchronized int weight() { - int weight = 0; - if (blockIndex != null) { - for (BlockIndexEntry blockIndexEntry : blockIndex) { - weight += blockIndexEntry.weight(); + public int weight() { + synchronized (lock) { + int weight = 0; + if (blockIndex != null) { + for (BlockIndexEntry blockIndexEntry : blockIndex) { + weight += blockIndexEntry.weight(); + } } - } - weight += - ClassSize.ATOMIC_INTEGER + ClassSize.OBJECT + 2 * ClassSize.REFERENCE + ClassSize.ARRAY; - return weight; + weight += + ClassSize.ATOMIC_INTEGER + ClassSize.OBJECT + 2 * ClassSize.REFERENCE + ClassSize.ARRAY; + return weight; + } } } diff --git a/core/src/main/java/org/apache/accumulo/core/file/rfile/bcfile/BCFile.java b/core/src/main/java/org/apache/accumulo/core/file/rfile/bcfile/BCFile.java index 31f32d64ae1..15f65542b41 100644 --- a/core/src/main/java/org/apache/accumulo/core/file/rfile/bcfile/BCFile.java +++ b/core/src/main/java/org/apache/accumulo/core/file/rfile/bcfile/BCFile.java @@ -473,6 +473,7 @@ private static final class RBlockState { private final InputStream rawInputStream; private final InputStream in; private volatile boolean closed; + private final Object closeLock = new Object(); public RBlockState( CompressionAlgorithm compressionAlgo, InputStreamType fsin, BlockRegion region, @@ -523,7 +524,7 @@ public void flushStats() { } public void finish() throws IOException { - synchronized (in) { + synchronized (closeLock) { if (!closed) { try { in.close(); diff --git a/core/src/main/java/org/apache/accumulo/core/metadata/schema/TabletMetadata.java b/core/src/main/java/org/apache/accumulo/core/metadata/schema/TabletMetadata.java index dd11a95df84..b26bbd1698f 100644 --- a/core/src/main/java/org/apache/accumulo/core/metadata/schema/TabletMetadata.java +++ b/core/src/main/java/org/apache/accumulo/core/metadata/schema/TabletMetadata.java @@ -133,6 +133,8 @@ public enum ColumnType { ECOMP } + private static final Object lock = new Object(); + public static class Location { private final TServerInstance tServerInstance; private final LocationType lt; @@ -534,18 +536,20 @@ public static TabletMetadata create(String id, String prevEndRow, String endRow) * Get the tservers that are live from ZK. Live servers will have a valid ZooLock. This method was * pulled from org.apache.accumulo.server.manager.LiveTServerSet */ - public static synchronized Set getLiveTServers(ClientContext context) { + public static Set getLiveTServers(ClientContext context) { + synchronized (lock) { - final String path = context.getZooKeeperRoot() + Constants.ZTSERVERS; - final List children = context.getZooCache().getChildren(path); - final Set liveServers = new HashSet<>(children.size()); + final String path = context.getZooKeeperRoot() + Constants.ZTSERVERS; + final List children = context.getZooCache().getChildren(path); + final Set liveServers = new HashSet<>(children.size()); - for (String child : children) { - checkServer(context, path, child).ifPresent(liveServers::add); - } - log.trace("Found {} live tservers at ZK path: {}", liveServers.size(), path); + for (String child : children) { + checkServer(context, path, child).ifPresent(liveServers::add); + } + log.trace("Found {} live tservers at ZK path: {}", liveServers.size(), path); - return liveServers; + return liveServers; + } } /** diff --git a/core/src/main/java/org/apache/accumulo/core/rpc/ThriftUtil.java b/core/src/main/java/org/apache/accumulo/core/rpc/ThriftUtil.java index 41c4650dd30..554845bda7c 100644 --- a/core/src/main/java/org/apache/accumulo/core/rpc/ThriftUtil.java +++ b/core/src/main/java/org/apache/accumulo/core/rpc/ThriftUtil.java @@ -60,6 +60,7 @@ public class ThriftUtil { private static final SecureRandom random = new SecureRandom(); private static final int RELOGIN_MAX_BACKOFF = 5000; + private static final Object lock = new Object(); /** * An instance of {@link TraceProtocolFactory} @@ -169,17 +170,19 @@ public static TTransport createTransport(HostAndPort address, ClientContext cont * @param maxFrameSize Maximum Thrift message frame size * @return A, possibly cached, TTransportFactory with the requested maximum frame size */ - public static synchronized TTransportFactory transportFactory(long maxFrameSize) { - if (maxFrameSize > Integer.MAX_VALUE || maxFrameSize < 1) { - throw new RuntimeException("Thrift transport frames are limited to " + Integer.MAX_VALUE); - } - int maxFrameSize1 = (int) maxFrameSize; - TTransportFactory factory = factoryCache.get(maxFrameSize1); - if (factory == null) { - factory = new AccumuloTFramedTransportFactory(maxFrameSize1); - factoryCache.put(maxFrameSize1, factory); + public static TTransportFactory transportFactory(long maxFrameSize) { + synchronized (lock) { + if (maxFrameSize > Integer.MAX_VALUE || maxFrameSize < 1) { + throw new RuntimeException("Thrift transport frames are limited to " + Integer.MAX_VALUE); + } + int maxFrameSize1 = (int) maxFrameSize; + TTransportFactory factory = factoryCache.get(maxFrameSize1); + if (factory == null) { + factory = new AccumuloTFramedTransportFactory(maxFrameSize1); + factoryCache.put(maxFrameSize1, factory); + } + return factory; } - return factory; } /** diff --git a/core/src/main/java/org/apache/accumulo/core/singletons/SingletonManager.java b/core/src/main/java/org/apache/accumulo/core/singletons/SingletonManager.java index 5f3e151b58f..b34139e074d 100644 --- a/core/src/main/java/org/apache/accumulo/core/singletons/SingletonManager.java +++ b/core/src/main/java/org/apache/accumulo/core/singletons/SingletonManager.java @@ -46,6 +46,7 @@ public class SingletonManager { private static final Logger log = LoggerFactory.getLogger(SingletonManager.class); + private static final Object lock = new Object(); /** * These enums determine the behavior of the SingletonManager. @@ -112,16 +113,18 @@ private static void disable(SingletonService service) { /** * Register a static singleton that should be disabled and enabled as needed. */ - public static synchronized void register(SingletonService service) { - if (enabled && !service.isEnabled()) { - enable(service); - } + public static void register(SingletonService service) { + synchronized (lock) { + if (enabled && !service.isEnabled()) { + enable(service); + } - if (!enabled && service.isEnabled()) { - disable(service); - } + if (!enabled && service.isEnabled()) { + disable(service); + } - services.add(service); + services.add(service); + } } /** @@ -132,11 +135,13 @@ public static synchronized void register(SingletonService service) { * * @return A reservation that must be closed when the AccumuloClient is closed. */ - public static synchronized SingletonReservation getClientReservation() { - Preconditions.checkState(reservations >= 0); - reservations++; - transition(); - return new SingletonReservation(); + public static SingletonReservation getClientReservation() { + synchronized (lock) { + Preconditions.checkState(reservations >= 0); + reservations++; + transition(); + return new SingletonReservation(); + } } static synchronized void releaseReservation() { @@ -153,38 +158,42 @@ public static long getReservationCount() { /** * Change how singletons are managed. The default mode is {@link Mode#CLIENT} */ - public static synchronized void setMode(Mode mode) { - if (SingletonManager.mode == mode) { - return; - } - if (SingletonManager.mode == Mode.CLOSED) { - throw new IllegalStateException("Cannot leave closed mode once entered"); - } - if (SingletonManager.mode == Mode.CLIENT && mode == Mode.CONNECTOR) { - if (transitionedFromClientToConnector) { - throw new IllegalStateException("Can only transition from " + Mode.CLIENT + " to " - + Mode.CONNECTOR + " once. This error indicates that " - + "org.apache.accumulo.core.util.CleanUp.shutdownNow() was called and then later a " - + "Connector was created. Connectors can not be created after CleanUp.shutdownNow()" - + " is called."); + public static void setMode(Mode mode) { + synchronized (lock) { + if (SingletonManager.mode == mode) { + return; + } + if (SingletonManager.mode == Mode.CLOSED) { + throw new IllegalStateException("Cannot leave closed mode once entered"); + } + if (SingletonManager.mode == Mode.CLIENT && mode == Mode.CONNECTOR) { + if (transitionedFromClientToConnector) { + throw new IllegalStateException("Can only transition from " + Mode.CLIENT + " to " + + Mode.CONNECTOR + " once. This error indicates that " + + "org.apache.accumulo.core.util.CleanUp.shutdownNow() was called and then later a " + + "Connector was created. Connectors can not be created after CleanUp.shutdownNow()" + + " is called."); + } + + transitionedFromClientToConnector = true; } - transitionedFromClientToConnector = true; - } - - /* - * Always allow transition to closed and only allow transition to client/connector when the - * current mode is not server. - */ - if (SingletonManager.mode != Mode.SERVER || mode == Mode.CLOSED) { - SingletonManager.mode = mode; + /* + * Always allow transition to closed and only allow transition to client/connector when the + * current mode is not server. + */ + if (SingletonManager.mode != Mode.SERVER || mode == Mode.CLOSED) { + SingletonManager.mode = mode; + } + transition(); } - transition(); } @VisibleForTesting - public static synchronized Mode getMode() { - return mode; + public static Mode getMode() { + synchronized (lock) { + return mode; + } } private static void transition() { diff --git a/core/src/main/java/org/apache/accumulo/core/util/ratelimit/SharedRateLimiterFactory.java b/core/src/main/java/org/apache/accumulo/core/util/ratelimit/SharedRateLimiterFactory.java index 7b086522d5e..329b4f35fa3 100644 --- a/core/src/main/java/org/apache/accumulo/core/util/ratelimit/SharedRateLimiterFactory.java +++ b/core/src/main/java/org/apache/accumulo/core/util/ratelimit/SharedRateLimiterFactory.java @@ -42,6 +42,7 @@ public class SharedRateLimiterFactory { private static final long REPORT_RATE = 60000; private static final long UPDATE_RATE = 1000; + private static final Object INSTANCE_LOCK = new Object(); private static SharedRateLimiterFactory instance = null; private static ScheduledFuture updateTaskFuture; private final Logger log = LoggerFactory.getLogger(SharedRateLimiterFactory.class); @@ -51,22 +52,23 @@ public class SharedRateLimiterFactory { private SharedRateLimiterFactory() {} /** Get the singleton instance of the SharedRateLimiterFactory. */ - public static synchronized SharedRateLimiterFactory - getInstance(ScheduledThreadPoolExecutor executor) { - if (instance == null) { - instance = new SharedRateLimiterFactory(); + public static SharedRateLimiterFactory getInstance(ScheduledThreadPoolExecutor executor) { + synchronized (INSTANCE_LOCK) { + if (instance == null) { + instance = new SharedRateLimiterFactory(); - updateTaskFuture = executor.scheduleWithFixedDelay(Threads - .createNamedRunnable("SharedRateLimiterFactory update polling", instance::updateAll), - UPDATE_RATE, UPDATE_RATE, MILLISECONDS); + updateTaskFuture = executor.scheduleWithFixedDelay(Threads + .createNamedRunnable("SharedRateLimiterFactory update polling", instance::updateAll), + UPDATE_RATE, UPDATE_RATE, MILLISECONDS); - ScheduledFuture future = executor.scheduleWithFixedDelay(Threads - .createNamedRunnable("SharedRateLimiterFactory report polling", instance::reportAll), - REPORT_RATE, REPORT_RATE, MILLISECONDS); - ThreadPools.watchNonCriticalScheduledTask(future); + ScheduledFuture future = executor.scheduleWithFixedDelay(Threads + .createNamedRunnable("SharedRateLimiterFactory report polling", instance::reportAll), + REPORT_RATE, REPORT_RATE, MILLISECONDS); + ThreadPools.watchNonCriticalScheduledTask(future); + } + return instance; } - return instance; } /** diff --git a/core/src/main/spotbugs/exclude-filter.xml b/core/src/main/spotbugs/exclude-filter.xml index c503d6f344f..fa213ce568d 100644 --- a/core/src/main/spotbugs/exclude-filter.xml +++ b/core/src/main/spotbugs/exclude-filter.xml @@ -43,14 +43,11 @@ - - - - - + diff --git a/core/src/test/java/org/apache/accumulo/core/fate/zookeeper/DistributedReadWriteLockTest.java b/core/src/test/java/org/apache/accumulo/core/fate/zookeeper/DistributedReadWriteLockTest.java index da30442d933..fbdcfbdf5b9 100644 --- a/core/src/test/java/org/apache/accumulo/core/fate/zookeeper/DistributedReadWriteLockTest.java +++ b/core/src/test/java/org/apache/accumulo/core/fate/zookeeper/DistributedReadWriteLockTest.java @@ -37,7 +37,7 @@ public class DistributedReadWriteLockTest { public static class MockQueueLock implements QueueLock { long next = 0L; - final SortedMap locks = new TreeMap<>(); + private final SortedMap locks = new TreeMap<>(); @Override public synchronized SortedMap getEarlierEntries(long entry) { @@ -47,7 +47,7 @@ public synchronized SortedMap getEarlierEntries(long entry) { } @Override - public synchronized void removeEntry(long entry) { + public void removeEntry(long entry) { synchronized (locks) { locks.remove(entry); locks.notifyAll(); @@ -55,7 +55,7 @@ public synchronized void removeEntry(long entry) { } @Override - public synchronized long addEntry(byte[] data) { + public long addEntry(byte[] data) { long result; synchronized (locks) { locks.put(result = next++, data);