Author: thomasm
Date: Thu Nov  9 16:52:49 2017
New Revision: 1814745

URL: http://svn.apache.org/viewvc?rev=1814745&view=rev
Log:
OAK-5519 Skip problematic binaries instead of blocking indexing

Modified:
    
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCache.java
    
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
    
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/TextExtractionStatsMBean.java
    
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/binary/BinaryTextExtractor.java
    
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCacheTest.java

Modified: 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCache.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCache.java?rev=1814745&r1=1814744&r2=1814745&view=diff
==============================================================================
--- 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCache.java
 (original)
+++ 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCache.java
 Thu Nov  9 16:52:49 2017
@@ -19,8 +19,23 @@
 
 package org.apache.jackrabbit.oak.plugins.index.lucene;
 
+import java.io.File;
+import java.io.FileInputStream;
+import java.io.FileOutputStream;
 import java.io.IOException;
+import java.util.Map.Entry;
+import java.util.Properties;
+import java.util.concurrent.Callable;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Future;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.ThreadPoolExecutor;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicInteger;
 
 import javax.annotation.CheckForNull;
 import javax.annotation.Nonnull;
@@ -32,6 +47,7 @@ import org.apache.jackrabbit.oak.api.Blo
 import org.apache.jackrabbit.oak.cache.CacheStats;
 import org.apache.jackrabbit.oak.commons.IOUtils;
 import org.apache.jackrabbit.oak.plugins.index.fulltext.ExtractedText;
+import 
org.apache.jackrabbit.oak.plugins.index.fulltext.ExtractedText.ExtractionResult;
 import 
