From c196a1de77a9c42c3600a938eed38ee84dd77663 Mon Sep 17 00:00:00 2001 From: gojomo Date: Tue, 4 Aug 2009 00:06:22 +0000 Subject: [PATCH] [HER-1658] CachedBdbMaps not expunging as expected (especially StatisticsTracker.processedSeedsRecords) * CachedBdbMap.java do expunge on put(), replace() add low-memory-sensitive 'canary' to force expunge even if otherwise untriggered * CachedBdbMapTest.java add test of idle expunge in low-memory conditions * CrawlURI.java (getURI) direct access to String URI (don't reuse toString() functionally) * SeedRecord.java support for updating record with later/repeat report * StatisticsTracker.java remove put()s for processedSeedRecords, hostsLastFinished * BdbFrontier.java (getQueueFor) remove put, do in concurrent-compliant manner --- .../java/org/archive/util/CachedBdbMap.java | 71 +++++++------------ .../org/archive/util/CachedBdbMapTest.java | 26 +++++++ .../archive/crawler/datamodel/CrawlURI.java | 8 ++- .../archive/crawler/frontier/BdbFrontier.java | 25 ++++--- .../archive/crawler/reporting/SeedRecord.java | 43 +++++++---- .../crawler/reporting/StatisticsTracker.java | 35 ++++++--- 6 files changed, 124 insertions(+), 84 deletions(-) diff --git a/commons/src/main/java/org/archive/util/CachedBdbMap.java b/commons/src/main/java/org/archive/util/CachedBdbMap.java index 0e072e9d..71678d67 100644 --- a/commons/src/main/java/org/archive/util/CachedBdbMap.java +++ b/commons/src/main/java/org/archive/util/CachedBdbMap.java @@ -190,7 +190,7 @@ implements ConcurrentMap, Serializable { * high-concurrency client of this class in terms of number of objects * stored and performance impact. */ - private transient ConcurrentHashMap> memMap; + protected transient ConcurrentHashMap> memMap; protected transient ReferenceQueue refQueue; @@ -315,7 +315,7 @@ implements ConcurrentMap, Serializable { * mutating an unmapped value instance that would not be persisted. * (warning is emitted at most once per instance) */ - final private boolean LOG_ERROR_ON_DESIGN_VIOLATING_METHODS=false; + final private boolean LOG_ERROR_ON_DESIGN_VIOLATING_METHODS=true; /** * Simple structure to keep needed information about a DB Environment. @@ -421,45 +421,6 @@ implements ConcurrentMap, Serializable { 64 // est. number of concurrent threads ); this.refQueue = new ReferenceQueue(); - startExpunger(); - } - - private void startExpunger() { - /* - SettingsHandler ambientSettingsHandler = null; - try { - ambientSettingsHandler - = SettingsHandler.getThreadContextSettingsHandler(); - } catch (RuntimeException absorbed) { - } - if (ambientSettingsHandler != null) { - this.expunger = new Expunger("Expunger_" + dbName, refQueue, - logger, ambientSettingsHandler); - this.expunger.setDaemon(true); - this.expunger.setPriority(Thread.MAX_PRIORITY - 1); - this.expunger.start(); - } - */ - } - - /** - * - * @return true if expunder thread was stopped, false if no expunger is - * running. - */ - private boolean stopExpunger() { - /* - if (expunger != null) { - expunger.interrupt(); - try { - expunger.join(); - } catch (InterruptedException ignored) { - } - expunger = null; - return true; - } - */ - return false; } @SuppressWarnings("unchecked") @@ -595,7 +556,6 @@ implements ConcurrentMap, Serializable { } public synchronized void close() throws DatabaseException { - stopExpunger(); // Close out my bdb db. if (this.db != null) { try { @@ -798,6 +758,7 @@ implements ConcurrentMap, Serializable { if (pu == 1 && LOG_ERROR_ON_DESIGN_VIOLATING_METHODS) { logger.warning("design violating put() used on dbName=" + dbName); } + expungeStaleEntries(); // catchup all clears; possible disk IO int attemptTolerance = BDB_LOCK_ATTEMPT_TOLERANCE; while (--attemptTolerance > 0) { try { @@ -873,6 +834,7 @@ implements ConcurrentMap, Serializable { // warn on non-accretive use logger.warning("design violating replace(,,) used on dbName=" + dbName); } + expungeStaleEntries(); // catchup all clears; possible disk IO // make the ref wrappers SoftEntry newEntry = new SoftEntry(key, newValue, refQueue); @@ -1303,7 +1265,6 @@ implements ConcurrentMap, Serializable { String dbName = null; // Sync. memory and disk. useStatsSyncUsed.incrementAndGet(); - boolean expungerWasRunning = stopExpunger(); long startTime = 0; if (logger.isLoggable(Level.INFO)) { dbName = getDatabaseName(); @@ -1347,9 +1308,6 @@ implements ConcurrentMap, Serializable { this.diskMapSize.get() + ", mem " + this.memMap.size()); dumpExtraStats(); } - if (expungerWasRunning) { - startExpunger(); - } } /** log at INFO level, interesting stats if non-zero. */ @@ -1398,7 +1356,7 @@ implements ConcurrentMap, Serializable { * See #Expunger for dedicated thread expunger. */ @SuppressWarnings("unchecked") - private void expungeStaleEntries() { + protected void expungeStaleEntries() { int c = 0; long startTime = System.currentTimeMillis(); for(SoftEntry entry; (entry = (SoftEntry)refQueuePoll()) != null;) { @@ -1814,4 +1772,23 @@ implements ConcurrentMap, Serializable { private SoftEntry refQueuePoll() { return (SoftEntry)refQueue.poll(); } + + // + // Crude, probably unreliable/fragile but harmless mechanism to + // trigger expunge of cleared SoftReferences in low-memory + // conditions even without any of the other get/put triggers. + // + + protected SoftReference canary = + new SoftReference(new LowMemoryCanary()); + protected class LowMemoryCanary { + /** When collected/finalized -- as should be expected in + * low-memory conditions -- trigger an expunge and a + * new 'canary' insertion. */ + public void finalize() { + CachedBdbMap.this.expungeStaleEntries(); + CachedBdbMap.this.canary = + new SoftReference(new LowMemoryCanary()); + } + } } diff --git a/commons/src/test/java/org/archive/util/CachedBdbMapTest.java b/commons/src/test/java/org/archive/util/CachedBdbMapTest.java index 6d602829..657ca444 100644 --- a/commons/src/test/java/org/archive/util/CachedBdbMapTest.java +++ b/commons/src/test/java/org/archive/util/CachedBdbMapTest.java @@ -83,6 +83,32 @@ public class CachedBdbMapTest extends TmpDirTestCase { } } + /** + * Test that in scarce memory conditions, the memory map is + * expunged of otherwise unreferenced entries as expected. + * + * NOTE: this test may be especially fragile with regard to + * GC/timing issues; relies on timely finalization, which is + * never guaranteed by JVM/GC. + * + * @throws InterruptedException + */ + public void testMemMapCleared() throws InterruptedException { + assertEquals(cache.memMap.size(), 0); + for(int i=0; i < 10000; i++) { + cache.putIfAbsent(""+i, new HashMap()); + } + assertEquals(cache.memMap.size(), 10000); + assertEquals(cache.size(), 10000); + TestUtils.forceScarceMemory(); + Thread.sleep(1000); + // The 'canary' trick makes this explicit expunge, or + // an expunge triggered by a get() or put...(), unnecessary + // cache.expungeStaleEntries(); + System.out.println(cache.size()+","+cache.memMap.size()); + assertEquals(0, cache.memMap.size()); + } + public static void main(String [] args) { junit.textui.TestRunner.run(CachedBdbMapTest.class); } diff --git a/engine/src/main/java/org/archive/crawler/datamodel/CrawlURI.java b/engine/src/main/java/org/archive/crawler/datamodel/CrawlURI.java index 327cfcb1..488b3343 100644 --- a/engine/src/main/java/org/archive/crawler/datamodel/CrawlURI.java +++ b/engine/src/main/java/org/archive/crawler/datamodel/CrawlURI.java @@ -1413,6 +1413,12 @@ implements ProcessorURI, MultiReporter, Serializable, OverlayContext { return this.uuri; } + /** + * @return String of URI + */ + public String getURI() { + return getUURI().toCustomString(); + } /** * @return path (hop-types) from seed */ @@ -1721,7 +1727,7 @@ implements ProcessorURI, MultiReporter, Serializable, OverlayContext { * @return The UURI this CandidateURI wraps as a string */ public String toString() { - return getUURI().toString(); + return getUURI().toCustomString(); } 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 03246ce1..5f78a013 100644 --- a/engine/src/main/java/org/archive/crawler/frontier/BdbFrontier.java +++ b/engine/src/main/java/org/archive/crawler/frontier/BdbFrontier.java @@ -144,20 +144,8 @@ implements Serializable, Checkpointable { * @return the found or created BdbWorkQueue */ protected WorkQueue getQueueFor(CrawlURI curi) { - WorkQueue wq; String classKey = curi.getClassKey(); - synchronized (allQueues) { - wq = (WorkQueue)allQueues.get(classKey); - if (wq == null) { - wq = new BdbWorkQueue(classKey, this); - //TODO:SPRINGY set overrides - wq.setTotalBudget(getQueueTotalBudget()); - //TODO:SPRINGY set overrides - getQueuePrecedencePolicy().queueCreated(wq); - allQueues.put(classKey, wq); - } - } - return wq; + return getQueueFor(classKey); } /** @@ -171,6 +159,17 @@ implements Serializable, Checkpointable { assert Thread.currentThread() == managerThread; 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; + } + } return wq; } diff --git a/engine/src/main/java/org/archive/crawler/reporting/SeedRecord.java b/engine/src/main/java/org/archive/crawler/reporting/SeedRecord.java index 50ab0c3d..6e39e386 100644 --- a/engine/src/main/java/org/archive/crawler/reporting/SeedRecord.java +++ b/engine/src/main/java/org/archive/crawler/reporting/SeedRecord.java @@ -20,11 +20,11 @@ package org.archive.crawler.reporting; import java.io.Serializable; +import java.util.logging.Logger; import org.archive.crawler.datamodel.CrawlURI; import org.archive.crawler.datamodel.CoreAttributeConstants; - /** * Record of all interesting info about the most-recent * processing of a specific seed. @@ -33,9 +33,11 @@ import org.archive.crawler.datamodel.CoreAttributeConstants; */ public class SeedRecord implements CoreAttributeConstants, Serializable { private static final long serialVersionUID = -8455358640509744478L; + private static Logger logger = + Logger.getLogger(SeedRecord.class.getName()); private final String uri; private int statusCode; - private final String disposition; + private String disposition; private String redirectUri; /** @@ -47,17 +49,8 @@ public class SeedRecord implements CoreAttributeConstants, Serializable { */ public SeedRecord(CrawlURI curi, String disposition) { super(); - this.uri = curi.toString(); - this.statusCode = curi.getFetchStatus(); - this.disposition = disposition; - if (statusCode==301 || statusCode == 302) { - for (CrawlURI cauri: curi.getOutCandidates()) { - if("location:".equalsIgnoreCase(cauri.getViaContext(). - toString())) { - redirectUri = cauri.toString(); - } - } - } + this.uri = curi.getURI(); + updateWith(curi,disposition); } /** @@ -89,6 +82,30 @@ public class SeedRecord implements CoreAttributeConstants, Serializable { this.redirectUri = redirectUri; } + /** + * A later/repeat report of the same seed has arrived; update with + * latest. Should be rare/never? + * + * @param curi + */ + public void updateWith(CrawlURI curi,String disposition) { + if(!this.uri.equals(curi.getURI())) { + logger.warning("SeedRecord URI changed: "+uri+"->"+curi.getURI()); + } + this.statusCode = curi.getFetchStatus(); + this.disposition = disposition; + if (statusCode==301 || statusCode == 302) { + for (CrawlURI cauri: curi.getOutCandidates()) { + if("location:".equalsIgnoreCase(cauri.getViaContext(). + toString())) { + redirectUri = cauri.toString(); + } + } + } else { + redirectUri = null; + } + } + /** * @return Returns the disposition. */ 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 2d73a84a..209005da 100644 --- a/engine/src/main/java/org/archive/crawler/reporting/StatisticsTracker.java +++ b/engine/src/main/java/org/archive/crawler/reporting/StatisticsTracker.java @@ -301,8 +301,8 @@ public class StatisticsTracker new ConcurrentHashMap(); // temp dummy protected ConcurrentMap hostsBytes = new ConcurrentHashMap(); // temp dummy - protected ConcurrentMap hostsLastFinished = - new ConcurrentHashMap(); // temp dummy + protected ConcurrentMap hostsLastFinished = + new ConcurrentHashMap(); // temp dummy /** Keep track of URL counts per host per seed */ protected @@ -346,7 +346,7 @@ public class StatisticsTracker this.hostsBytes = bdb.getBigMap("hostsBytes", false, String.class, AtomicLong.class); this.hostsLastFinished = bdb.getBigMap("hostsLastFinished", - false, String.class, Long.class); + false, String.class, AtomicLong.class); this.processedSeedsRecords = bdb.getBigMap("processedSeedsRecords", false, String.class, SeedRecord.class); @@ -638,9 +638,17 @@ public class StatisticsTracker * host was last finished processing. If no URI has been completed for host * -1 will be returned. */ - public long getHostLastFinished(String host){ - Long l = (Long)hostsLastFinished.get(host); - return (l != null)? l.longValue(): -1; + 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; + } + } + return fini; } /** @@ -682,8 +690,15 @@ public class StatisticsTracker */ private void handleSeed(CrawlURI curi, String disposition) { if(curi.isSeed()){ - SeedRecord sr = new SeedRecord(curi, disposition); - processedSeedsRecords.put(sr.getUri(), sr); + 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); + } + } } } @@ -740,7 +755,7 @@ public class StatisticsTracker getReportValue(hostsBytes, hostname)); long time = new Long(System.currentTimeMillis()); - hostsLastFinished.put(hostname, time); + getHostLastFinished(hostname).set(time); hostsLastFinishedTop.update(hostname, time); } @@ -857,7 +872,7 @@ public class StatisticsTracker if(f.exists() && !controller.isRunning() && controller.hasStarted()) { // controller already started and stopped and file exists // so, don't overwrite - logger.info("reusing report: " + f.getAbsolutePath()); + logger.warning("reusing report: " + f.getAbsolutePath()); return f; }