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 f1776b30607..fa213ce568d 100644 --- a/core/src/main/spotbugs/exclude-filter.xml +++ b/core/src/main/spotbugs/exclude-filter.xml @@ -43,4 +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); 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 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/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); 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/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", 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 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/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/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/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; 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); } } 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/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 @@ + + + + + +