org.apache.jackrabbit.oak.plugins.index.fulltext.PreExtractedTextProvider;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -39,6 +55,18 @@ import org.slf4j.LoggerFactory;
 import static org.apache.jackrabbit.oak.commons.PathUtils.concat;
 
 public class ExtractedTextCache {
+    private static final boolean CACHE_ONLY_SUCCESS =
+            Boolean.getBoolean("oak.extracted.cacheOnlySuccess");
+    private static final int EXTRACTION_TIMEOUT_SECONDS =
+            Integer.getInteger("oak.extraction.timeoutSeconds", 60);
+    private static final int EXTRACTION_MAX_THREADS =
+            Integer.getInteger("oak.extraction.maxThreads", 10);
+    private static final boolean EXTRACT_IN_CALLER_THREAD =
+            Boolean.getBoolean("oak.extraction.inCallerThread");
+    private static final boolean EXTRACT_FORGET_TIMEOUT =
+            Boolean.getBoolean("oak.extraction.forgetTimeout");
+
+    private static final String TIMEOUT_MAP = 
"textExtractionTimeout.properties";
     private static final String EMPTY_STRING = "";
     private static final Logger log = 
LoggerFactory.getLogger(ExtractedTextCache.class);
     private volatile PreExtractedTextProvider extractedTextProvider;
@@ -48,14 +76,20 @@ public class ExtractedTextCache {
     private long totalTime;
     private int preFetchedCount;
     private final Cache<String, String> cache;
+    private final ConcurrentHashMap<String, String> timeoutMap;
+    private final File indexDir;
     private final CacheStats cacheStats;
     private final boolean alwaysUsePreExtractedCache;
+    private volatile ExecutorService executorService;
+    private volatile int timeoutCount;
+    private long extractionTimeoutMillis = EXTRACTION_TIMEOUT_SECONDS * 1000;
 
     public ExtractedTextCache(long maxWeight, long expiryTimeInSecs){
-        this(maxWeight, expiryTimeInSecs, false);
+        this(maxWeight, expiryTimeInSecs, false, null);
     }
 
-    public ExtractedTextCache(long maxWeight, long expiryTimeInSecs, boolean 
alwaysUsePreExtractedCache) {
+    public ExtractedTextCache(long maxWeight, long expiryTimeInSecs, boolean 
alwaysUsePreExtractedCache,
+            File indexDir) {
         if (maxWeight > 0) {
             cache = CacheBuilder.newBuilder()
                     .weigher(EmpiricalWeigher.INSTANCE)
@@ -70,6 +104,9 @@ public class ExtractedTextCache {
             cacheStats = null;
         }
         this.alwaysUsePreExtractedCache = alwaysUsePreExtractedCache;
+        this.timeoutMap = new ConcurrentHashMap<String, String>();
+        this.indexDir = indexDir;
+        loadTimeoutMap();
     }
 
     /**
@@ -90,39 +127,53 @@ public class ExtractedTextCache {
                 ExtractedText text = 
extractedTextProvider.getText(propertyPath, blob);
                 if (text != null) {
                     preFetchedCount++;
-                    switch (text.getExtractionResult()) {
-                        case SUCCESS:
-                            result = text.getExtractedText().toString();
-                            break;
-                        case ERROR:
-                            result = LuceneIndexEditor.TEXT_EXTRACTION_ERROR;
-                            break;
-                        case EMPTY:
-                            result = EMPTY_STRING;
-                            break;
-                    }
+                    result = getText(text);
                 }
             } catch (IOException e) {
                 log.warn("Error occurred while fetching pre extracted text for 
{}", propertyPath, e);
             }
         }
-
         String id = blob.getContentIdentity();
         if (cache != null && id != null && result == null) {
             result = cache.getIfPresent(id);
         }
+        if (result == null && id != null) {
+            result = timeoutMap.get(id);
+        }
         return result;
     }
 
     public void put(@Nonnull Blob blob, @Nonnull ExtractedText extractedText) {
         String id = blob.getContentIdentity();
-        if (extractedText.getExtractionResult() == 
ExtractedText.ExtractionResult.SUCCESS
-                && cache != null
-                && id != null) {
-            cache.put(id, extractedText.getExtractedText().toString());
+        if (cache != null && id != null) {
+            if (extractedText.getExtractionResult() == ExtractionResult.SUCCESS
+                    || !CACHE_ONLY_SUCCESS) {
+                cache.put(id, getText(extractedText));
+            }
         }
     }
 
+    public void putTimeout(@Nonnull Blob blob, @Nonnull ExtractedText 
extractedText) {
+        if (EXTRACT_FORGET_TIMEOUT) {
+            return;
+        }
+        String id = blob.getContentIdentity();
+        timeoutMap.put(id, getText(extractedText));
+        storeTimeoutMap();
+    }
+
+    private static String getText(ExtractedText text) {
+        switch (text.getExtractionResult()) {
+        case SUCCESS:
+            return text.getExtractedText().toString();
+        case ERROR:
+            return LuceneIndexEditor.TEXT_EXTRACTION_ERROR;
+        case EMPTY:
+            return EMPTY_STRING;
+        }
+        throw new IllegalArgumentException();
+    }
+
     public void addStats(int count, long timeInMillis, long bytesRead, long 
textLength){
         this.textExtractionCount += count;
         this.totalTime += timeInMillis;
@@ -130,7 +181,7 @@ public class ExtractedTextCache {
         this.totalTextSize += textLength;
     }
 
-    public TextExtractionStatsMBean getStatsMBean(){
+    public TextExtractionStatsMBean getStatsMBean() {
         return new TextExtractionStatsMBean() {
             @Override
             public boolean isPreExtractedTextProviderConfigured() {
@@ -166,6 +217,11 @@ public class ExtractedTextCache {
             public boolean isAlwaysUsePreExtractedCache() {
                 return alwaysUsePreExtractedCache;
             }
+
+            @Override
+            public int getTimeoutCount() {
+                return timeoutCount;
+            }
         };
     }
 
@@ -216,4 +272,133 @@ public class ExtractedTextCache {
             return (int) size;
         }
     }
+
+    public void close() {
+        resetCache();
+        // don't clean the persistent map on purpose, so we don't re-try
+        // after restarting the service or so
+        closeExecutorService();
+    }
+
+    public void process(String name, Callable<Void> callable) throws 
InterruptedException, Throwable {
+        Callable<Void> callable2 = new Callable<Void>() {
+            @Override
+            public Void call() throws Exception {
+                Thread t = Thread.currentThread();
+                String oldThreadName = t.getName();
+                t.setName(oldThreadName + ": " + name);
+                try {
+                    return callable.call();
+                } finally {
+                    Thread.currentThread().setName(oldThreadName);
+                }
+            }
+        };
+        try {
+            if (EXTRACT_IN_CALLER_THREAD) {
+                callable2.call();
+            } else {
+                Future<Void> future = getExecutor().submit(callable2);
+                future.get(extractionTimeoutMillis, TimeUnit.MILLISECONDS);
+            }
+        } catch (TimeoutException e) {
+            timeoutCount++;
+            throw e;
+        } catch (InterruptedException e) {
+            throw e;
+        } catch (ExecutionException e) {
+            throw e.getCause();
+        }
+    }
+
+    public void setExtractionTimeoutMillis(int extractionTimeoutMillis) {
+        this.extractionTimeoutMillis = extractionTimeoutMillis;
+    }
+
+    private ExecutorService getExecutor() {
+        if (executorService == null) {
+            createExecutor();
+        }
+        return executorService;
+    }
+
+    private synchronized void createExecutor() {
+        if (executorService != null) {
+            return;
+        }
+        log.debug("ExtractedTextCache createExecutor " + this);
+        ThreadPoolExecutor executor = new ThreadPoolExecutor(1, 
EXTRACTION_MAX_THREADS,
+                60L, TimeUnit.SECONDS,
+                new LinkedBlockingQueue<Runnable>(), new ThreadFactory() {
+            private final AtomicInteger counter = new AtomicInteger();
+            private final Thread.UncaughtExceptionHandler handler = new 
Thread.UncaughtExceptionHandler() {
+                @Override
+                public void uncaughtException(Thread t, Throwable e) {
+                    log.warn("Error occurred in asynchronous processing ", e);
+                }
+            };
+            @Override
+            public Thread newThread(@Nonnull Runnable r) {
+                Thread thread = new Thread(r, createName());
+                thread.setDaemon(true);
+                thread.setPriority(Thread.MIN_PRIORITY);
+                thread.setUncaughtExceptionHandler(handler);
+                return thread;
+            }
+
+            private String createName() {
+                int index = counter.getAndIncrement();
+                return "oak binary text extractor" + (index == 0 ? "" : " " + 
index);
+            }
+        });
+        executor.setKeepAliveTime(1, TimeUnit.MINUTES);
+        executor.allowCoreThreadTimeOut(true);
+        executorService = executor;
+    }
+
+    private synchronized void closeExecutorService() {
+        if (executorService != null) {
+            log.debug("ExtractedTextCache closeExecutorService " + this);
+            executorService.shutdown();
+            try {
+                executorService.awaitTermination(1, TimeUnit.MINUTES);
+            } catch (InterruptedException e) {
+                log.warn("Interrupted", e);
+            }
+            executorService = null;
+        }
+    }
+
+    private synchronized void loadTimeoutMap() {
+        if (indexDir == null || !indexDir.exists()) {
+            return;
+        }
+        try (FileInputStream in = new FileInputStream(
+                new File(indexDir, TIMEOUT_MAP))) {
+            Properties prop = new Properties();
+            prop.load(in);
+            for(Entry<Object, Object> e : prop.entrySet()) {
+                timeoutMap.put(e.getKey().toString(), e.getValue().toString());
+            }
+        } catch (Exception e) {
+            log.warn("Could not load timeout map {} from {}",
+                    TIMEOUT_MAP, indexDir, e);
+        }
+    }
+
+    private synchronized void storeTimeoutMap() {
+        if (indexDir == null || !indexDir.exists()) {
+            return;
+        }
+        try (FileOutputStream out = new FileOutputStream(
+                    new File(indexDir, TIMEOUT_MAP))) {
+            Properties prop = new Properties();
+            prop.putAll(timeoutMap);
+            prop.store(out, "Text extraction timed out for the following 
binaries, and will not be retried");
+        } catch (Exception e) {
+            log.warn("Could not store timeout map {} from {}",
+                    TIMEOUT_MAP, indexDir, e);
+        }
+    }
+
 }

Modified: 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java?rev=1814745&r1=1814744&r2=1814745&view=diff
==============================================================================
--- 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
 (original)
+++ 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
 Thu Nov  9 16:52:49 2017
@@ -434,6 +434,10 @@ public class LuceneIndexProviderService
             executorService.awaitTermination(1, TimeUnit.MINUTES);
         }
 
+        if (extractedTextCache != null) {
+            extractedTextCache.close();
+        }
+
         InfoStream.setDefault(InfoStream.NO_OUTPUT);
     }
 
@@ -673,7 +677,11 @@ public class LuceneIndexProviderService
         boolean alwaysUsePreExtractedCache = 
PropertiesUtil.toBoolean(config.get(PROP_PRE_EXTRACTED_TEXT_ALWAYS_USE),
                 PROP_PRE_EXTRACTED_TEXT_ALWAYS_USE_DEFAULT);
 
-        extractedTextCache = new ExtractedTextCache(cacheSizeInMB * ONE_MB, 
cacheExpiryInSecs, alwaysUsePreExtractedCache);
+        extractedTextCache = new ExtractedTextCache(
+                cacheSizeInMB * ONE_MB,
+                cacheExpiryInSecs,
+                alwaysUsePreExtractedCache,
+                indexDir);
         if (extractedTextProvider != null){
             registerExtractedTextProvider(extractedTextProvider);
         }

Modified: 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/TextExtractionStatsMBean.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/TextExtractionStatsMBean.java?rev=1814745&r1=1814744&r2=1814745&view=diff
==============================================================================
--- 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/TextExtractionStatsMBean.java
 (original)
+++ 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/TextExtractionStatsMBean.java
 Thu Nov  9 16:52:49 2017
@@ -38,4 +38,6 @@ public interface TextExtractionStatsMBea
     String getExtractedTextSize();
 
     String getBytesRead();
+
+    int getTimeoutCount();
 }

Modified: 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/binary/BinaryTextExtractor.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/binary/BinaryTextExtractor.java?rev=1814745&r1=1814744&r2=1814745&view=diff
==============================================================================
--- 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/binary/BinaryTextExtractor.java
 (original)
+++ 
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/binary/BinaryTextExtractor.java
 Thu Nov  9 16:52:49 2017
@@ -26,6 +26,8 @@ import java.util.ArrayList;
 import java.util.Collections;
 import java.util.List;
 import java.util.Set;
+import java.util.concurrent.Callable;
+import java.util.concurrent.TimeoutException;
 
 import javax.annotation.Nullable;
 
@@ -140,16 +142,21 @@ public class BinaryTextExtractor {
         if (log.isDebugEnabled()) {
             log.debug("Extracting {}, {} bytes, id {}", path, length, 
v.getContentIdentity());
         }
-        String oldThreadName = null;
-        if (length > SMALL_BINARY) {
-            Thread t = Thread.currentThread();
-            oldThreadName = t.getName();
-            t.setName(oldThreadName + ": Extracting " + path + ", " + length + 
" bytes");
-        }
         try {
             CountingInputStream stream = new CountingInputStream(new 
LazyInputStream(new BlobByteSource(v)));
             try {
-                getParser().parse(stream, handler, metadata, new 
ParseContext());
+                if (length > SMALL_BINARY) {
+                    String name = "Extracting " + path + ", " + length + " 
bytes";
+                    extractedTextCache.process(name, new Callable<Void>() {
+                        @Override
+                        public Void call() throws Exception {
+                            getParser().parse(stream, handler, metadata, new 
ParseContext());
+                            return null;
+                        }
+                    });
+                } else {
+                    getParser().parse(stream, handler, metadata, new 
ParseContext());
+                }
             } finally {
                 bytesRead = stream.getCount();
                 stream.close();
@@ -159,6 +166,13 @@ public class BinaryTextExtractor {
             // not being present. This is equivalent to disabling
             // selected media types in configuration, so we can simply
             // ignore these errors.
+        } catch (TimeoutException t) {
+            log.warn(
+                    "[{}] Failed to extract text from a binary property due to 
timeout: {}.",
+                    getIndexName(), path);
+            extractedTextCache.put(v, ExtractedText.ERROR);
+            extractedTextCache.putTimeout(v, ExtractedText.ERROR);
+            return TEXT_EXTRACTION_ERROR;
         } catch (Throwable t) {
             // Capture and report any other full text extraction problems.
             // The special STOP exception is used for normal termination.
@@ -172,10 +186,6 @@ public class BinaryTextExtractor {
                 extractedTextCache.put(v, ExtractedText.ERROR);
                 return TEXT_EXTRACTION_ERROR;
             }
-        } finally {
-            if (oldThreadName != null) {
-                Thread.currentThread().setName(oldThreadName);
-            }
         }
         String result = handler.toString();
         if (bytesRead > 0) {

Modified: 
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCacheTest.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCacheTest.java?rev=1814745&r1=1814744&r2=1814745&view=diff
==============================================================================
--- 
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCacheTest.java
 (original)
+++ 
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/ExtractedTextCacheTest.java
 Thu Nov  9 16:52:49 2017
@@ -27,15 +27,20 @@ import org.apache.jackrabbit.oak.plugins
 import org.apache.jackrabbit.oak.plugins.memory.ArrayBasedBlob;
 import org.junit.Test;
 
+import static org.junit.Assert.fail;
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertNotNull;
 import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertTrue;
 import static org.mockito.Matchers.any;
 import static org.mockito.Matchers.anyString;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.verifyZeroInteractions;
 import static org.mockito.Mockito.when;
 
+import java.util.concurrent.Callable;
+import java.util.concurrent.TimeoutException;
+
 public class ExtractedTextCacheTest {
 
     @Test
@@ -84,7 +89,7 @@ public class ExtractedTextCacheTest {
         cache.put(b, new ExtractedText(ExtractionResult.ERROR, "test hello"));
 
         text = cache.get("/a", "foo", b, false);
-        assertNull(text);
+        assertEquals(LuceneIndexEditor.TEXT_EXTRACTION_ERROR, text);
     }
 
     @Test
@@ -124,7 +129,7 @@ public class ExtractedTextCacheTest {
 
     @Test
     public void preExtractionAlwaysUse() throws Exception{
-        ExtractedTextCache cache = new ExtractedTextCache(10 * 
FileUtils.ONE_MB, 100, true);
+        ExtractedTextCache cache = new ExtractedTextCache(10 * 
FileUtils.ONE_MB, 100, true, null);
         PreExtractedTextProvider provider = 
mock(PreExtractedTextProvider.class);
 
         cache.setExtractedTextProvider(provider);
@@ -135,6 +140,51 @@ public class ExtractedTextCacheTest {
         assertEquals("bar", text);
     }
 
+    @Test
+    public void rememberTimeout() throws Exception{
+        ExtractedTextCache cache = new ExtractedTextCache(0, 0, false, null);
+        Blob b = new IdBlob("hello", "a");
+        cache.put(b, ExtractedText.ERROR);
+        assertNull(cache.get("/a", "foo", b, false));
+        cache.putTimeout(b, ExtractedText.ERROR);
+        assertEquals(LuceneIndexEditor.TEXT_EXTRACTION_ERROR, cache.get("/a", 
"foo", b, false));
+    }
+
+    @Test
+    public void process() throws Throwable {
+        ExtractedTextCache cache = new ExtractedTextCache(0, 0, false, null);
+        try {
+            cache.process("test", new Callable<Void>() {
+                @Override
+                public Void call() throws Exception {
+                    throw new OutOfMemoryError();
+                }
+            });
+            fail();
+        } catch (OutOfMemoryError e) {
+            // expected
+        }
+        assertEquals(0, cache.getStatsMBean().getTimeoutCount());
+        cache.setExtractionTimeoutMillis(10);
+        long time = System.currentTimeMillis();
+        try {
+            cache.process("test", new Callable<Void>() {
+                @Override
+                public Void call() throws Exception {
+                    // this happens in the background, so doesn't block the 
test
+                    Thread.sleep(10000);
+                    return null;
+                }
+            });
+            fail();
+        } catch (TimeoutException e) {
+            // expected
+        }
+        time = System.currentTimeMillis() - time;
+        assertTrue("" + time, time < 5000);
+        assertEquals(1, cache.getStatsMBean().getTimeoutCount());
+    }
+
     private static class IdBlob extends ArrayBasedBlob {
         final String id;
 


Reply via email to