diff --git a/commons/src/main/java/org/archive/bdb/BdbModule.java b/commons/src/main/java/org/archive/bdb/BdbModule.java index ddb4591a..a49eb7a5 100644 --- a/commons/src/main/java/org/archive/bdb/BdbModule.java +++ b/commons/src/main/java/org/archive/bdb/BdbModule.java @@ -47,6 +47,8 @@ import org.archive.checkpointing.Checkpointable; import org.archive.checkpointing.RecoverAction; import org.archive.spring.ConfigPath; import org.archive.util.CachedBdbMap; +import org.archive.util.ObjectIdentityBdbCache; +import org.archive.util.ObjectIdentityCache; import org.archive.util.bdbje.EnhancedEnvironment; import org.springframework.context.Lifecycle; @@ -230,10 +232,9 @@ Serializable, Closeable { private transient StoredClassCatalog classCatalog; - @SuppressWarnings("unchecked") - private Map bigMaps = - new ConcurrentHashMap(); - + private Map oiCaches = + new ConcurrentHashMap(); + private Map databases = new ConcurrentHashMap(); @@ -361,15 +362,24 @@ Serializable, Closeable { } - public CachedBdbMap getBigMap(String dbName, boolean recycle, + /** + * Get a CachedBdbMap, backed by a BDB Database of the given name, + * with the given key and value class types. If 'recycle' is true, + * reuse values already in the database; otherwise start with an + * empty map. + * + * @param + * @param + * @param dbName + * @param recycle + * @param key + * @param value + * @return + * @throws DatabaseException + */ + public CachedBdbMap getCBMMap(String dbName, boolean recycle, Class key, Class value) throws DatabaseException { - @SuppressWarnings("unchecked") - CachedBdbMap r = bigMaps.get(dbName); - if (r != null) { - return r; - } - if (!recycle) { try { bdbEnvironment.truncateDatabase(null, dbName, false); @@ -377,20 +387,78 @@ Serializable, Closeable { // ignored } } - - r = new CachedBdbMap(dbName); - + CachedBdbMap r = new CachedBdbMap(dbName); r.initialize(bdbEnvironment, key, value, classCatalog); - bigMaps.put(dbName, r); return r; } + /** + * Get an ObjectIdentityBdbCache, backed by a BDB Database of the + * given name, with the given value class type. If 'recycle' is true, + * reuse values already in the database; otherwise start with an + * empty cache. + * + * @param + * @param dbName + * @param recycle + * @param valueClass + * @return + * @throws DatabaseException + */ + public ObjectIdentityBdbCache getOIBCCache(String dbName, boolean recycle, + Class valueClass) + throws DatabaseException { + if (!recycle) { + try { + bdbEnvironment.truncateDatabase(null, dbName, false); + } catch (DatabaseNotFoundException e) { + // ignored + } + } + ObjectIdentityBdbCache oic = new ObjectIdentityBdbCache(); + oic.initialize(bdbEnvironment, dbName, valueClass, classCatalog); + return oic; + } + /** controls which alternate ObjectIdentityCache implementation to use */ + private static boolean USE_OIBC = true; + + + + /** + * Get an ObjectIdentityCache, backed by a BDB Database of the given + * name, with objects of the given valueClass type. If 'recycle' is + * true, reuse values already in the database; otherwise start with + * an empty cache. + * + * @param + * @param dbName + * @param recycle + * @param valueClass + * @return + * @throws DatabaseException + */ + public ObjectIdentityCache getObjectCache(String dbName, boolean recycle, + Class valueClass) + throws DatabaseException { + @SuppressWarnings("unchecked") + ObjectIdentityCache oic = oiCaches.get(dbName); + if(oic!=null) { + return oic; + } + if(USE_OIBC) { + oic = getOIBCCache(dbName, recycle, valueClass); + } else { + oic = getCBMMap(dbName, recycle, String.class, valueClass); + } + oiCaches.put(dbName, oic); + return oic; + } + private void writeObject(ObjectOutputStream out) throws IOException { out.defaultWriteObject(); } - // TODO:FIXME: restore functionality @SuppressWarnings("unchecked") private void readObject(ObjectInputStream in) @@ -403,15 +471,15 @@ Serializable, Closeable { } try { setUp(getDir().getFile(), getCachePercent(), false, getUseSharedCache()); - for (CachedBdbMap map: bigMaps.values()) { - map.initialize( - this.bdbEnvironment, - null, - null, +// for (CachedBdbMap map: bigMaps.values()) { +// map.initialize( +// this.bdbEnvironment, +// null, +// null, // map.getKeyClass(), // map.getValueClass(), - this.classCatalog); - } +// this.classCatalog); +// } for (DatabasePlusConfig dpc: databases.values()) { if (!(dpc.config instanceof SecondaryBdbConfig)) { dpc.database = bdbEnvironment.openDatabase(null, @@ -442,9 +510,9 @@ Serializable, Closeable { if (checkpointCopyLogs) { actions.add(new BdbRecover(getDir().getFile().getAbsolutePath())); } - // First sync bigMaps - for (Map.Entry me: bigMaps.entrySet()) { - me.getValue().sync(); + // First sync objectCaches + for (ObjectIdentityCache oic : oiCaches.values()) { + oic.sync(); } EnvironmentConfig envConfig; @@ -600,10 +668,13 @@ Serializable, Closeable { if (classCatalog == null) { return; } - for (Map.Entry me: bigMaps.entrySet()) try { - me.getValue().close(); - } catch (Exception e) { - LOGGER.log(Level.SEVERE, "Error closing bigMap " + me.getKey(), e); + + for(ObjectIdentityCache cache : oiCaches.values()) { + try { + cache.close(); + } catch (Exception e) { + LOGGER.log(Level.SEVERE, "Error closing oiCache " + cache, e); + } } List dbNames = new ArrayList(databases.keySet()); diff --git a/commons/src/main/java/org/archive/util/CachedBdbMap.java b/commons/src/main/java/org/archive/util/CachedBdbMap.java index b4cb67b7..52fb74c7 100644 --- a/commons/src/main/java/org/archive/util/CachedBdbMap.java +++ b/commons/src/main/java/org/archive/util/CachedBdbMap.java @@ -21,7 +21,6 @@ package org.archive.util; import java.io.Closeable; import java.io.File; -import java.io.IOException; import java.io.Serializable; import java.lang.ref.PhantomReference; import java.lang.ref.Reference; @@ -42,6 +41,9 @@ import java.util.concurrent.atomic.AtomicLong; import java.util.logging.Level; import java.util.logging.Logger; +import org.archive.crawler.framework.CrawlController; +import org.archive.modules.net.ServerCache; + import com.sleepycat.bind.EntryBinding; import com.sleepycat.bind.serial.SerialBinding; import com.sleepycat.bind.serial.StoredClassCatalog; @@ -92,7 +94,7 @@ import com.sleepycat.je.Environment; * */ public class CachedBdbMap extends AbstractMap -implements ConcurrentMap, Serializable, Closeable { +implements ConcurrentMap, ObjectIdentityCache, Serializable, Closeable { private static final long serialVersionUID = -8655539411367047332L; private static final Logger logger = @@ -478,7 +480,7 @@ implements ConcurrentMap, Serializable, Closeable { return environment.openDatabase(null, dbName, dbConfig); } - public synchronized void close() throws IOException { + public synchronized void close() { // Close out my bdb db. if (this.db != null) { try { @@ -512,6 +514,22 @@ implements ConcurrentMap, Serializable, Closeable { // maintain identity guarantees, so skipping throw new UnsupportedOperationException(); } + + /** + * ObjectIdentityCache get-or-atomic-create method. + */ + public V getOrUse(K key, Supplier supplierOrNull) { + V val = get(key); + if(val!=null || supplierOrNull == null) { + return val; + } + val = supplierOrNull.get(); + V prevVal = putIfAbsent(key, val); + if(prevVal!=null) { + return prevVal; + } + return val; + } public V get(final Object object) { K key = toKey(object); diff --git a/commons/src/main/java/org/archive/util/bdbje/EnhancedEnvironment.java b/commons/src/main/java/org/archive/util/bdbje/EnhancedEnvironment.java index 821e833a..86b488f6 100644 --- a/commons/src/main/java/org/archive/util/bdbje/EnhancedEnvironment.java +++ b/commons/src/main/java/org/archive/util/bdbje/EnhancedEnvironment.java @@ -82,6 +82,21 @@ public class EnhancedEnvironment extends Environment { super.close(); } - - + /** + * Create a temporary test environment in the given directory. + * @param dir target directory + * @return EnhancedEnvironment + */ + public static EnhancedEnvironment getTestEnvironment(File dir) { + EnvironmentConfig envConfig = new EnvironmentConfig(); + envConfig.setAllowCreate(true); + envConfig.setTransactional(false); + EnhancedEnvironment env; + try { + env = new EnhancedEnvironment(dir, envConfig); + } catch (DatabaseException e) { + throw new RuntimeException(e); + } + return env; + } } diff --git a/commons/src/test/java/org/archive/settings/file/BdbModuleTest.java b/commons/src/test/java/org/archive/settings/file/BdbModuleTest.java index eda23d90..8589a65a 100644 --- a/commons/src/test/java/org/archive/settings/file/BdbModuleTest.java +++ b/commons/src/test/java/org/archive/settings/file/BdbModuleTest.java @@ -77,7 +77,7 @@ public class BdbModuleTest extends TmpDirTestCase { config.setAllowCreate(true); bdb.openDatabase("testOpen", config, false); - Map testData = bdb.getBigMap("testData", false, + Map testData = bdb.getCBMMap("testData", false, String.class, String.class); for (int i = 0; i < 1000; i++) { testData.put(String.valueOf(i), String.valueOf(i * 2)); diff --git a/commons/src/test/java/org/archive/util/CachedBdbMapTest.java b/commons/src/test/java/org/archive/util/CachedBdbMapTest.java index 62c15565..ed0b6163 100644 --- a/commons/src/test/java/org/archive/util/CachedBdbMapTest.java +++ b/commons/src/test/java/org/archive/util/CachedBdbMapTest.java @@ -49,7 +49,7 @@ public class CachedBdbMapTest extends TmpDirTestCase { bdb.getDir().setBase(null); bdb.getDir().setPath(envDir.getAbsolutePath()); bdb.start(); - this.cache = bdb.getBigMap( + this.cache = bdb.getCBMMap( this.getClass().getName(), false, String.class, HashMap.class); } diff --git a/engine/src/main/java/org/archive/crawler/frontier/BdbFrontier.java b/engine/src/main/java/org/archive/crawler/frontier/BdbFrontier.java index e5d75d34..0ecb24d6 100644 --- a/engine/src/main/java/org/archive/crawler/frontier/BdbFrontier.java +++ b/engine/src/main/java/org/archive/crawler/frontier/BdbFrontier.java @@ -41,6 +41,7 @@ import org.archive.checkpointing.Checkpointable; import org.archive.checkpointing.RecoverAction; import org.archive.modules.CrawlURI; import org.archive.queue.StoredQueue; +import org.archive.util.Supplier; import org.springframework.beans.factory.annotation.Autowired; import com.sleepycat.collections.StoredIterator; @@ -148,20 +149,17 @@ implements Serializable, Checkpointable { * @param classKey key to look for * @return the found WorkQueue */ - protected WorkQueue getQueueFor(String classKey) { - - WorkQueue wq = (WorkQueue)allQueues.get(classKey); - if(wq == null) { - String qKey = new String(classKey); // ensure private minimal key - wq = new BdbWorkQueue(qKey, this); - wq.setTotalBudget(getQueueTotalBudget()); - getQueuePrecedencePolicy().queueCreated(wq); - WorkQueue prevVal = allQueues.putIfAbsent(qKey, wq); - if(prevVal!=null) { - // lost race; prefer earlier object - wq = prevVal; - } - } + protected WorkQueue getQueueFor(final String classKey) { + WorkQueue wq = allQueues.getOrUse( + classKey, + new Supplier() { + public WorkQueue get() { + String qKey = new String(classKey); // ensure private minimal key + WorkQueue q = new BdbWorkQueue(qKey, BdbFrontier.this); + q.setTotalBudget(getQueueTotalBudget()); + getQueuePrecedencePolicy().queueCreated(q); + return q; + }}); return wq; } @@ -228,8 +226,7 @@ implements Serializable, Checkpointable { @Override protected void initAllQueues() throws DatabaseException { - this.allQueues = bdb.getBigMap("allqueues", false, - String.class, WorkQueue.class); + this.allQueues = bdb.getObjectCache("allqueues", false, WorkQueue.class); if (logger.isLoggable(Level.FINE)) { Iterator i = this.allQueues.keySet().iterator(); try { diff --git a/engine/src/main/java/org/archive/crawler/frontier/WorkQueueFrontier.java b/engine/src/main/java/org/archive/crawler/frontier/WorkQueueFrontier.java index 17a934d8..7ea39f12 100644 --- a/engine/src/main/java/org/archive/crawler/frontier/WorkQueueFrontier.java +++ b/engine/src/main/java/org/archive/crawler/frontier/WorkQueueFrontier.java @@ -40,8 +40,6 @@ import java.util.Map; import java.util.Queue; import java.util.SortedMap; import java.util.concurrent.BlockingQueue; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; import java.util.concurrent.DelayQueue; import java.util.concurrent.Delayed; import java.util.concurrent.TimeUnit; @@ -63,6 +61,8 @@ import org.archive.crawler.frontier.precedence.UriPrecedencePolicy; import org.archive.modules.CrawlURI; import org.archive.spring.KeyedProperties; import org.archive.util.ArchiveUtils; +import org.archive.util.ObjectIdentityCache; +import org.archive.util.ObjectIdentityMemCache; import org.archive.util.Transform; import org.archive.util.Transformer; import org.springframework.beans.BeansException; @@ -215,7 +215,7 @@ ApplicationContextAware { /** All known queues. */ - protected ConcurrentMap allQueues = null; + protected ObjectIdentityCache allQueues = null; // of classKey -> ClassKeyQueue /** @@ -287,7 +287,7 @@ ApplicationContextAware { && getQueueAssignmentPolicy().maximumNumberOfKeys() <= MAX_QUEUES_TO_HOLD_ALLQUEUES_IN_MEMORY) { this.allQueues = - new ConcurrentHashMap(701, .9f, 100); + new ObjectIdentityMemCache(701, .9f, 100); } else { this.initAllQueues(); } @@ -316,19 +316,15 @@ ApplicationContextAware { // references. if (this.uriUniqFilter != null) { this.uriUniqFilter.close(); -// this.alreadyIncluded = null; } - -// this.queueAssignmentPolicy = null; - + try { closeQueue(); } catch (IOException e) { - // FIXME exception handling - e.printStackTrace(); + logger.log(Level.WARNING,"closeQueue problem",e); } - this.allQueues.clear(); + this.allQueues.close(); } /** @@ -561,7 +557,7 @@ ApplicationContextAware { // TODO: Only do this when necessary. - Object key = getRetiredQueues().poll(); + String key = getRetiredQueues().poll(); while (key != null) { WorkQueue q = (WorkQueue)this.allQueues.get(key); if(q != null) { @@ -754,7 +750,7 @@ ApplicationContextAware { Queue inactiveQueues = inactiveQueuesByPrecedence.get( targetPrecedence); - Object key = inactiveQueues.poll(); + String key = inactiveQueues.poll(); assert key != null : "empty precedence queue in map"; if(inactiveQueues.isEmpty()) { @@ -1282,7 +1278,7 @@ ApplicationContextAware { q = ((DelayedWorkQueue)obj).getWorkQueue(); } else { try { - q = (WorkQueue)this.allQueues.get(obj); + q = this.allQueues.get((String)obj); } catch (ClassCastException cce) { logger.log(Level.SEVERE,"not convertible to workqueue:"+obj,cce); q = null; @@ -1474,9 +1470,9 @@ ApplicationContextAware { if (obj == null) { continue; } - q = (obj instanceof WorkQueue)? - (WorkQueue)obj: - (WorkQueue)this.allQueues.get(obj); + q = (obj instanceof WorkQueue) + ? (WorkQueue)obj + : this.allQueues.get((String)obj); if(q == null) { w.print("WARNING: No report for queue "+obj); } diff --git a/engine/src/main/java/org/archive/crawler/reporting/StatisticsTracker.java b/engine/src/main/java/org/archive/crawler/reporting/StatisticsTracker.java index c766edc7..d1a55502 100644 --- a/engine/src/main/java/org/archive/crawler/reporting/StatisticsTracker.java +++ b/engine/src/main/java/org/archive/crawler/reporting/StatisticsTracker.java @@ -50,8 +50,8 @@ import org.archive.crawler.event.StatSnapshotEvent; import org.archive.crawler.framework.CrawlController; import org.archive.crawler.framework.Engine; import org.archive.crawler.util.CrawledBytesHistotable; -import org.archive.crawler.util.TopNSet; -import org.archive.modules.CrawlURI; +import org.archive.crawler.util.TopNSet; +import org.archive.modules.CrawlURI; import org.archive.modules.ProcessorURI; import org.archive.modules.net.ServerCache; import org.archive.modules.net.ServerCacheUtil; @@ -60,7 +60,10 @@ import org.archive.modules.seeds.SeedModule; import org.archive.spring.ConfigPath; import org.archive.util.ArchiveUtils; import org.archive.util.MimetypeUtils; +import org.archive.util.ObjectIdentityCache; +import org.archive.util.ObjectIdentityMemCache; import org.archive.util.PaddingStringBuffer; +import org.archive.util.Supplier; import org.springframework.beans.BeansException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ApplicationContext; @@ -230,6 +233,13 @@ public class StatisticsTracker FrontierSummaryReport.class, ToeThreadsReport.class, }; + + /** reusable Supplier for initial zero AtomicLong instances */ + private static final Supplier ATOMIC_ZERO_SUPPLIER = + new Supplier() { + public AtomicLong get() { + return new AtomicLong(0); + }}; /** * The interval between writing progress information to log. @@ -285,28 +295,27 @@ public class StatisticsTracker // TODO: fortify these against key explosion with bigmaps like other tallies /** Keep track of the file types we see (mime type -> count) */ - protected ConcurrentMap mimeTypeDistribution - = new ConcurrentHashMap(); - protected ConcurrentMap mimeTypeBytes - = new ConcurrentHashMap(); + protected ObjectIdentityCache mimeTypeDistribution + = new ObjectIdentityMemCache(); + protected ObjectIdentityCache mimeTypeBytes + = new ObjectIdentityMemCache(); /** Keep track of fetch status codes */ - protected ConcurrentMap statusCodeDistribution - = new ConcurrentHashMap(); + protected ObjectIdentityCache statusCodeDistribution + = new ObjectIdentityMemCache(); /** Keep track of hosts. */ - protected ConcurrentMap hostsDistribution = - new ConcurrentHashMap(); // temp dummy - protected ConcurrentMap hostsBytes = - new ConcurrentHashMap(); // temp dummy - protected ConcurrentMap hostsLastFinished = - new ConcurrentHashMap(); // temp dummy + protected ObjectIdentityCache hostsDistribution = + new ObjectIdentityMemCache(); // temp dummy + protected ObjectIdentityCache hostsBytes = + new ObjectIdentityMemCache(); // temp dummy + protected ObjectIdentityCache hostsLastFinished = + new ObjectIdentityMemCache(); // temp dummy /** Keep track of URL counts per host per seed */ - protected - ConcurrentMap> sourceHostDistribution = - new ConcurrentHashMap>(); // temp dummy; + protected ObjectIdentityCache> sourceHostDistribution = + new ObjectIdentityMemCache>(); // temp dummy; /* Keep track of 'top' hosts for live reports */ protected TopNSet hostsDistributionTop; @@ -316,8 +325,8 @@ public class StatisticsTracker /** * Record of seeds and latest results */ - protected ConcurrentMap processedSeedsRecords = - new ConcurrentHashMap(); + protected ObjectIdentityCache processedSeedsRecords = + new ObjectIdentityMemCache(); long seedsTotal = -1; long seedsCrawled = -1; @@ -338,16 +347,16 @@ public class StatisticsTracker public void start() { isRunning = true; try { - this.sourceHostDistribution = bdb.getBigMap("sourceHostDistribution", - false, String.class, ConcurrentMap.class); - this.hostsDistribution = bdb.getBigMap("hostsDistribution", - false, String.class, AtomicLong.class); - this.hostsBytes = bdb.getBigMap("hostsBytes", false, String.class, + this.sourceHostDistribution = bdb.getObjectCache("sourceHostDistribution", + false, ConcurrentMap.class); + this.hostsDistribution = bdb.getObjectCache("hostsDistribution", + false, AtomicLong.class); + this.hostsBytes = bdb.getObjectCache("hostsBytes", false, AtomicLong.class); - this.hostsLastFinished = bdb.getBigMap("hostsLastFinished", - false, String.class, AtomicLong.class); - this.processedSeedsRecords = bdb.getBigMap("processedSeedsRecords", - false, String.class, SeedRecord.class); + this.hostsLastFinished = bdb.getObjectCache("hostsLastFinished", + false, AtomicLong.class); + this.processedSeedsRecords = bdb.getObjectCache("processedSeedsRecords", + false, SeedRecord.class); this.hostsDistributionTop = new TopNSet(getLiveHostReportSize()); this.hostsBytesTop = new TopNSet(getLiveHostReportSize()); @@ -520,7 +529,7 @@ public class StatisticsTracker * Note: All the values are wrapped with a {@link AtomicLong AtomicLong} * @return mimeTypeDistribution */ - public ConcurrentMap getFileDistribution() { + public ObjectIdentityCache getFileDistribution() { return mimeTypeDistribution; } @@ -539,6 +548,42 @@ public class StatisticsTracker incrementMapCount(map,key,1); } + /** + * Increment a counter for a key in a given cache. Used for various + * aggregate data. + * + * @param cache the ObjectIdentityCache + * @param key The key for the counter to be incremented, if it does not + * exist it will be added (set to 1). If null it will + * increment the counter "unknown". + */ + protected static void incrementCacheCount(ObjectIdentityCache cache, + String key) { + incrementCacheCount(cache,key,1); + } + /** + * Increment a counter for a key in a given cache by an arbitrary amount. + * Used for various aggregate data. The increment amount can be negative. + * + * + * @param cache + * The ObjectIdentityCache + * @param key + * The key for the counter to be incremented, if it does not exist + * it will be added (set to equal to increment). + * If null it will increment the counter "unknown". + * @param increment + * The amount to increment counter related to the key. + */ + protected static void incrementCacheCount(ObjectIdentityCache cache, + String key, long increment) { + if (key == null) { + key = "unknown"; + } + AtomicLong lw = cache.getOrUse(key, ATOMIC_ZERO_SUPPLIER); + lw.addAndGet(increment); + } + /** * Increment a counter for a key in a given HashMap by an arbitrary amount. * Used for various aggregate data. The increment amount can be negative. @@ -615,8 +660,46 @@ public class StatisticsTracker } /** - * Return a HashMap representing the distribution of status codes for - * successfully fetched curis, as represented by a hashmap where key -> + * Sort the entries of the given ObjectIdentityCache in descending order by their + * values, which must be longs wrapped with AtomicLong. + *

+ * Elements are sorted by value from largest to smallest. Equal values are + * sorted in an arbitrary, but consistent manner by their keys. Only items + * with identical value and key are considered equal. + * + * If the passed-in map requires access to be synchronized, the caller + * should ensure this synchronization. + * + * @param mapOfAtomicLongValues + * Assumes values are wrapped with AtomicLong. + * @return a sorted set containing the same elements as the map. + */ + public TreeMap getReverseSortedCopy( + final ObjectIdentityCache cacheOfAtomicLongValues) { + TreeMap sortedMap = + new TreeMap(new Comparator() { + public int compare(String e1, String e2) { + long firstVal = cacheOfAtomicLongValues.get(e1).get(); + long secondVal = cacheOfAtomicLongValues.get(e2).get(); + if (firstVal < secondVal) { + return 1; + } + if (secondVal < firstVal) { + return -1; + } + // If the values are the same, sort by keys. + return e1.compareTo(e2); + } + }); + for(String key : cacheOfAtomicLongValues.keySet()) { + sortedMap.put(key, cacheOfAtomicLongValues.get(key)); + } + return sortedMap; + } + + /** + * Return a objectCache representing the distribution of status codes for + * successfully fetched curis, as represented by a cache where key -> * val represents (string)code -> (integer)count. * * Note: All the values are wrapped with a @@ -624,7 +707,7 @@ public class StatisticsTracker * * @return statusCodeDistribution */ - public ConcurrentMap getStatusCodeDistribution() { + public ObjectIdentityCache getStatusCodeDistribution() { return statusCodeDistribution; } @@ -638,15 +721,7 @@ public class StatisticsTracker * -1 will be returned. */ public AtomicLong getHostLastFinished(String host){ - AtomicLong fini = hostsLastFinished.get(host); - if(fini==null) { - String hkey = new String(host); // ensure private minimal key - fini = new AtomicLong(-1); - AtomicLong prevVal = hostsLastFinished.putIfAbsent(hkey, fini); - if(prevVal!=null) { - fini = prevVal; - } - } + AtomicLong fini = hostsLastFinished.getOrUse(host, ATOMIC_ZERO_SUPPLIER); return fini; } @@ -682,22 +757,20 @@ public class StatisticsTracker } /** - * If the curi is a seed, we update the processedSeeds table. + * If the curi is a seed, we update the processedSeeds cache. * * @param curi The CrawlURI that may be a seed. * @param disposition The disposition of the CrawlURI. */ - private void handleSeed(CrawlURI curi, String disposition) { + private void handleSeed(final CrawlURI curi, final String disposition) { if(curi.isSeed()){ - SeedRecord sr = processedSeedsRecords.get(curi.getURI()); - if(sr==null) { - sr = new SeedRecord(curi, disposition); - SeedRecord prevVal = processedSeedsRecords.putIfAbsent(sr.getUri(), sr); - if(prevVal!=null) { - sr = prevVal; - sr.updateWith(curi,disposition); - } - } + SeedRecord sr = processedSeedsRecords.getOrUse( + curi.getURI(), + new Supplier() { + public SeedRecord get() { + return new SeedRecord(curi, disposition); + }}); + sr.updateWith(curi,disposition); } } @@ -707,13 +780,13 @@ public class StatisticsTracker crawledBytes.accumulate(curi); // Save status codes - incrementMapCount(statusCodeDistribution, + incrementCacheCount(statusCodeDistribution, Integer.toString(curi.getFetchStatus())); // Save mime types String mime = MimetypeUtils.truncate(curi.getContentType()); - incrementMapCount(mimeTypeDistribution, mime); - incrementMapCount(mimeTypeBytes, mime, curi.getContentSize()); + incrementCacheCount(mimeTypeDistribution, mime); + incrementCacheCount(mimeTypeBytes, mime, curi.getContentSize()); // Save hosts stats. ServerCache sc = serverCache; @@ -731,25 +804,22 @@ public class StatisticsTracker protected void saveSourceStats(String source, String hostname) { synchronized(sourceHostDistribution) { ConcurrentMap hostUriCount = - sourceHostDistribution.get(source); - if (hostUriCount == null) { - hostUriCount = new ConcurrentHashMap(); - ConcurrentMap prevVal = - sourceHostDistribution.putIfAbsent(source, hostUriCount); - if(prevVal != null) { - hostUriCount = prevVal; - } - } + sourceHostDistribution.getOrUse( + source, + new Supplier>() { + public ConcurrentMap get() { + return new ConcurrentHashMap(); + }}); incrementMapCount(hostUriCount, hostname); } } protected void saveHostStats(String hostname, long size) { - incrementMapCount(hostsDistribution, hostname); + incrementCacheCount(hostsDistribution, hostname); hostsDistributionTop.update( hostname, getReportValue(hostsDistribution, hostname)); - incrementMapCount(hostsBytes, hostname, size); + incrementCacheCount(hostsBytes, hostname, size); hostsBytesTop.update(hostname, getReportValue(hostsBytes, hostname)); @@ -909,7 +979,7 @@ public class StatisticsTracker logNote("CRAWL CHECKPOINTING TO " + cpDir.toString()); } - private long getReportValue(Map map, String key) { + private long getReportValue(ObjectIdentityCache map, String key) { if (key == null) { return -1; } diff --git a/modules/src/main/java/org/archive/modules/fetcher/DefaultServerCache.java b/modules/src/main/java/org/archive/modules/fetcher/DefaultServerCache.java index 57036671..4371d3fc 100644 --- a/modules/src/main/java/org/archive/modules/fetcher/DefaultServerCache.java +++ b/modules/src/main/java/org/archive/modules/fetcher/DefaultServerCache.java @@ -21,8 +21,6 @@ package org.archive.modules.fetcher; import java.io.Closeable; import java.io.Serializable; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; import java.util.logging.Logger; import org.apache.commons.collections.Closure; @@ -31,6 +29,9 @@ import org.archive.modules.net.CrawlHost; import org.archive.modules.net.CrawlServer; import org.archive.modules.net.ServerCache; import org.archive.net.UURI; +import org.archive.util.ObjectIdentityCache; +import org.archive.util.ObjectIdentityMemCache; +import org.archive.util.Supplier; /** @@ -49,27 +50,27 @@ public class DefaultServerCache implements ServerCache, Closeable, Serializable * hostname[:port] -> CrawlServer. * Set in the initialization. */ - protected ConcurrentMap servers = null; + protected ObjectIdentityCache servers = null; /** * hostname -> CrawlHost. * Set in the initialization. */ - protected ConcurrentMap hosts = null; + protected ObjectIdentityCache hosts = null; /** * Constructor. */ public DefaultServerCache() { this( - new ConcurrentHashMap(), - new ConcurrentHashMap()); + new ObjectIdentityMemCache(), + new ObjectIdentityMemCache()); } - public DefaultServerCache(ConcurrentMap servers, - ConcurrentMap hosts) { + public DefaultServerCache(ObjectIdentityCache servers, + ObjectIdentityCache hosts) { this.servers = servers; this.hosts = hosts; } @@ -79,16 +80,14 @@ public class DefaultServerCache implements ServerCache, Closeable, Serializable * @param serverKey Server name we're to return server for. * @return CrawlServer instance that matches the passed server name. */ - public synchronized CrawlServer getServerFor(String serverKey) { - CrawlServer cserver = servers.get(serverKey); - if(cserver==null) { - String skey = new String(serverKey); // ensure private minimal key - cserver = new CrawlServer(skey); - CrawlServer prevVal = servers.putIfAbsent(skey, cserver); - if(prevVal!=null) { - cserver = prevVal; - } - } + public synchronized CrawlServer getServerFor(final String serverKey) { + CrawlServer cserver = servers.getOrUse( + serverKey, + new Supplier() { + public CrawlServer get() { + String skey = new String(serverKey); // ensure private minimal key + return new CrawlServer(skey); + }}); return cserver; } @@ -121,19 +120,17 @@ public class DefaultServerCache implements ServerCache, Closeable, Serializable * @param hostname Host name we're to return Host for. * @return CrawlHost instance that matches the passed Host name. */ - public synchronized CrawlHost getHostFor(String hostname) { + public synchronized CrawlHost getHostFor(final String hostname) { if (hostname == null || hostname.length() == 0) { return null; } - CrawlHost host = hosts.get(hostname); - if(host == null) { - String hkey = new String(hostname); // ensure private minimal key - host = new CrawlHost(hkey); - CrawlHost prevVal = hosts.putIfAbsent(hkey, host); - if(prevVal!=null) { - host = prevVal; - } - } + CrawlHost host = hosts.getOrUse( + hostname, + new Supplier() { + public CrawlHost get() { + String hkey = new String(hostname); // ensure private minimal key + return new CrawlHost(hkey); + }}); return host; } @@ -175,11 +172,11 @@ public class DefaultServerCache implements ServerCache, Closeable, Serializable if (this.hosts != null) { // If we're using a bdb bigmap, the call to clear will // close down the bdb database. - this.hosts.clear(); + this.hosts.close(); this.hosts = null; } if (this.servers != null) { - this.servers.clear(); + this.servers.close(); this.servers = null; } } diff --git a/modules/src/main/java/org/archive/modules/net/BdbServerCache.java b/modules/src/main/java/org/archive/modules/net/BdbServerCache.java index 1ff270c7..a810b508 100644 --- a/modules/src/main/java/org/archive/modules/net/BdbServerCache.java +++ b/modules/src/main/java/org/archive/modules/net/BdbServerCache.java @@ -50,8 +50,8 @@ implements Lifecycle { return; } try { - this.servers = bdb.getBigMap("servers", false, String.class, CrawlServer.class); - this.hosts = bdb.getBigMap("hosts", false, String.class, CrawlHost.class); + this.servers = bdb.getObjectCache("servers", false, CrawlServer.class); + this.hosts = bdb.getObjectCache("hosts", false, CrawlHost.class); } catch (DatabaseException e) { throw new IllegalStateException(e); }