Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.ThreadLocalRandom;
import java.util.function.Function;
import java.util.stream.IntStream;
import org.openjdk.jmh.annotations.Benchmark;
Expand Down Expand Up @@ -105,7 +106,13 @@ public void multiThreaded(Blackhole blackhole, ThreadData threadData) {
generateKeys();
}

private static final Function<Object, Object> allocator = key -> new byte[1024];
// simulates a small, but non-trivial, value-construction cost
private static final Function<Object, Object> allocator =
key -> {
byte[] value = new byte[1024];
ThreadLocalRandom.current().nextBytes(value);
return value;
};

@Setup
public void setup() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,10 +9,9 @@
import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.lang.ref.ReferenceQueue;
import java.lang.ref.WeakReference;
import java.util.HashSet;
import java.util.Map;
import java.util.Set;
import java.util.Iterator;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Function;
import javax.annotation.Nullable;

Expand All @@ -26,22 +25,34 @@ public final class GlobalObjectStore {
/** Never allow more than this number of objects in the global store. */
private static final int GLOBAL_HARD_LIMIT = 100_000;

/** Threshold at which we age the young generation into the old generation. */
private static final int AGEING_THRESHOLD = GLOBAL_HARD_LIMIT / 2;

/** Temporarily allow more than this number of objects, but start removing old content. */
private static final int GLOBAL_SOFT_LIMIT = 50_000;
private static final int GLOBAL_SOFT_LIMIT = (GLOBAL_HARD_LIMIT + AGEING_THRESHOLD) / 2;

/** Threshold at which we start doing limited cleanup at the same time as put operations. */
private static final int INLINE_CLEANUP_THRESHOLD = 5_000;

/** Threshold at which we start sampling keys to track old content. */
private static final int OLD_KEYS_THRESHOLD = 512;
/** Bounds any inline eviction attempts. */
private static final int MAX_INLINE_EVICTION_ATTEMPTS = 10;

private static final Object staleEntriesLock = new Object();

private static final Map<StoreKey, Object> weakMap = new ConcurrentHashMap<>();
/** Tracks young and old content in the object store. */
private static final class Generations {
final ConcurrentHashMap<StoreKey, Object> young;
final ConcurrentHashMap<StoreKey, Object> old;

Generations(ConcurrentHashMap<StoreKey, Object> old) {
this.young = new ConcurrentHashMap<>(old.size());
this.old = old;
}
}

private static final Set<StoreKey> oldKeys = new HashSet<>();
private static volatile Generations generations = new Generations(new ConcurrentHashMap<>());

private static int previousEstimate = 0;
private static final AtomicBoolean ageing = new AtomicBoolean();

private GlobalObjectStore() {}

Expand All @@ -55,50 +66,21 @@ private GlobalObjectStore() {}
*/
public static int removeStaleEntries() {
synchronized (staleEntriesLock) {
Map<StoreKey, Object> weakMap = GlobalObjectStore.weakMap;
int estimatedSize = weakMap.size(); // capture size before any cleanup
StoreKey key;
while ((key = StoreKey.pollStaleKeys()) != null) {
if (weakMap.remove(key) != null) {
estimatedSize--;
}
removeEntry(key);
}
Comment thread
mcculls marked this conversation as resolved.

// The following code handles proactively removing old content in an attempt to guide the
// store below its soft limit. We remove older objects before recent additions, assuming
// that older objects are less likely to be used. For performance reasons this only runs
// after observed periods of growth or reduction, or if the store is near its hard limit.

// We deliberately avoid tracking exact age, and instead regularly sample keys to maintain
// a small set that we know are still alive after a couple of calls to removeStaleEntries.

if (Math.abs(estimatedSize - previousEstimate) > OLD_KEYS_THRESHOLD
|| estimatedSize >= (GLOBAL_HARD_LIMIT + GLOBAL_SOFT_LIMIT) / 2) {

if (estimatedSize >= GLOBAL_SOFT_LIMIT) {
// start proactively removing old content to keep growth in check
for (StoreKey oldKey : oldKeys) {
if (weakMap.remove(oldKey) != null) {
estimatedSize--;
}
}
oldKeys.clear();
} else {
// have any of the old previously sampled keys been collected?
oldKeys.removeIf(StoreKey::isStale);
}
Generations g = generations;

int refill = OLD_KEYS_THRESHOLD - oldKeys.size();
if (refill > 0) {
// sample of keys at this time, don't need strict age ordering
for (StoreKey sampleKey : weakMap.keySet()) {
if (oldKeys.add(sampleKey) && --refill == 0) {
break;
}
}
}
int estimatedSize = g.young.size() + g.old.size();

previousEstimate = estimatedSize;
// proactively remove old content when above the soft limit
Iterator<StoreKey> itr = g.old.keySet().iterator();
while (estimatedSize >= GLOBAL_SOFT_LIMIT && itr.hasNext()) {
itr.next();
itr.remove();
estimatedSize--;
}
Comment thread
mcculls marked this conversation as resolved.

return estimatedSize;
Expand All @@ -116,8 +98,14 @@ public static int removeStaleEntries() {
public static Object get(Object key, int storeId) {
LookupKey lookupKey = LookupKey.with(key, storeId);
try {
Generations g = generations;
//noinspection All: intentionally use lookup key without reference overhead
return weakMap.get(lookupKey);
Object value = g.young.get(lookupKey);
if (value == null) {
//noinspection All: intentionally use lookup key without reference overhead
value = g.old.get(lookupKey);
}
return value;
} finally {
lookupKey.reset();
}
Expand All @@ -133,8 +121,9 @@ public static Object get(Object key, int storeId) {
public static void put(Object key, int storeId, @Nullable Object value) {
if (value == null) {
remove(key, storeId);
} else if (checkCapacity()) {
weakMap.put(new StoreKey(key, storeId), value);
} else {
enforceCapacity();
generations.young.put(new StoreKey(key, storeId), value);
}
}

Expand All @@ -151,10 +140,11 @@ public static Object getOrPut(Object key, int storeId, @Nullable Object value) {
Object existing = get(key, storeId);
if (existing != null || value == null) {
return existing;
} else if (checkCapacity()) {
existing = weakMap.putIfAbsent(new StoreKey(key, storeId), value);
} else {
enforceCapacity();
existing = generations.young.putIfAbsent(new StoreKey(key, storeId), value);
return existing != null ? existing : value;
}
return existing != null ? existing : value;
}

/**
Expand All @@ -171,11 +161,11 @@ public static Object getOrCompute(Object key, int storeId, Function valueFunctio
Object existing = get(key, storeId);
if (existing != null) {
return existing;
} else if (checkCapacity()) {
return weakMap.computeIfAbsent(
new StoreKey(key, storeId), storeKey -> valueFunction.apply(storeKey.get()));
} else {
enforceCapacity();
return generations.young.computeIfAbsent(
new StoreKey(key, storeId), unused -> valueFunction.apply(key));
}
return valueFunction.apply(key);
}

/**
Expand All @@ -189,27 +179,72 @@ public static Object getOrCompute(Object key, int storeId, Function valueFunctio
public static Object remove(Object key, int storeId) {
LookupKey lookupKey = LookupKey.with(key, storeId);
try {
//noinspection All: intentionally use lookup key without reference overhead
return weakMap.remove(lookupKey);
return removeEntry(lookupKey);
} finally {
lookupKey.reset();
}
}

/**
* @return {@code true} if there is space to add new objects.
*/
private static boolean checkCapacity() {
int estimatedSize = weakMap.size();
if (estimatedSize > INLINE_CLEANUP_THRESHOLD) {
// periodic cleanup may not be enough, start performing inline cleanup
StoreKey staleKey = StoreKey.pollStaleKeys();
if (staleKey == null) {
return estimatedSize < GLOBAL_HARD_LIMIT;
private static Object removeEntry(Object key) {
Generations g = generations;
//noinspection All: intentionally use lookup key without reference overhead
Object youngValue = g.young.remove(key);
//noinspection All: intentionally use lookup key without reference overhead
Object oldValue = g.old.remove(key);
return youngValue != null ? youngValue : oldValue;
}

private static void enforceCapacity() {
Generations g = generations;
int youngSize = g.young.size();
int totalSize = youngSize + g.old.size();
if (totalSize < INLINE_CLEANUP_THRESHOLD) {
return; // skip inline eviction
}

// attempt a single stale eviction
StoreKey staleKey;
int attempts = MAX_INLINE_EVICTION_ATTEMPTS;
while (attempts-- > 0 && (staleKey = StoreKey.pollStaleKeys()) != null) {
if (removeEntry(staleKey) != null) {
return;
}
}

// when the young generation maxes out, age it so it becomes old
if (youngSize >= AGEING_THRESHOLD && ageGenerations(g)) {
return;
}

// attempt a single old eviction
if (totalSize >= GLOBAL_HARD_LIMIT) {
evictOldKey();
}
}

private static boolean ageGenerations(Generations g) {
if (ageing.compareAndSet(false, true)) {
try {
if (g == generations) {
generations = new Generations(g.young);
return true;
}
} finally {
ageing.set(false);
}
}
return false;
}

private static void evictOldKey() {
Generations g = generations;
int attempts = MAX_INLINE_EVICTION_ATTEMPTS;
Iterator<StoreKey> itr = g.old.keySet().iterator();
while (attempts-- > 0 && itr.hasNext()) {
if (g.old.remove(itr.next()) != null) {
return;
}
weakMap.remove(staleKey);
}
return true;
}

/** Key used to weakly associate a non-injected key and store-id with a value. */
Expand All @@ -231,10 +266,6 @@ static StoreKey pollStaleKeys() {
return (StoreKey) staleKeys.poll();
}

boolean isStale() {
return get() == null;
}

@Override
public int hashCode() {
return hash;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -219,6 +219,45 @@ void removeStaleEntriesIsIdempotent() {
assertEquals("stable", store.get(key));
}

// --- generational capacity / eviction ---

// Mirrors the private thresholds in GlobalObjectStore; kept here so the test intent is clear
// without exposing internals. If those thresholds change this test should be revisited.
private static final int AGEING_THRESHOLD = 50_000;
private static final int GLOBAL_SOFT_LIMIT = 75_000;
private static final int GLOBAL_HARD_LIMIT = 100_000;

@Test
void sustainedInsertionIsBoundedByAgeingAndSoftLimitTrim() {
ObjectStore<Object, Integer> capStore =
ObjectStore.of("test.Capacity.Key", "test.Capacity.Value");

// Insert enough distinct, strongly-referenced keys to drive the young generation past
// AGEING_THRESHOLD several times over, landing mid-cycle (comfortably above the soft
// limit) so both inline ageing and removeStaleEntries' soft-limit trim get exercised.
int totalInserts = (AGEING_THRESHOLD * 4) + 40_000;
List<Object> keys = new ArrayList<>(totalInserts);
for (int i = 0; i < totalInserts; i++) {
Object key = new Object();
keys.add(key);
capStore.put(key, i);
}

// Inline enforceCapacity keeps young+old from ever exceeding the hard limit by ageing
// young into old before that point is reached, so recently inserted keys must still be
// retrievable even after hundreds of thousands of insertions.
Object lastKey = keys.get(keys.size() - 1);
assertEquals(totalInserts - 1, capStore.get(lastKey));

int finalSize = ObjectStore.removeStaleEntries();
assertTrue(
finalSize < GLOBAL_HARD_LIMIT,
"Sustained insertion should have triggered eviction rather than unbounded growth");
assertTrue(
finalSize <= GLOBAL_SOFT_LIMIT,
"removeStaleEntries should trim content back to the soft limit, observed " + finalSize);
}

// --- concurrency ---

@Test
Expand Down
Loading