This is an automated email from the ASF dual-hosted git repository.
DomGarguilo pushed a commit to branch 2.1
in repository https://gitbox.apache.org/repos/asf/accumulo.git
The following commit(s) were added to refs/heads/2.1 by this push:
new ed7e1fd7d3 Enables SharedVariableAtomicityDetector (#6025)
ed7e1fd7d3 is described below
commit ed7e1fd7d339ff5883b18f6861fde23ffd39e223
Author: Kevin Rathbun <[email protected]>
AuthorDate: Fri Aug 28 12:44:12 2026 -0400
Enables SharedVariableAtomicityDetector (#6025)
* Enables SharedVariableAtomicityDetector
* Fix concurrency issues that are related to the affected sites
---------
Co-authored-by: Dom Garguilo <[email protected]>
---
core/pom.xml | 5 ++++
.../accumulo/core/client/rfile/RFileScanner.java | 3 +++
.../accumulo/core/clientImpl/ScannerOptions.java | 3 +++
.../TabletServerBatchReaderIterator.java | 5 +++-
.../core/clientImpl/ThriftTransportPool.java | 4 ++++
.../core/crypto/streams/BlockedInputStream.java | 3 +++
.../fate/zookeeper/DistributedReadWriteLock.java | 5 ++++
.../accumulo/core/file/BloomFilterLayer.java | 3 +++
.../impl/SeekableByteArrayInputStream.java | 3 +++
.../accumulo/core/file/map/MapFileOperations.java | 3 +++
.../org/apache/accumulo/core/file/rfile/RFile.java | 4 ++++
.../rfile/bcfile/SimpleBufferedOutputStream.java | 3 +++
.../file/streams/BoundedRangeFileInputStream.java | 3 +++
.../system/ColumnFamilySkippingIterator.java | 3 +++
.../core/iteratorsImpl/system/StatsIterator.java | 3 +++
.../core/metadata/schema/TabletsMetadata.java | 3 +++
.../accumulo/core/singletons/SingletonManager.java | 4 ++++
.../spi/compaction/DefaultCompactionPlanner.java | 6 ++---
.../accumulo/core/spi/crypto/AESCryptoService.java | 2 ++
.../accumulo/core/util/CountingInputStream.java | 3 +++
pom.xml | 2 +-
server/base/pom.xml | 5 ++++
.../server/compaction/CountingIterator.java | 3 +++
.../org/apache/accumulo/server/fs/FileManager.java | 3 +++
.../server/problems/ProblemReportingIterator.java | 3 +++
.../delegation/AuthenticationTokenKeyManager.java | 8 +++++++
.../apache/accumulo/server/tablets/TabletTime.java | 14 ++++-------
.../main/java/org/apache/accumulo/gc/GCRun.java | 27 +++++++++++-----------
server/manager/pom.xml | 5 ++++
.../manager/tableOps/bulkVer2/LoadFiles.java | 3 +++
.../java/org/apache/accumulo/monitor/Monitor.java | 20 +++++++---------
server/tserver/pom.xml | 5 ++++
.../org/apache/accumulo/tserver/InMemoryMap.java | 3 +++
.../org/apache/accumulo/tserver/NativeMap.java | 3 +++
.../tserver/TabletServerResourceManager.java | 2 +-
.../accumulo/tserver/tablet/ScanDataSource.java | 3 +++
.../apache/accumulo/tserver/tablet/Scanner.java | 26 ++++++++++++++++-----
.../org/apache/accumulo/tserver/tablet/Tablet.java | 4 ++--
.../tablet/CompactableImplFileManagerTest.java | 3 +++
test/pom.xml | 5 ++++
.../test/functional/ErrorThrowingIterator.java | 3 +++
41 files changed, 173 insertions(+), 48 deletions(-)
diff --git a/core/pom.xml b/core/pom.xml
index 15951126d6..395c8bf8c0 100644
--- a/core/pom.xml
+++ b/core/pom.xml
@@ -43,6 +43,11 @@
<groupId>com.github.ben-manes.caffeine</groupId>
<artifactId>caffeine</artifactId>
</dependency>
+ <dependency>
+ <groupId>com.google.code.findbugs</groupId>
+ <artifactId>jsr305</artifactId>
+ <optional>true</optional>
+ </dependency>
<dependency>
<groupId>com.google.code.gson</groupId>
<artifactId>gson</artifactId>
diff --git
a/core/src/main/java/org/apache/accumulo/core/client/rfile/RFileScanner.java
b/core/src/main/java/org/apache/accumulo/core/client/rfile/RFileScanner.java
index 49e424b456..f926a4fffc 100644
--- a/core/src/main/java/org/apache/accumulo/core/client/rfile/RFileScanner.java
+++ b/core/src/main/java/org/apache/accumulo/core/client/rfile/RFileScanner.java
@@ -28,6 +28,8 @@ import java.util.Map.Entry;
import java.util.Set;
import java.util.SortedSet;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.client.IteratorSetting;
import org.apache.accumulo.core.client.Scanner;
import org.apache.accumulo.core.client.TableNotFoundException;
@@ -77,6 +79,7 @@ import org.apache.hadoop.io.Text;
import com.google.common.base.Preconditions;
+@NotThreadSafe
class RFileScanner extends ScannerOptions implements Scanner {
private static class RFileScannerEnvironmentImpl extends
ClientServiceEnvironmentImpl {
diff --git
a/core/src/main/java/org/apache/accumulo/core/clientImpl/ScannerOptions.java
b/core/src/main/java/org/apache/accumulo/core/clientImpl/ScannerOptions.java
index 2b044f6fd8..fb288a519f 100644
--- a/core/src/main/java/org/apache/accumulo/core/clientImpl/ScannerOptions.java
+++ b/core/src/main/java/org/apache/accumulo/core/clientImpl/ScannerOptions.java
@@ -35,6 +35,8 @@ import java.util.SortedSet;
import java.util.TreeSet;
import java.util.concurrent.TimeUnit;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.client.IteratorSetting;
import org.apache.accumulo.core.client.ScannerBase;
import org.apache.accumulo.core.client.sample.SamplerConfiguration;
@@ -46,6 +48,7 @@ import org.apache.accumulo.core.security.Authorizations;
import org.apache.accumulo.core.util.TextUtil;
import org.apache.hadoop.io.Text;
+@NotThreadSafe
public class ScannerOptions implements ScannerBase {
protected List<IterInfo> serverSideIteratorList = Collections.emptyList();
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 ccd5ec6e08..51e777a2e2 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
@@ -46,6 +46,8 @@ import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Supplier;
import java.util.stream.Collectors;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.client.AccumuloException;
import org.apache.accumulo.core.client.AccumuloSecurityException;
import org.apache.accumulo.core.client.SampleNotPresentException;
@@ -360,6 +362,7 @@ public class TabletServerBatchReaderIterator implements
Iterator<Entry<Key,Value
return context.getPrintableTableInfoFromId(tableId);
}
+ @NotThreadSafe
private class QueryTask implements Runnable {
private final String tsLocation;
@@ -759,7 +762,7 @@ public class TabletServerBatchReaderIterator implements
Iterator<Entry<Key,Value
}
class Session {
- long activityTime;
+ volatile long activityTime;
void check() throws IOException {
if (System.currentTimeMillis() - activityTime > timeOut) {
diff --git
a/core/src/main/java/org/apache/accumulo/core/clientImpl/ThriftTransportPool.java
b/core/src/main/java/org/apache/accumulo/core/clientImpl/ThriftTransportPool.java
index 2f543eed31..a5c603cc50 100644
---
a/core/src/main/java/org/apache/accumulo/core/clientImpl/ThriftTransportPool.java
+++
b/core/src/main/java/org/apache/accumulo/core/clientImpl/ThriftTransportPool.java
@@ -578,6 +578,10 @@ public class ThriftTransportPool {
private static final long serialVersionUID = 1L;
}
+ @SuppressFBWarnings(
+ value = {"AT_NONATOMIC_OPERATIONS_ON_SHARED_VARIABLE",
"AT_STALE_THREAD_WRITE_OF_PRIMITIVE"},
+ justification = "stuck IO counters are best effort diagnostics read by
the pool's "
+ + "background thread; only the reserving thread increments ioCount")
private static class CachedTTransport extends TTransport {
private final ThriftTransportKey cacheKey;
diff --git
a/core/src/main/java/org/apache/accumulo/core/crypto/streams/BlockedInputStream.java
b/core/src/main/java/org/apache/accumulo/core/crypto/streams/BlockedInputStream.java
index 3e37f6f404..4777b34153 100644
---
a/core/src/main/java/org/apache/accumulo/core/crypto/streams/BlockedInputStream.java
+++
b/core/src/main/java/org/apache/accumulo/core/crypto/streams/BlockedInputStream.java
@@ -23,10 +23,13 @@ import java.io.EOFException;
import java.io.IOException;
import java.io.InputStream;
+import javax.annotation.concurrent.NotThreadSafe;
+
/**
* Reader corresponding to BlockedOutputStream. Expects all data to be in the
form of size (int)
* data (size bytes) junk (however many bytes it takes to complete a block)
*/
+@NotThreadSafe
public class BlockedInputStream extends InputStream {
byte[] array;
// ReadPos is where to start reading
diff --git
a/core/src/main/java/org/apache/accumulo/core/fate/zookeeper/DistributedReadWriteLock.java
b/core/src/main/java/org/apache/accumulo/core/fate/zookeeper/DistributedReadWriteLock.java
index da27d408a0..50ec14df8f 100644
---
a/core/src/main/java/org/apache/accumulo/core/fate/zookeeper/DistributedReadWriteLock.java
+++
b/core/src/main/java/org/apache/accumulo/core/fate/zookeeper/DistributedReadWriteLock.java
@@ -30,6 +30,8 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.util.UtilWaitThread;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -107,6 +109,9 @@ public class DistributedReadWriteLock implements
java.util.concurrent.locks.Read
private static final Logger log =
LoggerFactory.getLogger(DistributedReadWriteLock.class);
+ // Intended usage of instances of this class is that it's only used by a
single thread and
+ // supports locking across threads/processes
+ @NotThreadSafe
static class ReadLock implements Lock {
QueueLock qlock;
diff --git
a/core/src/main/java/org/apache/accumulo/core/file/BloomFilterLayer.java
b/core/src/main/java/org/apache/accumulo/core/file/BloomFilterLayer.java
index e8f53d47bf..81dfc3b85d 100644
--- a/core/src/main/java/org/apache/accumulo/core/file/BloomFilterLayer.java
+++ b/core/src/main/java/org/apache/accumulo/core/file/BloomFilterLayer.java
@@ -30,6 +30,8 @@ import java.util.Set;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.atomic.AtomicBoolean;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.bloomfilter.DynamicBloomFilter;
import org.apache.accumulo.core.classloader.ClassLoaderUtil;
import org.apache.accumulo.core.conf.AccumuloConfiguration;
@@ -333,6 +335,7 @@ public class BloomFilterLayer {
}
}
+ @NotThreadSafe
public static class Reader implements FileSKVIterator {
private final BloomFilterLoader bfl;
diff --git
a/core/src/main/java/org/apache/accumulo/core/file/blockfile/impl/SeekableByteArrayInputStream.java
b/core/src/main/java/org/apache/accumulo/core/file/blockfile/impl/SeekableByteArrayInputStream.java
index c231e88b94..732123d199 100644
---
a/core/src/main/java/org/apache/accumulo/core/file/blockfile/impl/SeekableByteArrayInputStream.java
+++
b/core/src/main/java/org/apache/accumulo/core/file/blockfile/impl/SeekableByteArrayInputStream.java
@@ -23,12 +23,15 @@ import static java.util.Objects.requireNonNull;
import java.io.IOException;
import java.io.InputStream;
+import javax.annotation.concurrent.NotThreadSafe;
+
import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
/**
* This class is like byte array input stream with two differences. It
supports seeking and avoids
* synchronization.
*/
+@NotThreadSafe
public class SeekableByteArrayInputStream extends InputStream {
// making this volatile for the following case
diff --git
a/core/src/main/java/org/apache/accumulo/core/file/map/MapFileOperations.java
b/core/src/main/java/org/apache/accumulo/core/file/map/MapFileOperations.java
index c2da8c0066..26d24e5846 100644
---
a/core/src/main/java/org/apache/accumulo/core/file/map/MapFileOperations.java
+++
b/core/src/main/java/org/apache/accumulo/core/file/map/MapFileOperations.java
@@ -24,6 +24,8 @@ import java.util.Collection;
import java.util.Map;
import java.util.concurrent.atomic.AtomicBoolean;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.data.ByteSequence;
import org.apache.accumulo.core.data.Key;
import org.apache.accumulo.core.data.Range;
@@ -43,6 +45,7 @@ import org.apache.hadoop.io.MapFile;
public class MapFileOperations extends FileOperations {
private static final String MSG = "Map files are not supported";
+ @NotThreadSafe
public static class RangeIterator implements FileSKVIterator {
SortedKeyValueIterator<Key,Value> reader;
diff --git a/core/src/main/java/org/apache/accumulo/core/file/rfile/RFile.java
b/core/src/main/java/org/apache/accumulo/core/file/rfile/RFile.java
index 0e941129b8..7b6f4fcd2b 100644
--- a/core/src/main/java/org/apache/accumulo/core/file/rfile/RFile.java
+++ b/core/src/main/java/org/apache/accumulo/core/file/rfile/RFile.java
@@ -40,6 +40,8 @@ import java.util.Set;
import java.util.TreeMap;
import java.util.concurrent.atomic.AtomicBoolean;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.client.SampleNotPresentException;
import org.apache.accumulo.core.client.sample.Sampler;
import org.apache.accumulo.core.client.sample.SamplerConfiguration;
@@ -557,6 +559,7 @@ public class RFile {
}
}
+ @NotThreadSafe
public static class Writer implements FileSKVWriter {
public static final int MAX_CF_IN_DLG = 1000;
@@ -754,6 +757,7 @@ public class RFile {
}
}
+ @NotThreadSafe
private static class LocalityGroupReader extends LocalityGroup implements
FileSKVIterator {
private final CachableBlockFile.Reader reader;
diff --git
a/core/src/main/java/org/apache/accumulo/core/file/rfile/bcfile/SimpleBufferedOutputStream.java
b/core/src/main/java/org/apache/accumulo/core/file/rfile/bcfile/SimpleBufferedOutputStream.java
index 87e96a602d..e1bd64ed78 100644
---
a/core/src/main/java/org/apache/accumulo/core/file/rfile/bcfile/SimpleBufferedOutputStream.java
+++
b/core/src/main/java/org/apache/accumulo/core/file/rfile/bcfile/SimpleBufferedOutputStream.java
@@ -22,10 +22,13 @@ import java.io.FilterOutputStream;
import java.io.IOException;
import java.io.OutputStream;
+import javax.annotation.concurrent.NotThreadSafe;
+
/**
* A simplified BufferedOutputStream with borrowed buffer, and allow users to
see how much data have
* been buffered.
*/
+@NotThreadSafe
class SimpleBufferedOutputStream extends FilterOutputStream {
protected byte[] buf; // the borrowed buffer
protected int count = 0; // bytes used in buffer.
diff --git
a/core/src/main/java/org/apache/accumulo/core/file/streams/BoundedRangeFileInputStream.java
b/core/src/main/java/org/apache/accumulo/core/file/streams/BoundedRangeFileInputStream.java
index d9f41862ae..ff6a7be71b 100644
---
a/core/src/main/java/org/apache/accumulo/core/file/streams/BoundedRangeFileInputStream.java
+++
b/core/src/main/java/org/apache/accumulo/core/file/streams/BoundedRangeFileInputStream.java
@@ -21,6 +21,8 @@ package org.apache.accumulo.core.file.streams;
import java.io.IOException;
import java.io.InputStream;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.hadoop.fs.Seekable;
/**
@@ -28,6 +30,7 @@ import org.apache.hadoop.fs.Seekable;
* regular input stream. One can create multiple BoundedRangeFileInputStream
on top of the same
* FSDataInputStream and they would not interfere with each other.
*/
+@NotThreadSafe
public class BoundedRangeFileInputStream extends InputStream {
private volatile boolean closed = false;
diff --git
a/core/src/main/java/org/apache/accumulo/core/iteratorsImpl/system/ColumnFamilySkippingIterator.java
b/core/src/main/java/org/apache/accumulo/core/iteratorsImpl/system/ColumnFamilySkippingIterator.java
index ae26b0dd81..ba2d912b72 100644
---
a/core/src/main/java/org/apache/accumulo/core/iteratorsImpl/system/ColumnFamilySkippingIterator.java
+++
b/core/src/main/java/org/apache/accumulo/core/iteratorsImpl/system/ColumnFamilySkippingIterator.java
@@ -25,6 +25,8 @@ import java.util.Set;
import java.util.TreeSet;
import java.util.concurrent.atomic.AtomicBoolean;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.data.ByteSequence;
import org.apache.accumulo.core.data.Column;
import org.apache.accumulo.core.data.Key;
@@ -36,6 +38,7 @@ import org.apache.accumulo.core.iterators.IteratorEnvironment;
import org.apache.accumulo.core.iterators.ServerSkippingIterator;
import org.apache.accumulo.core.iterators.SortedKeyValueIterator;
+@NotThreadSafe
public class ColumnFamilySkippingIterator extends ServerSkippingIterator
implements InterruptibleIterator {
diff --git
a/core/src/main/java/org/apache/accumulo/core/iteratorsImpl/system/StatsIterator.java
b/core/src/main/java/org/apache/accumulo/core/iteratorsImpl/system/StatsIterator.java
index 87215d34b5..66923002ad 100644
---
a/core/src/main/java/org/apache/accumulo/core/iteratorsImpl/system/StatsIterator.java
+++
b/core/src/main/java/org/apache/accumulo/core/iteratorsImpl/system/StatsIterator.java
@@ -26,6 +26,8 @@ import java.util.List;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.LongAdder;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.data.ByteSequence;
import org.apache.accumulo.core.data.Key;
import org.apache.accumulo.core.data.Range;
@@ -34,6 +36,7 @@ import org.apache.accumulo.core.iterators.IteratorEnvironment;
import org.apache.accumulo.core.iterators.ServerWrappingIterator;
import org.apache.accumulo.core.iterators.SortedKeyValueIterator;
+@NotThreadSafe
public class StatsIterator extends ServerWrappingIterator {
private int numRead = 0;
diff --git
a/core/src/main/java/org/apache/accumulo/core/metadata/schema/TabletsMetadata.java
b/core/src/main/java/org/apache/accumulo/core/metadata/schema/TabletsMetadata.java
index 6cff9ae59f..90095cc5bb 100644
---
a/core/src/main/java/org/apache/accumulo/core/metadata/schema/TabletsMetadata.java
+++
b/core/src/main/java/org/apache/accumulo/core/metadata/schema/TabletsMetadata.java
@@ -44,6 +44,8 @@ import java.util.function.Function;
import java.util.stream.Stream;
import java.util.stream.StreamSupport;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.client.AccumuloClient;
import org.apache.accumulo.core.client.BatchScanner;
import org.apache.accumulo.core.client.IsolatedScanner;
@@ -89,6 +91,7 @@ import com.google.common.collect.Iterators;
*/
public class TabletsMetadata implements Iterable<TabletMetadata>,
AutoCloseable {
+ @NotThreadSafe
public static class Builder implements TableRangeOptions, TableOptions,
RangeOptions, Options {
private final List<Text> families = new ArrayList<>();
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 5f3e151b58..a8dee7fffb 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
@@ -27,6 +27,8 @@ import org.slf4j.LoggerFactory;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Preconditions;
+import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
+
/**
* This class automates management of static singletons that maintain state
for Accumulo clients.
* Historically, Accumulo client code that used Connector had no control over
these singletons. The
@@ -81,6 +83,8 @@ public class SingletonManager {
private static List<SingletonService> services;
@VisibleForTesting
+ @SuppressFBWarnings(value = "AT_NONATOMIC_64BIT_PRIMITIVE",
+ justification = "only called in static init block and testing, no sync
needed")
static void reset() {
reservations = 0;
mode = Mode.CLIENT;
diff --git
a/core/src/main/java/org/apache/accumulo/core/spi/compaction/DefaultCompactionPlanner.java
b/core/src/main/java/org/apache/accumulo/core/spi/compaction/DefaultCompactionPlanner.java
index 0103243ef4..5b67ed6642 100644
---
a/core/src/main/java/org/apache/accumulo/core/spi/compaction/DefaultCompactionPlanner.java
+++
b/core/src/main/java/org/apache/accumulo/core/spi/compaction/DefaultCompactionPlanner.java
@@ -165,9 +165,9 @@ public class DefaultCompactionPlanner implements
CompactionPlanner {
}
}
- private List<Executor> executors;
- private int maxFilesToCompact;
- private double lowestRatio;
+ private volatile List<Executor> executors;
+ private volatile int maxFilesToCompact;
+ private volatile double lowestRatio;
@SuppressFBWarnings(value = {"UWF_UNWRITTEN_FIELD", "NP_UNWRITTEN_FIELD"},
justification = "Field is written by Gson")
diff --git
a/core/src/main/java/org/apache/accumulo/core/spi/crypto/AESCryptoService.java
b/core/src/main/java/org/apache/accumulo/core/spi/crypto/AESCryptoService.java
index 555103fd98..86aa95b26d 100644
---
a/core/src/main/java/org/apache/accumulo/core/spi/crypto/AESCryptoService.java
+++
b/core/src/main/java/org/apache/accumulo/core/spi/crypto/AESCryptoService.java
@@ -117,6 +117,8 @@ public class AESCryptoService implements CryptoService {
};
@Override
+ @SuppressFBWarnings(value = "AT_STALE_THREAD_WRITE_OF_PRIMITIVE",
+ justification = "Fields modified here are initialized once, and
read-only after.")
public void init(Map<String,String> conf) throws CryptoException {
ensureNotInit();
String keyLocation = Objects.requireNonNull(conf.get(KEY_URI_PROPERTY),
diff --git
a/core/src/main/java/org/apache/accumulo/core/util/CountingInputStream.java
b/core/src/main/java/org/apache/accumulo/core/util/CountingInputStream.java
index 15d41b18d7..5f6a03871e 100644
--- a/core/src/main/java/org/apache/accumulo/core/util/CountingInputStream.java
+++ b/core/src/main/java/org/apache/accumulo/core/util/CountingInputStream.java
@@ -18,6 +18,8 @@ import java.io.IOException;
import java.io.InputStream;
import java.util.Objects;
+import javax.annotation.concurrent.NotThreadSafe;
+
/**
* This class was copied from Guava and modified. If this class was not final
in Guava it could have
* been extended. Guava has issue 590 open about this.
@@ -26,6 +28,7 @@ import java.util.Objects;
*
* @author Chris Nokleberg
*/
+@NotThreadSafe
public final class CountingInputStream extends FilterInputStream {
private long count;
diff --git a/pom.xml b/pom.xml
index 331e4ddbca..38064f712a 100644
--- a/pom.xml
+++ b/pom.xml
@@ -788,7 +788,7 @@ under the License.
<includeTests>true</includeTests>
<maxHeap>1024</maxHeap>
<maxRank>16</maxRank>
-
<omitVisitors>ConstructorThrow,SharedVariableAtomicityDetector</omitVisitors>
+ <omitVisitors>ConstructorThrow</omitVisitors>
<jvmArgs>-Dcom.overstock.findbugs.ignore=com.google.common.util.concurrent.RateLimiter,com.google.common.hash.Hasher,com.google.common.hash.HashCode,com.google.common.hash.HashFunction,com.google.common.hash.Hashing,com.google.common.cache.Cache,com.google.common.io.CountingOutputStream,com.google.common.io.ByteStreams,com.google.common.cache.LoadingCache,com.google.common.base.Stopwatch,com.google.common.cache.RemovalNotification,com.google.common.util.concurrent.Uninterrupt
[...]
<plugins combine.children="append">
<plugin>
diff --git a/server/base/pom.xml b/server/base/pom.xml
index 3a200701ba..ed184918b0 100644
--- a/server/base/pom.xml
+++ b/server/base/pom.xml
@@ -39,6 +39,11 @@
<groupId>com.github.ben-manes.caffeine</groupId>
<artifactId>caffeine</artifactId>
</dependency>
+ <dependency>
+ <groupId>com.google.code.findbugs</groupId>
+ <artifactId>jsr305</artifactId>
+ <optional>true</optional>
+ </dependency>
<dependency>
<groupId>com.google.code.gson</groupId>
<artifactId>gson</artifactId>
diff --git
a/server/base/src/main/java/org/apache/accumulo/server/compaction/CountingIterator.java
b/server/base/src/main/java/org/apache/accumulo/server/compaction/CountingIterator.java
index 4c768ba6f0..f029f02284 100644
---
a/server/base/src/main/java/org/apache/accumulo/server/compaction/CountingIterator.java
+++
b/server/base/src/main/java/org/apache/accumulo/server/compaction/CountingIterator.java
@@ -23,12 +23,15 @@ import java.util.ArrayList;
import java.util.Map;
import java.util.concurrent.atomic.AtomicLong;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.data.Key;
import org.apache.accumulo.core.data.Value;
import org.apache.accumulo.core.iterators.IteratorEnvironment;
import org.apache.accumulo.core.iterators.SortedKeyValueIterator;
import org.apache.accumulo.core.iterators.WrappingIterator;
+@NotThreadSafe
public class CountingIterator extends WrappingIterator {
private long count;
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 4e81899ed2..24ff167019 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
@@ -33,6 +33,8 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.stream.Collectors;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.client.SampleNotPresentException;
import org.apache.accumulo.core.conf.Property;
import org.apache.accumulo.core.data.Key;
@@ -388,6 +390,7 @@ public class FileManager {
}
+ @NotThreadSafe
static class FileDataSource implements DataSource {
private SortedKeyValueIterator<Key,Value> iter;
diff --git
a/server/base/src/main/java/org/apache/accumulo/server/problems/ProblemReportingIterator.java
b/server/base/src/main/java/org/apache/accumulo/server/problems/ProblemReportingIterator.java
index 4344c6ba97..e72cd6dab4 100644
---
a/server/base/src/main/java/org/apache/accumulo/server/problems/ProblemReportingIterator.java
+++
b/server/base/src/main/java/org/apache/accumulo/server/problems/ProblemReportingIterator.java
@@ -23,6 +23,8 @@ import java.util.Collection;
import java.util.Map;
import java.util.concurrent.atomic.AtomicBoolean;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.data.ByteSequence;
import org.apache.accumulo.core.data.Key;
import org.apache.accumulo.core.data.Range;
@@ -33,6 +35,7 @@ import
org.apache.accumulo.core.iterators.SortedKeyValueIterator;
import org.apache.accumulo.core.iteratorsImpl.system.InterruptibleIterator;
import org.apache.accumulo.server.ServerContext;
+@NotThreadSafe
public class ProblemReportingIterator implements InterruptibleIterator {
private final SortedKeyValueIterator<Key,Value> source;
private boolean sawError = false;
diff --git
a/server/base/src/main/java/org/apache/accumulo/server/security/delegation/AuthenticationTokenKeyManager.java
b/server/base/src/main/java/org/apache/accumulo/server/security/delegation/AuthenticationTokenKeyManager.java
index c98c861ac5..8ebe302db4 100644
---
a/server/base/src/main/java/org/apache/accumulo/server/security/delegation/AuthenticationTokenKeyManager.java
+++
b/server/base/src/main/java/org/apache/accumulo/server/security/delegation/AuthenticationTokenKeyManager.java
@@ -26,6 +26,8 @@ import org.slf4j.LoggerFactory;
import com.google.common.annotations.VisibleForTesting;
+import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
+
/**
* Service that handles generation of the secret key used to create delegation
tokens.
*/
@@ -94,6 +96,9 @@ public class AuthenticationTokenKeyManager implements
Runnable {
}
@VisibleForTesting
+ @SuppressFBWarnings(
+ value = {"AT_STALE_THREAD_WRITE_OF_PRIMITIVE",
"AT_NONATOMIC_64BIT_PRIMITIVE"},
+ justification = "only called from run() and testing")
void updateStateFromCurrentKeys() {
try {
List<AuthenticationKey> currentKeys = keyDistributor.getCurrentKeys();
@@ -142,6 +147,9 @@ public class AuthenticationTokenKeyManager implements
Runnable {
*
* @param now The current time in millis since epoch.
*/
+ @SuppressFBWarnings(
+ value = {"AT_STALE_THREAD_WRITE_OF_PRIMITIVE",
"AT_NONATOMIC_64BIT_PRIMITIVE"},
+ justification = "only called from run() and testing")
void _run(long now) {
// clear any expired keys
int removedKeys = secretManager.removeExpiredKeys(keyDistributor);
diff --git
a/server/base/src/main/java/org/apache/accumulo/server/tablets/TabletTime.java
b/server/base/src/main/java/org/apache/accumulo/server/tablets/TabletTime.java
index 45d3870db1..ead338961d 100644
---
a/server/base/src/main/java/org/apache/accumulo/server/tablets/TabletTime.java
+++
b/server/base/src/main/java/org/apache/accumulo/server/tablets/TabletTime.java
@@ -77,7 +77,7 @@ public abstract class TabletTime {
}
@Override
- public MetadataTime getMetadataTime() {
+ public synchronized MetadataTime getMetadataTime() {
return getMetadataTime(lastTime);
}
@@ -87,7 +87,7 @@ public abstract class TabletTime {
}
@Override
- public void useMaxTimeFromWALog(long time) {
+ public synchronized void useMaxTimeFromWALog(long time) {
if (time > lastTime) {
lastTime = time;
}
@@ -113,7 +113,7 @@ public abstract class TabletTime {
return currTime;
}
- private long updateTime(long currTime) {
+ private synchronized long updateTime(long currTime) {
if (currTime < lastTime) {
if (currTime - lastUpdateTime > 0) {
// not in same millisecond as last call
@@ -131,7 +131,7 @@ public abstract class TabletTime {
}
@Override
- public long getTime() {
+ public synchronized long getTime() {
return lastTime;
}
@@ -157,11 +157,7 @@ public abstract class TabletTime {
@Override
public void useMaxTimeFromWALog(long time) {
- time++;
-
- if (this.nextTime.get() < time) {
- this.nextTime.set(time);
- }
+ this.nextTime.getAndUpdate(val -> Math.max(val, time + 1));
}
@Override
diff --git a/server/gc/src/main/java/org/apache/accumulo/gc/GCRun.java
b/server/gc/src/main/java/org/apache/accumulo/gc/GCRun.java
index cbd968245e..567682c720 100644
--- a/server/gc/src/main/java/org/apache/accumulo/gc/GCRun.java
+++ b/server/gc/src/main/java/org/apache/accumulo/gc/GCRun.java
@@ -40,6 +40,7 @@ import java.util.SortedMap;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.LongAdder;
import java.util.stream.Stream;
import org.apache.accumulo.core.Constants;
@@ -100,10 +101,10 @@ public class GCRun implements
GarbageCollectionEnvironment {
private final ServerContext context;
private final AccumuloConfiguration config;
private final Duration loggingInterval = Duration.ofMinutes(1);
- private long candidates = 0;
- private long inUse = 0;
- private long deleted = 0;
- private long errors = 0;
+ private final LongAdder candidates = new LongAdder();
+ private final LongAdder inUse = new LongAdder();
+ private final LongAdder deleted = new LongAdder();
+ private final LongAdder errors = new LongAdder();
private AtomicInteger batchCount;
public GCRun(Ample.DataLevel level, ServerContext context) {
@@ -337,18 +338,18 @@ public class GCRun implements
GarbageCollectionEnvironment {
if (moveToTrash(pathToDel) || fs.deleteRecursively(pathToDel)) {
// delete succeeded, still want to delete
removeFlag = true;
- deleted++;
+ deleted.increment();
} else if (fs.exists(pathToDel)) {
// leave the entry in the metadata; we'll try again later
removeFlag = false;
- errors++;
+ errors.increment();
log.warn("{} File exists, but was not deleted for an unknown
reason: {}",
fileActionPrefix, pathToDel);
break;
} else {
// this failure, we still want to remove the metadata entry
removeFlag = true;
- errors++;
+ errors.increment();
String[] parts =
pathToDel.toString().split(Constants.ZTABLES)[1].split("/");
if (parts.length > 2) {
TableId tableId = TableId.of(parts[1]);
@@ -425,12 +426,12 @@ public class GCRun implements
GarbageCollectionEnvironment {
@Override
public void incrementCandidatesStat(long i) {
- candidates += i;
+ candidates.add(i);
}
@Override
public void incrementInUseStat(long i) {
- inUse += i;
+ inUse.add(i);
}
@Override
@@ -576,19 +577,19 @@ public class GCRun implements
GarbageCollectionEnvironment {
}
public long getInUseStat() {
- return inUse;
+ return inUse.sum();
}
public long getDeletedStat() {
- return deleted;
+ return deleted.sum();
}
public long getErrorsStat() {
- return errors;
+ return errors.sum();
}
public long getCandidatesStat() {
- return candidates;
+ return candidates.sum();
}
/**
diff --git a/server/manager/pom.xml b/server/manager/pom.xml
index b7176a8356..ebe4130e92 100644
--- a/server/manager/pom.xml
+++ b/server/manager/pom.xml
@@ -35,6 +35,11 @@
<groupId>com.github.ben-manes.caffeine</groupId>
<artifactId>caffeine</artifactId>
</dependency>
+ <dependency>
+ <groupId>com.google.code.findbugs</groupId>
+ <artifactId>jsr305</artifactId>
+ <optional>true</optional>
+ </dependency>
<dependency>
<groupId>com.google.code.gson</groupId>
<artifactId>gson</artifactId>
diff --git
a/server/manager/src/main/java/org/apache/accumulo/manager/tableOps/bulkVer2/LoadFiles.java
b/server/manager/src/main/java/org/apache/accumulo/manager/tableOps/bulkVer2/LoadFiles.java
index e26db3c4e6..ee52348569 100644
---
a/server/manager/src/main/java/org/apache/accumulo/manager/tableOps/bulkVer2/LoadFiles.java
+++
b/server/manager/src/main/java/org/apache/accumulo/manager/tableOps/bulkVer2/LoadFiles.java
@@ -38,6 +38,8 @@ import java.util.concurrent.ExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.client.BatchWriter;
import org.apache.accumulo.core.client.MutationsRejectedException;
import org.apache.accumulo.core.clientImpl.bulk.Bulk;
@@ -159,6 +161,7 @@ class LoadFiles extends ManagerRepo {
abstract long finish() throws Exception;
}
+ @NotThreadSafe
static class OnlineLoader extends Loader {
private final int maxConnections;
diff --git
a/server/monitor/src/main/java/org/apache/accumulo/monitor/Monitor.java
b/server/monitor/src/main/java/org/apache/accumulo/monitor/Monitor.java
index d3558f7606..391528906d 100644
--- a/server/monitor/src/main/java/org/apache/accumulo/monitor/Monitor.java
+++ b/server/monitor/src/main/java/org/apache/accumulo/monitor/Monitor.java
@@ -130,14 +130,14 @@ public class Monitor extends AbstractServer implements
HighlyAvailableService {
}
private final AtomicLong lastRecalc = new AtomicLong(0L);
- private double totalIngestRate = 0.0;
- private double totalQueryRate = 0.0;
- private double totalScanRate = 0.0;
- private long totalEntries = 0L;
- private int totalTabletCount = 0;
- private long totalHoldTime = 0;
- private long totalLookups = 0;
- private int totalTables = 0;
+ private volatile double totalIngestRate = 0.0;
+ private volatile double totalQueryRate = 0.0;
+ private volatile double totalScanRate = 0.0;
+ private volatile long totalEntries = 0L;
+ private volatile int totalTabletCount = 0;
+ private volatile long totalHoldTime = 0;
+ private volatile long totalLookups = 0;
+ private volatile int totalTables = 0;
private final AtomicBoolean monitorInitialized = new AtomicBoolean(false);
private static <T> List<Pair<Long,T>> newMaxList() {
@@ -1008,10 +1008,6 @@ public class Monitor extends AbstractServer implements
HighlyAvailableService {
return coordinatorHost;
}
- public int getLivePort() {
- return livePort;
- }
-
@Override
public ServiceLock getLock() {
return monitorLock;
diff --git a/server/tserver/pom.xml b/server/tserver/pom.xml
index f70b264db0..728fba4451 100644
--- a/server/tserver/pom.xml
+++ b/server/tserver/pom.xml
@@ -39,6 +39,11 @@
<groupId>com.github.ben-manes.caffeine</groupId>
<artifactId>caffeine</artifactId>
</dependency>
+ <dependency>
+ <groupId>com.google.code.findbugs</groupId>
+ <artifactId>jsr305</artifactId>
+ <optional>true</optional>
+ </dependency>
<dependency>
<groupId>com.google.guava</groupId>
<artifactId>guava</artifactId>
diff --git
a/server/tserver/src/main/java/org/apache/accumulo/tserver/InMemoryMap.java
b/server/tserver/src/main/java/org/apache/accumulo/tserver/InMemoryMap.java
index 306a4e7d27..c8182e6e1f 100644
--- a/server/tserver/src/main/java/org/apache/accumulo/tserver/InMemoryMap.java
+++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/InMemoryMap.java
@@ -37,6 +37,8 @@ import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.client.SampleNotPresentException;
import org.apache.accumulo.core.client.sample.Sampler;
import org.apache.accumulo.core.conf.AccumuloConfiguration;
@@ -531,6 +533,7 @@ public class InMemoryMap {
private final Set<MemoryIterator> activeIters =
Collections.synchronizedSet(new HashSet<>());
+ @NotThreadSafe
class MemoryDataSource implements DataSource {
private boolean switched = false;
diff --git
a/server/tserver/src/main/java/org/apache/accumulo/tserver/NativeMap.java
b/server/tserver/src/main/java/org/apache/accumulo/tserver/NativeMap.java
index 0b8a6816c2..0173b2b95b 100644
--- a/server/tserver/src/main/java/org/apache/accumulo/tserver/NativeMap.java
+++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/NativeMap.java
@@ -34,6 +34,8 @@ import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.client.SampleNotPresentException;
import org.apache.accumulo.core.data.ByteSequence;
import org.apache.accumulo.core.data.ColumnUpdate;
@@ -534,6 +536,7 @@ public class NativeMap implements
Iterable<Map.Entry<Key,Value>> {
}
}
+ @NotThreadSafe
private static class NMSKVIter implements InterruptibleIterator {
private ConcurrentIterator iter;
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 7b3ec10e86..be327e0f46 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
@@ -412,7 +412,7 @@ public class TabletServerResourceManager {
public static class AssignmentWatcher implements Runnable {
private static final Logger log =
LoggerFactory.getLogger(AssignmentWatcher.class);
- private static long longAssignments = 0;
+ private static volatile long longAssignments = 0;
private final Map<KeyExtent,RunnableStartedAt> activeAssignments;
private final AccumuloConfiguration conf;
diff --git
a/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/ScanDataSource.java
b/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/ScanDataSource.java
index 602f61c845..edb71dffbc 100644
---
a/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/ScanDataSource.java
+++
b/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/ScanDataSource.java
@@ -27,6 +27,8 @@ import java.util.Map;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.client.AccumuloException;
import org.apache.accumulo.core.client.IteratorSetting;
import org.apache.accumulo.core.conf.Property;
@@ -67,6 +69,7 @@ import com.google.common.base.Preconditions;
import io.opentelemetry.api.trace.Span;
+@NotThreadSafe
class ScanDataSource implements DataSource {
private static final Logger log =
LoggerFactory.getLogger(ScanDataSource.class);
diff --git
a/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Scanner.java
b/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Scanner.java
index d382fbd171..5e170fdba7 100644
---
a/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Scanner.java
+++
b/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Scanner.java
@@ -94,10 +94,13 @@ public class Scanner {
Batch results = null;
+ boolean locked = false;
+
try {
try {
lock.lockInterruptibly();
+ locked = true;
Preconditions.checkState(!readInProgress);
// Simple check to ensure the same thread never calls this method
recursively. This code
// would not handle that well.
@@ -192,8 +195,10 @@ public class Scanner {
tablet.updateQueryStats(results.getResults().size(),
results.getNumBytes());
}
} finally {
- readInProgress = false;
- lock.unlock();
+ if (locked) {
+ readInProgress = false;
+ lock.unlock();
+ }
}
}
}
@@ -228,10 +233,7 @@ public class Scanner {
return false;
}
- scanClosed = true;
- if (isolatedDataSource != null) {
- isolatedDataSource.close(false);
- }
+ lockAndClose();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return false;
@@ -242,4 +244,16 @@ public class Scanner {
}
return true;
}
+
+ private void lockAndClose() {
+ lock.lock();
+ try {
+ scanClosed = true;
+ if (isolatedDataSource != null) {
+ isolatedDataSource.close(false);
+ }
+ } finally {
+ lock.unlock();
+ }
+ }
}
diff --git
a/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Tablet.java
b/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Tablet.java
index 63b3d95866..9fef32b6fd 100644
---
a/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Tablet.java
+++
b/server/tserver/src/main/java/org/apache/accumulo/tserver/tablet/Tablet.java
@@ -226,8 +226,8 @@ public class Tablet extends TabletBase {
private final Rate ingestByteRate = new Rate(0.95);
private final Rate scannedRate = new Rate(0.95);
- private long lastMinorCompactionFinishTime = 0;
- private long lastMapFileImportTime = 0;
+ private volatile long lastMinorCompactionFinishTime = 0;
+ private volatile long lastMapFileImportTime = 0;
private volatile long numEntries = 0;
private volatile long numEntriesInMemory = 0;
diff --git
a/server/tserver/src/test/java/org/apache/accumulo/tserver/tablet/CompactableImplFileManagerTest.java
b/server/tserver/src/test/java/org/apache/accumulo/tserver/tablet/CompactableImplFileManagerTest.java
index 743ac8feca..aa9d818907 100644
---
a/server/tserver/src/test/java/org/apache/accumulo/tserver/tablet/CompactableImplFileManagerTest.java
+++
b/server/tserver/src/test/java/org/apache/accumulo/tserver/tablet/CompactableImplFileManagerTest.java
@@ -37,6 +37,8 @@ import java.util.concurrent.atomic.AtomicLong;
import java.util.stream.Collectors;
import java.util.stream.Stream;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.client.admin.compaction.CompactableFile;
import org.apache.accumulo.core.data.TableId;
import org.apache.accumulo.core.dataImpl.KeyExtent;
@@ -478,6 +480,7 @@ public class CompactableImplFileManagerTest {
}
+ @NotThreadSafe
static class TestFileManager extends CompactableImpl.FileManager {
public static final Duration SELECTION_EXPIRATION = Duration.ofMinutes(2);
diff --git a/test/pom.xml b/test/pom.xml
index 81d1ae4de5..cb4dbb45e5 100644
--- a/test/pom.xml
+++ b/test/pom.xml
@@ -42,6 +42,11 @@
<groupId>com.github.ben-manes.caffeine</groupId>
<artifactId>caffeine</artifactId>
</dependency>
+ <dependency>
+ <groupId>com.google.code.findbugs</groupId>
+ <artifactId>jsr305</artifactId>
+ <optional>true</optional>
+ </dependency>
<dependency>
<groupId>com.google.code.gson</groupId>
<artifactId>gson</artifactId>
diff --git
a/test/src/main/java/org/apache/accumulo/test/functional/ErrorThrowingIterator.java
b/test/src/main/java/org/apache/accumulo/test/functional/ErrorThrowingIterator.java
index 32dcdbf9ec..60c0e443ec 100644
---
a/test/src/main/java/org/apache/accumulo/test/functional/ErrorThrowingIterator.java
+++
b/test/src/main/java/org/apache/accumulo/test/functional/ErrorThrowingIterator.java
@@ -24,6 +24,8 @@ import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
+import javax.annotation.concurrent.NotThreadSafe;
+
import org.apache.accumulo.core.data.ByteSequence;
import org.apache.accumulo.core.data.Key;
import org.apache.accumulo.core.data.Range;
@@ -40,6 +42,7 @@ import com.google.common.base.Preconditions;
* Iterator used in tests *and* the test class must spawn a new MAC instance
for each test since the
* timesThrown variable is static.
*/
+@NotThreadSafe
public class ErrorThrowingIterator extends WrappingIterator {
private static final Logger log =
LoggerFactory.getLogger(ErrorThrowingIterator.class);